Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions backend/agents/budget_guard.py
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,7 @@ def __init__(

self.config_path = Path(config_path)
self.budget_config: dict[str, Any] = {}
self.downgrade_paths: dict[str, str] = {}
self.state = BudgetState()
self._load_config()

Expand All @@ -134,6 +135,8 @@ def _load_config(self) -> None:
with open(self.config_path) as f:
self.budget_config = yaml.safe_load(f)

self.downgrade_paths = self.budget_config.get("downgrade_paths", {})

budget = self.budget_config.get("budget", {})
self.state.hard_limit = budget.get("hard_limit", 200.0)
logger.debug("budget_guard.config_loaded", limit=self.state.hard_limit)
Expand Down Expand Up @@ -269,6 +272,10 @@ def can_spend(self, estimated_cost: float) -> tuple[bool, str | None]:

return True, None

def downgrade_model_for(self, current_model: str) -> str | None:
"""Return the cheaper model for `current_model`, or None if there is no successor."""
return self.downgrade_paths.get(current_model)

def get_downgrade_targets(self) -> list[tuple[str, str]]:
"""
Get list of agents to downgrade based on current spend.
Expand Down
12 changes: 0 additions & 12 deletions backend/agents/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -451,18 +451,6 @@ def reload(self) -> None:
with self._lock:
self._instances.clear()

def get_budget_config(self) -> dict[str, Any]:
"""Get budget configuration from the YAML file."""
with open(self.config_path) as f:
raw_config = yaml.safe_load(f)
return raw_config.get("budget", {})

def get_downgrade_paths(self) -> dict[str, str]:
"""Get model downgrade paths from the YAML file."""
with open(self.config_path) as f:
raw_config = yaml.safe_load(f)
return raw_config.get("downgrade_paths", {})


# Global registry instance (singleton pattern)
_registry: AgentRegistry | None = None
Expand Down
10 changes: 10 additions & 0 deletions backend/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,11 @@ class Config:
sqlite_path: str
max_clarifying_questions: int
enable_phase4: bool
engine_lease_ttl: float = 120.0
engine_heartbeat_interval: float = 20.0
engine_reaper_interval: float = 30.0
engine_worker_count: int = 4
engine_max_attempts: int = 3
agents_yaml: dict[str, Any] = field(default_factory=dict)
budget_yaml: dict[str, Any] = field(default_factory=dict)
llm_yaml: dict[str, Any] = field(default_factory=dict)
Expand All @@ -74,6 +79,11 @@ def load(cls) -> Config:
sqlite_path=os.getenv("SQLITE_PATH", default_sqlite),
max_clarifying_questions=_env_int("MAX_CLARIFYING_QUESTIONS", 6),
enable_phase4=_env_bool("ENABLE_PHASE4", True),
engine_lease_ttl=_env_float("ENGINE_LEASE_TTL", 120.0),
engine_heartbeat_interval=_env_float("ENGINE_HEARTBEAT_INTERVAL", 20.0),
engine_reaper_interval=_env_float("ENGINE_REAPER_INTERVAL", 30.0),
engine_worker_count=_env_int("ENGINE_WORKER_COUNT", 4),
engine_max_attempts=_env_int("ENGINE_MAX_ATTEMPTS", 3),
agents_yaml=_load_yaml(CONFIG_DIR / "agents.yaml"),
budget_yaml=_load_yaml(CONFIG_DIR / "budget.yaml"),
llm_yaml=_load_yaml(CONFIG_DIR / "llm.yaml"),
Expand Down
Empty file added backend/engine/__init__.py
Empty file.
83 changes: 83 additions & 0 deletions backend/engine/models.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
"""SQL schema + typed result models for the engine store."""
from __future__ import annotations

from pydantic import BaseModel

PHASE_STATUS = {"blocked", "open", "complete"}
TASK_STATUS = {"blocked", "ready", "claimed", "running", "done", "failed"}
GATE = {"none", "pending", "approved", "rejected"}

SCHEMA_SQL = """
CREATE TABLE IF NOT EXISTS runs (
run_id TEXT PRIMARY KEY,
idea TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'running',
current_phase INTEGER NOT NULL DEFAULT 0,
budget_limit REAL NOT NULL DEFAULT 200.0,
created_at REAL NOT NULL
);
CREATE TABLE IF NOT EXISTS phases (
run_id TEXT NOT NULL,
name TEXT NOT NULL,
phase_order INTEGER NOT NULL,
status TEXT NOT NULL DEFAULT 'blocked',
gate TEXT NOT NULL DEFAULT 'none',
seeded INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (run_id, name)
);
CREATE TABLE IF NOT EXISTS tasks (
task_id TEXT PRIMARY KEY,
run_id TEXT NOT NULL,
phase TEXT NOT NULL,
phase_order INTEGER NOT NULL,
agent_id TEXT NOT NULL,
input TEXT NOT NULL DEFAULT '{}',
depends_on TEXT NOT NULL DEFAULT '[]',
status TEXT NOT NULL DEFAULT 'blocked',
owner TEXT,
version INTEGER NOT NULL DEFAULT 0,
attempts INTEGER NOT NULL DEFAULT 0,
lease_expires REAL,
created_at REAL NOT NULL,
claimed_at REAL,
model TEXT,
sim_cost REAL NOT NULL DEFAULT 0.0,
result TEXT
);
CREATE INDEX IF NOT EXISTS idx_tasks_ready ON tasks (run_id, status, phase_order, created_at);
CREATE TABLE IF NOT EXISTS state (
run_id TEXT NOT NULL,
key TEXT NOT NULL,
value TEXT NOT NULL,
version INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (run_id, key)
);
CREATE TABLE IF NOT EXISTS spend (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id TEXT NOT NULL,
task_id TEXT,
agent_id TEXT,
cost REAL NOT NULL,
model TEXT,
ts REAL NOT NULL
);
CREATE TABLE IF NOT EXISTS events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id TEXT NOT NULL,
type TEXT NOT NULL,
payload TEXT NOT NULL,
worker_pid INTEGER,
ts REAL NOT NULL
);
"""


class ClaimResult(BaseModel):
task_id: str
run_id: str
phase: str
phase_order: int
agent_id: str
input: dict
model: str | None
version: int
86 changes: 86 additions & 0 deletions backend/engine/phases.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
"""Loader/validator for config/phases.yaml — the six-phase source of truth."""
from __future__ import annotations

from dataclasses import dataclass
from pathlib import Path

import yaml


@dataclass(frozen=True)
class AgentSpec:
agent_id: str
reads: list[str]
writes: str
sim_cost: float
depends_on: list[str]


@dataclass(frozen=True)
class PhaseSpec:
name: str
order: int
gate: str
agents: dict[str, AgentSpec]


class PhasesConfig:
def __init__(self, phases: list[PhaseSpec]):
self._phases = sorted(phases, key=lambda p: p.order)
self._by_name = {p.name: p for p in self._phases}

@classmethod
def load(cls, path: str = "config/phases.yaml") -> "PhasesConfig":
raw = yaml.safe_load(Path(path).read_text(encoding="utf-8"))
phases: list[PhaseSpec] = []
for p in raw["phases"]:
agents = {
aid: AgentSpec(
agent_id=aid,
reads=list(a.get("reads", [])),
writes=a["writes"],
sim_cost=float(a.get("sim_cost", 0.0)),
depends_on=list(a.get("depends_on", [])),
)
for aid, a in p["agents"].items()
}
phases.append(
PhaseSpec(name=p["name"], order=int(p["order"]), gate=p.get("gate", "none"), agents=agents)
)
cfg = cls(phases)
cfg._validate()
return cfg

def _validate(self) -> None:
orders = [p.order for p in self._phases]
if orders != list(range(len(self._phases))):
raise ValueError(f"phase orders must be 0..N-1, got {orders}")
for p in self._phases:
for aid, spec in p.agents.items():
for dep in spec.depends_on:
if dep == aid:
raise ValueError(f"{p.name}.{aid} depends on itself")
if dep not in p.agents:
raise ValueError(
f"{p.name}.{aid} depends on unknown intra-phase agent {dep!r}"
)

@property
def phase_names(self) -> list[str]:
return [p.name for p in self._phases]

def order_of(self, name: str) -> int:
return self._by_name[name].order

def gate_of(self, name: str) -> str:
return self._by_name[name].gate

def agents_of(self, name: str) -> dict[str, AgentSpec]:
return dict(self._by_name[name].agents)

def all_agent_ids(self) -> list[str]:
return [aid for p in self._phases for aid in p.agents]


def load_downgrade_paths(path: str = "config/budget.yaml") -> dict[str, str]:
return yaml.safe_load(Path(path).read_text(encoding="utf-8")).get("downgrade_paths", {})
83 changes: 83 additions & 0 deletions backend/engine/scheduler.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
"""Pure scheduling logic. No DB, no I/O — takes plain dicts, returns plans."""
from __future__ import annotations

from backend.engine.phases import PhasesConfig


def task_id(run_id: str, phase: str, agent_id: str) -> str:
return f"{run_id}:{phase}:{agent_id}"


def seed_specs_for_phase(
cfg: PhasesConfig, run_id: str, phase_name: str, base_models: dict[str, str]
) -> list[dict]:
order = cfg.order_of(phase_name)
specs: list[dict] = []
for aid, spec in cfg.agents_of(phase_name).items():
specs.append(
{
"task_id": task_id(run_id, phase_name, aid),
"agent_id": aid,
"phase": phase_name,
"phase_order": order,
"depends_on": [task_id(run_id, phase_name, dep) for dep in spec.depends_on],
"sim_cost": spec.sim_cost,
"model": base_models.get(aid),
"input_keys": list(spec.reads),
}
)
return specs


def compute_ready(tasks: list[dict], phases: list[dict]) -> list[str]:
open_phases = {p["name"] for p in phases if p["status"] == "open"}
done = {t["task_id"] for t in tasks if t["status"] == "done"}
ready: list[str] = []
for t in tasks:
if t["status"] != "blocked" or t["phase"] not in open_phases:
continue
if all(dep in done for dep in t["depends_on"]):
ready.append(t["task_id"])
return ready


def advance(phases: list[dict], tasks: list[dict], cfg: PhasesConfig) -> dict:
by_phase: dict[str, list[dict]] = {}
for t in tasks:
by_phase.setdefault(t["phase"], []).append(t)

complete_phases: list[str] = []
open_gates: list[str] = []
open_phases: list[str] = []

ordered = sorted(phases, key=lambda p: p["phase_order"])
for p in ordered:
if p["status"] != "open":
continue
pts = by_phase.get(p["name"], [])
# A phase completes only if seeded, non-empty, and all its tasks are done.
if not p.get("seeded") or not pts:
continue
if all(t["status"] == "done" for t in pts):
complete_phases.append(p["name"])
gate = cfg.gate_of(p["name"])
if gate != "none":
open_gates.append(gate) # next phase waits for submit_approval
else:
nxt = next((q for q in ordered if q["phase_order"] == p["phase_order"] + 1), None)
if nxt is not None:
open_phases.append(nxt["name"])
return {"complete_phases": complete_phases, "open_gates": open_gates, "open_phases": open_phases}


def resolve_model(
agent_id: str,
base_model: str | None,
spend_ratio: float,
downgrade_paths: dict[str, str],
skip_list: set[str],
threshold: float = 0.85,
) -> str | None:
if base_model is None or agent_id in skip_list or spend_ratio < threshold:
return base_model
return downgrade_paths.get(base_model, base_model)
Loading
Loading