diff --git a/argus/apps/cli/_parser.py b/argus/apps/cli/_parser.py index e0e439af4..823c0abd1 100644 --- a/argus/apps/cli/_parser.py +++ b/argus/apps/cli/_parser.py @@ -53,6 +53,21 @@ def format_help(self) -> str: return _PUBLIC_HELP +def _mission_width(value: str) -> int | str: + text = str(value or "").strip().lower() + if text == "auto": + return "auto" + try: + width = int(text) + except ValueError as exc: + raise argparse.ArgumentTypeError( + "mission width must be a whole number or 'auto'" + ) from exc + if width < 0: + raise argparse.ArgumentTypeError("mission width must be zero or more") + return width + + def _tcp_port(value: str) -> int: """Reject a port the kernel can never bind, before anything offers a URL. @@ -244,10 +259,11 @@ def build_parser() -> argparse.ArgumentParser: ) daemon_grp.add_argument( "--mission-width", - type=int, + type=_mission_width, default=2, help="concurrent mission workers: 0 pauses, 1 is serial, N enables " - "path-disjoint parallel Planner tasks (default: 2)", + "path-disjoint parallel Planner tasks, auto is one per GPU on " + "this machine (at least 2, at most 4) (default: 2)", ) cockpit_grp = parser.add_argument_group("cockpit") diff --git a/argus/daemon/config.py b/argus/daemon/config.py index 761b05046..e7d088e60 100644 --- a/argus/daemon/config.py +++ b/argus/daemon/config.py @@ -21,6 +21,24 @@ def _default_subagent_family_failure_window_hours() -> float: return LifeSupervisorConfig.subagent_family_failure_window_hours +AUTO_MISSION_WIDTH_MIN = 2 +AUTO_MISSION_WIDTH_MAX = 4 + + +def resolve_auto_mission_width() -> int: + """The width ``--mission-width auto`` stands for on this machine. + + One mission worker per GPU, so comparison arms that each hold a card can + run side by side, with a floor of two (a campaign always has work that + needs no GPU) and a ceiling of four (more workers share one provider + session budget and one operator's attention). + """ + from ..tools.gpu_lease import gpu_snapshot + + gpus = len(gpu_snapshot()) + return max(AUTO_MISSION_WIDTH_MIN, min(AUTO_MISSION_WIDTH_MAX, gpus)) + + @dataclass class LifeWorkerConfig: """How the worker drains the backlog. @@ -41,7 +59,9 @@ class LifeWorkerConfig: engineer_reasoning_effort: str = "xhigh" reviewer_reasoning_effort: str = "high" global_daily_cap_usd: float = 0.0 - mission_width: int = 2 + # An integer, or "auto": one mission worker per GPU on this machine, + # at least two and at most four, resolved when the config is built. + mission_width: int | str = 2 planner_task_iteration_max_cycles: int = 0 # See LifeSupervisorConfig.subagent_family_failure_streak_limit / # ..._window_hours (life/supervisor/_config.py) for the circuit breaker @@ -82,6 +102,8 @@ class LifeWorkerConfig: last_spawn_error: str = field(default="", init=False, repr=False, compare=False) def __post_init__(self) -> None: + if isinstance(self.mission_width, str) and self.mission_width.strip().lower() == "auto": + self.mission_width = resolve_auto_mission_width() self.mission_width = int(self.mission_width) if self.mission_width < 0: raise ValueError("mission_width must be zero or a positive integer") diff --git a/argus/life/memory.py b/argus/life/memory.py index c41049e33..1d40e3466 100644 --- a/argus/life/memory.py +++ b/argus/life/memory.py @@ -891,6 +891,67 @@ class PlanRevisionResult: added_ids: tuple[str, ...] +# --------------------------------------------------------------------------- +# GPUs against the backlog +# --------------------------------------------------------------------------- + +_GPU_PROBE_TTL_S = 10.0 +_GPU_PROBE_LOCK = threading.Lock() +_GPU_PROBE_CACHE: tuple[float, tuple[int, int] | None] | None = None +_GPU_ACTIVE_STATUSES = frozenset({"running", "paused_external_work"}) + + +def _gpu_capacity() -> tuple[int, int] | None: + """``(total, busy)`` GPUs on this machine, or ``None`` when none is visible. + + Busy uses the thresholds the GPU lease uses for "not idle": over 5 % + utilisation or over 2 GiB in use. The probe is a subprocess and this is + consulted on every claim, so its answer is kept for ten seconds. + """ + global _GPU_PROBE_CACHE + now = time.monotonic() + with _GPU_PROBE_LOCK: + cached = _GPU_PROBE_CACHE + if cached is not None and now - cached[0] < _GPU_PROBE_TTL_S: + return cached[1] + from ..tools.gpu_lease import gpu_snapshot + + snapshot = gpu_snapshot() + value: tuple[int, int] | None = None + if snapshot: + busy = sum( + 1 for gpu in snapshot + if gpu["util_pct"] > 5 or gpu["mem_used_mib"] > 2048 + ) + value = (len(snapshot), busy) + with _GPU_PROBE_LOCK: + _GPU_PROBE_CACHE = (now, value) + return value + + +def gpu_reservation(items: Iterable[Any]) -> int: + """GPUs held by running tasks and by tasks parked on their own jobs.""" + return sum( + max(0, int(getattr(item, "gpu_count", 0) or 0)) + for item in items + if getattr(item, "status", "") in _GPU_ACTIVE_STATUSES + ) + + +def gpu_capacity_summary(items: Iterable[Any]) -> dict[str, int] | None: + capacity = _gpu_capacity() + if capacity is None: + return None + total, busy = capacity + reserved = gpu_reservation(list(items)) + return { + "total": total, + "busy": busy, + "reserved": reserved, + "free": max(0, total - max(busy, reserved)), + } + + @dataclass class BacklogItem: id: str @@ -995,6 +1056,9 @@ class BacklogItem: # workers. The primary worker remains able to execute every backlog item. parallel_safe: bool = False owns_paths: list[str] = field(default_factory=list) + # GPUs this task holds while it runs. A task is claimed only when that + # many are free: not busy now and not reserved by another active task. + gpu_count: int = 0 outcome: dict[str, Any] = field(default_factory=dict) # Optional durable return receipt; kept separate from public outcome dimensions. mission_result: dict[str, Any] | None = None @@ -1027,6 +1091,7 @@ def new( execution_workdir: str = "", parallel_safe: bool = False, owns_paths: list[str] | None = None, + gpu_count: int = 0, acceptance_check: str = "", plan_hypothesis: str = "", goal_contribution: str = "", @@ -1069,6 +1134,7 @@ def new( for path in (owns_paths or []) if str(path).strip() ], + gpu_count=max(0, int(gpu_count or 0)), acceptance_check=str(acceptance_check or "").strip(), plan_hypothesis=str(plan_hypothesis or "").strip(), goal_contribution=str(goal_contribution or "").strip(), @@ -1173,6 +1239,7 @@ def from_jsonable(cls, row: dict[str, Any]) -> "BacklogItem": for path in (row.get("owns_paths", []) or []) if str(path).strip() ], + gpu_count=max(0, int(row.get("gpu_count", 0) or 0)), outcome=( {str(key): value for key, value in row.get("outcome", {}).items()} if isinstance(row.get("outcome"), dict) @@ -1401,6 +1468,27 @@ def _done_ids(items: Iterable[BacklogItem]) -> set[str]: """ return {it.id for it in items if it.status == "done"} + @staticmethod + def _fit_gpus( + candidates: list[BacklogItem], items: Iterable[BacklogItem], + ) -> list[BacklogItem]: + """Drop candidates whose GPUs are not free right now. + + Free means total minus the larger of "busy now" and "reserved by + active tasks", so a running task's cards count once whether or not + its job has started. With no GPU visible nothing is gated: the + Engineer then reports the missing hardware itself. + """ + if not any(getattr(item, "gpu_count", 0) > 0 for item in candidates): + return candidates + summary = gpu_capacity_summary(items) + if summary is None: + return candidates + return [ + item for item in candidates + if getattr(item, "gpu_count", 0) <= summary["free"] + ] + @staticmethod def _is_ready(item: BacklogItem, done: set[str]) -> bool: """A pending item is ready iff every dep is in ``done``. @@ -2704,10 +2792,16 @@ def ready(self) -> list[BacklogItem]: items = self._load() history = self._dependency_history(items) done = self._done_ids([*history, *items]) - out = [it for it in items if self._is_ready(it, done)] + out = self._fit_gpus([it for it in items if self._is_ready(it, done)], items) out.sort(key=lambda it: (it.priority, it.ts)) return out + def gpu_summary(self) -> dict[str, int] | None: + """This machine's GPUs against the backlog: total, busy, reserved, free.""" + with self._locked(): + items = self._load() + return gpu_capacity_summary(items) + def next_pending( self, *, @@ -2730,7 +2824,9 @@ def next_pending( history = self._dependency_history(items) changed = self._cascade_blocked(items, history=history) done = self._done_ids([*history, *items]) - ready = [item for item in items if self._is_ready(item, done)] + ready = self._fit_gpus( + [item for item in items if self._is_ready(item, done)], items, + ) if parallel_only or ( respect_running and any( diff --git a/argus/life/supervisor/_planner_orchestration.py b/argus/life/supervisor/_planner_orchestration.py index 0f439dfe6..d674b9a1e 100644 --- a/argus/life/supervisor/_planner_orchestration.py +++ b/argus/life/supervisor/_planner_orchestration.py @@ -226,6 +226,18 @@ def _planner_current_reality_note(self) -> str: f"- mission_slots: {mission_slots} total; " f"{len(running_rows)} running; {free_slots} free" ) + gpu_summary = self.memory.backlog.gpu_summary() + if gpu_summary is not None: + slot_lines.append( + f"- gpus: {gpu_summary['total']} on this machine; " + f"{gpu_summary['busy']} busy now; " + f"{gpu_summary['reserved']} reserved by active tasks; " + f"{gpu_summary['free']} claimable. TASK_GPUS= is how many " + "GPUs a task holds while it runs; it is claimed only when that " + "many are free. Arms that can run side by side (model sizes, " + "training methods, seeds) belong in separate parallel-safe tasks " + "with their own TASK_GPUS, not in one task that runs them in turn." + ) # The claim gate also refuses everything while a paused external # job declares no owned paths, so those rows block a "free" slot # exactly like an unowned running mission does. @@ -270,7 +282,8 @@ def _active_item_line(item: Any) -> str: safe = "true" if getattr(item, "parallel_safe", False) else "false" return ( f"{base}; parallel_safe={safe}; " - f"owns_paths=[{', '.join(owns)}]" + f"owns_paths=[{', '.join(owns)}]; " + f"gpus={int(getattr(item, 'gpu_count', 0) or 0)}" ) # A subagent event wait is bound by matching the Planner's own words diff --git a/argus/life/supervisor/_planning_cycle_enqueue.py b/argus/life/supervisor/_planning_cycle_enqueue.py index 03be6af9f..66ae2f8ed 100644 --- a/argus/life/supervisor/_planning_cycle_enqueue.py +++ b/argus/life/supervisor/_planning_cycle_enqueue.py @@ -603,6 +603,7 @@ def _pc_build_pending_items(self, state: _PlanCycleState) -> Any | None: allow_skill_changes=False, parallel_safe=canonical_parallel, owns_paths=canonical_owns_paths, + gpu_count=max(0, int(getattr(task, "gpu_count", 0) or 0)), ) from ...skills.stage_machine import current_stage from ...skills.vertical_select import resolve_vertical @@ -1002,6 +1003,7 @@ def _pc_build_pending_items(self, state: _PlanCycleState) -> Any | None: ), parallel_safe=bool(getattr(task, "parallel_safe", False)), owns_paths=list(getattr(task, "owns_paths", []) or []), + gpu_count=max(0, int(getattr(task, "gpu_count", 0) or 0)), non_goals=list(getattr(task, "non_goals", []) or []), original_objective=str( getattr(self.config, "continuous_objective", "") or "" diff --git a/argus/planner/planner.py b/argus/planner/planner.py index 0be21a576..836bee9c5 100644 --- a/argus/planner/planner.py +++ b/argus/planner/planner.py @@ -136,6 +136,8 @@ class TaskSpec: allow_skill_changes: bool = False parallel_safe: bool = False owns_paths: list[str] = field(default_factory=list) + # GPUs the task holds while it runs; claimed only when that many are free. + gpu_count: int = 0 # Mission-level role selected by Planner. Empty inherits the campaign # vertical chosen by Manager at the front door. vertical: str = "" @@ -635,6 +637,7 @@ def _repair_no_task_verdict( "SCOPE", "PARALLEL_SAFE", "OWNS_PATHS", + "GPUS", "VERTICAL", "REQUIRE_INDEPENDENT_REVIEW", ) @@ -971,6 +974,12 @@ def boolean(source: Mapping[str, Any], name: str) -> bool: raise TypeError(f"{name} must be true or false") return value + def integer(source: Mapping[str, Any], name: str) -> int: + value = source.get(name, 0) + if isinstance(value, bool) or not isinstance(value, int) or value < 0: + raise TypeError(f"{name} must be a non-negative integer") + return value + def review_boolean(source: Mapping[str, Any], name: str) -> bool: # Mirrors the bounded-DAG validation contract: a structured boolean or # the literal strings "true"/"false"; anything else is a metadata error. @@ -1107,6 +1116,7 @@ def review_boolean(source: Mapping[str, Any], name: str) -> bool: owns_paths=items( raw_task.get("owns_paths", []), "owns_paths" ), + gpu_count=integer(raw_task, "gpu_count"), vertical=text(raw_task, "vertical").strip(), ) ) @@ -1310,6 +1320,7 @@ def _planner_verdict_from_fields( for path in row.get("TASK_OWNS_PATHS", "").split("|") if path.strip() ], + gpu_count=max(0, _key_value_int(row.get("TASK_GPUS", ""))), vertical=row.get("TASK_VERTICAL", "").strip(), ) ) diff --git a/argus/roles/prompts/planner.py b/argus/roles/prompts/planner.py index 1135a4f68..76a0fe824 100644 --- a/argus/roles/prompts/planner.py +++ b/argus/roles/prompts/planner.py @@ -82,7 +82,8 @@ Optional `ADVANCE_TO_STAGE` must be Host-valid; omit to hold. `TASK_SCOPE` defaults to `bounded`. Also: `TASK_KEY`/`TASK_DEPS`, `TASK_HYPOTHESIS`, `TASK_GOAL_CONTRIBUTION`, `TASK_EXPECTED_REGRESSIONS`, `TASK_DECISION_RULE`, - `TASK_ACCEPTANCE_CHECK`, `TASK_PARALLEL_SAFE`, `TASK_OWNS_PATHS`, and `TASK_VERTICAL`. + `TASK_ACCEPTANCE_CHECK`, `TASK_PARALLEL_SAFE`, `TASK_OWNS_PATHS`, `TASK_GPUS`, + and `TASK_VERTICAL`. - Tasks co-run only when each sets `TASK_PARALLEL_SAFE=true` with disjoint, literal, relative `TASK_OWNS_PATHS` (no wildcards); stage-closing and framework-maintenance work runs alone. The digest shows slots and ownership. diff --git a/tests/apps/test_cli_parser.py b/tests/apps/test_cli_parser.py index eabb9182b..9f15cea73 100644 --- a/tests/apps/test_cli_parser.py +++ b/tests/apps/test_cli_parser.py @@ -867,3 +867,14 @@ def test_a_bindable_web_port_is_accepted(value: str) -> None: """Zero is legal: it asks the kernel for any free port.""" args = build_parser().parse_args(["--web", "--web-port", value]) assert args.web_port == int(value) + + +def test_mission_width_accepts_a_count_or_auto() -> None: + parser = build_parser() + assert parser.parse_args(["--daemon", "--mission-width", "auto"]).mission_width == "auto" + assert parser.parse_args(["--daemon", "--mission-width", "3"]).mission_width == 3 + assert parser.parse_args(["--daemon"]).mission_width == 2 + with pytest.raises(SystemExit): + parser.parse_args(["--daemon", "--mission-width", "-1"]) + with pytest.raises(SystemExit): + parser.parse_args(["--daemon", "--mission-width", "many"]) diff --git a/tests/daemon/test_protocol.py b/tests/daemon/test_protocol.py index 43bff5c7e..1d5aead6a 100644 --- a/tests/daemon/test_protocol.py +++ b/tests/daemon/test_protocol.py @@ -283,3 +283,18 @@ def test_daemon_source_ownership_requires_same_installation( assert daemon_runtime_owned_by_current_source(owned) is True assert daemon_runtime_owned_by_current_source(foreign) is False + + +def test_mission_width_auto_is_one_worker_per_gpu_within_bounds( + tmp_path: Path, monkeypatch, +) -> None: + from argus.tools import gpu_lease + + for gpus, expected in ((0, 2), (1, 2), (4, 4), (8, 4)): + monkeypatch.setattr( + gpu_lease, "gpu_snapshot", + lambda gpus=gpus: [{"index": i, "mem_used_mib": 1, "mem_total_mib": 46068, "util_pct": 0} for i in range(gpus)], + ) + config = LifeWorkerConfig(life_dir=tmp_path, backend="memory", mission_width="auto") + assert config.mission_width == expected + assert config_from_payload(config_payload(config)).mission_width == expected diff --git a/tests/life/test_memory.py b/tests/life/test_memory.py index da8e04404..b0eb9e85a 100644 --- a/tests/life/test_memory.py +++ b/tests/life/test_memory.py @@ -1049,3 +1049,61 @@ def test_operator_reply_continuation_starts_streak_tracked( stored = next(row for row in b.all() if row.id == continuation.id) assert stored.replan_streak_tracked is True assert stored.consecutive_replans == 0 + + +# --------------------------------------------------------------------------- +# GPUs against the backlog +# --------------------------------------------------------------------------- + + +def _gpu_backlog(tmp_path, monkeypatch, capacity): + from argus.life import memory as memory_module + + monkeypatch.setattr(memory_module, "_gpu_capacity", lambda: capacity) + memory = LifeMemory.open(tmp_path / "life") + return memory.backlog + + +def test_a_task_that_holds_gpus_is_claimable_only_while_that_many_are_free(tmp_path, monkeypatch): + # Four cards; one is busy with someone else's process right now. + backlog = _gpu_backlog(tmp_path, monkeypatch, (4, 1)) + training = backlog.add(BacklogItem.new( + title="train the 27B arm", objective="train", gpu_count=2, + parallel_safe=True, owns_paths=["runs/27b"], + )) + assert [item.id for item in backlog.ready()] == [training.id] + claimed = backlog.claim_next() + assert claimed is not None and claimed.id == training.id + # Reserved 2, busy 1: two cards remain for other work. + assert backlog.gpu_summary() == {"total": 4, "busy": 1, "reserved": 2, "free": 2} + fits = backlog.add(BacklogItem.new( + title="train the small arm", objective="train", gpu_count=2, + parallel_safe=True, owns_paths=["runs/small"], + )) + too_big = backlog.add(BacklogItem.new( + title="train the huge arm", objective="train", gpu_count=3, + parallel_safe=True, owns_paths=["runs/huge"], + )) + no_gpu = backlog.add(BacklogItem.new(title="write the related work", objective="write")) + assert {item.id for item in backlog.ready()} == {fits.id, no_gpu.id} + assert backlog.next_pending(parallel_only=True).id == fits.id + # Once the running job actually occupies its cards, "busy" covers them and + # nothing is counted twice. + monkeypatch.setattr("argus.life.memory._gpu_capacity", lambda: (4, 3)) + assert backlog.gpu_summary()["free"] == 1 + assert {item.id for item in backlog.ready()} == {no_gpu.id} + assert too_big.id not in {item.id for item in backlog.ready()} + + +def test_without_a_visible_gpu_nothing_is_gated(tmp_path, monkeypatch): + backlog = _gpu_backlog(tmp_path, monkeypatch, None) + item = backlog.add(BacklogItem.new(title="train", objective="train", gpu_count=8)) + assert backlog.gpu_summary() is None + assert [ready.id for ready in backlog.ready()] == [item.id] + + +def test_gpu_count_survives_the_journal(tmp_path, monkeypatch): + backlog = _gpu_backlog(tmp_path, monkeypatch, None) + item = backlog.add(BacklogItem.new(title="train", objective="train", gpu_count=2)) + reloaded = next(row for row in LifeMemory.open(tmp_path / "life").backlog.all() if row.id == item.id) + assert reloaded.gpu_count == 2 diff --git a/tests/planner/test_planner.py b/tests/planner/test_planner.py index 456fc9431..f5132df36 100644 --- a/tests/planner/test_planner.py +++ b/tests/planner/test_planner.py @@ -1492,3 +1492,25 @@ def test_parse_numbered_planner_objective_keeps_its_following_lines() -> None: "## Claim\nProduce the theorem.\n- one lemma per file" ) assert verdict.new_tasks[0].acceptance_check == "Run the verifier." + + +def test_parse_planner_reads_the_gpus_a_task_holds() -> None: + verdict = parse_planner_text( + "\n".join([ + "PROJECT_DONE=false", + "REASON=two arms can train side by side", + "TASK_KEY=arm-27b", + "TASK_DEPS=", + "TASK_TITLE=Train the 27B arm", + "TASK_OBJECTIVE=Train it.", + "TASK_PARALLEL_SAFE=true", + "TASK_OWNS_PATHS=runs/27b", + "TASK_GPUS=2", + "TASK_KEY=writeup", + "TASK_DEPS=", + "TASK_TITLE=Draft the method section", + "TASK_OBJECTIVE=Write it.", + "TASK_GPUS=not-a-number", + ]) + ) + assert [task.gpu_count for task in verdict.new_tasks] == [2, 0]