From 463b149bdb58ed255e5b5e071a492b46564df852 Mon Sep 17 00:00:00 2001 From: Ahmed Allam Date: Tue, 29 Sep 2026 23:41:58 +0000 Subject: [PATCH] fix(budget): parked agents count as active; park never overwrites a stop active_agents_except (finish_scan, wait_for_message) treats budget_paused as active, so a root cannot finish the scan over a parked child. park_for_budget only transitions a running agent, and the wake back to running happens under the coordinator lock. --- strix/core/agents.py | 66 ++++++++++++++++++++----------- strix/tools/agents_graph/tools.py | 7 +--- tests/test_budget_pause_policy.py | 30 ++++++++++++++ 3 files changed, 75 insertions(+), 28 deletions(-) diff --git a/strix/core/agents.py b/strix/core/agents.py index 095d5fba8..eb1b76bc3 100644 --- a/strix/core/agents.py +++ b/strix/core/agents.py @@ -28,6 +28,10 @@ BudgetPolicy = Literal["stop", "pause"] TERMINAL_STATUSES: frozenset[str] = frozenset({"completed", "stopped", "crashed", "failed"}) +# Agents that still have work to do or to wake up for. A parked agent is mid-task: +# a scan is not finished while one exists, and it can be stopped like any other. +ACTIVE_STATUSES: frozenset[str] = frozenset({"running", "waiting", "budget_paused"}) + # Why an agent parked. The user can message any agent, so this - not the agent's # position in the tree - decides whether waiting is bounded: only an agent waiting # on other agents is re-checked on a timer. @@ -123,9 +127,19 @@ class AgentCoordinator: self._budget_paused = True await self.set_status(agent_id, "budget_paused") - async def park_for_budget(self, agent_id: str) -> None: - """Record that ``agent_id`` parked before an LLM call (pause policy).""" - await self.set_status(agent_id, "budget_paused") + async def park_for_budget(self, agent_id: str) -> bool: + """Park ``agent_id`` before an LLM call (pause policy). + + Only a ``running`` agent parks; returns False when it was stopped in the + meantime, so the stop is not overwritten. + """ + async with self._lock: + if self.statuses.get(agent_id) != "running": + return False + self._set_status_locked(agent_id, "budget_paused") + logger.info("agent.status %s=budget_paused", agent_id) + await self._maybe_snapshot() + return True async def pause_budget(self) -> None: """Operator pause: every agent parks before its next LLM call. @@ -164,22 +178,23 @@ class AgentCoordinator: ``parked_epoch`` is the ``resume_epoch`` the agent read when it decided to park, so a resume that lands between that decision and this wait is not - missed. ``agent_id`` is ``running`` again on return unless it was stopped. + missed. The switch back to ``running`` happens under the lock, so a stop + can never be overwritten by it; on return the agent is ``running`` unless + it was stopped. """ while True: async with self._lock: runtime = self.runtimes.setdefault(agent_id, AgentRuntime()) - if ( - self._budget_stopped - or self._resume_epoch != parked_epoch - or self.statuses.get(agent_id) != "budget_paused" - ): + if self._budget_stopped or self.statuses.get(agent_id) != "budget_paused": + return + if self._resume_epoch != parked_epoch: + self._set_status_locked(agent_id, "running") break wake = runtime.wake wake.clear() await wake.wait() - if not self._budget_stopped and self.statuses.get(agent_id) == "budget_paused": - await self.set_status(agent_id, "running") + logger.info("agent.status %s=running", agent_id) + await self._maybe_snapshot() async def resume_from_budget_pause(self, *, exclude: str | None = None) -> None: """Legacy interactive resume: extend by the original budget and nudge agents.""" @@ -337,20 +352,25 @@ class AgentCoordinator: async with self._lock: if agent_id not in self.statuses: return - self.statuses[agent_id] = status # type: ignore[assignment] - if error is not None: - self.errors[agent_id] = error - elif status == "running": - self.errors.pop(agent_id, None) - if status == "running": - # Running again means a fresh stint that owes its parent its own notice. - self._parent_notified.discard(agent_id) - runtime = self.runtimes.setdefault(agent_id, AgentRuntime()) - runtime.user_wake_required = status in {"failed", "crashed"} - runtime.wake.set() + self._set_status_locked(agent_id, status, error=error) logger.info("agent.status %s=%s", agent_id, status) await self._maybe_snapshot() + def _set_status_locked( + self, agent_id: str, status: Status | str, *, error: str | None = None + ) -> None: + self.statuses[agent_id] = status # type: ignore[assignment] + if error is not None: + self.errors[agent_id] = error + elif status == "running": + self.errors.pop(agent_id, None) + if status == "running": + # Running again means a fresh stint that owes its parent its own notice. + self._parent_notified.discard(agent_id) + runtime = self.runtimes.setdefault(agent_id, AgentRuntime()) + runtime.user_wake_required = status in {"failed", "crashed"} + runtime.wake.set() + async def claim_parent_notice(self, agent_id: str) -> bool: """Reserve the one notice a child owes its parent when it stops running. @@ -541,7 +561,7 @@ class AgentCoordinator: "parent_id": self.parent_of.get(aid), } for aid, status in self.statuses.items() - if aid != agent_id and status in {"running", "waiting"} + if aid != agent_id and status in ACTIVE_STATUSES ] async def graph_snapshot( diff --git a/strix/tools/agents_graph/tools.py b/strix/tools/agents_graph/tools.py index 2807b1c92..051fc1fc5 100644 --- a/strix/tools/agents_graph/tools.py +++ b/strix/tools/agents_graph/tools.py @@ -12,16 +12,13 @@ from typing import Any, Literal, get_args from agents import RunContextWrapper, function_tool -from strix.core.agents import Status, coordinator_from_context +from strix.core.agents import ACTIVE_STATUSES, Status, coordinator_from_context from strix.core.execution import notify_parent_on_terminal from strix.core.hooks import LLM_TURN_KEY from strix.report.state import get_global_report_state from strix.skills import validate_requested_skills -_ACTIVE_STATUSES: frozenset[str] = frozenset({"running", "waiting", "budget_paused"}) - - logger = logging.getLogger(__name__) @@ -810,7 +807,7 @@ async def stop_agent( ) current_status = statuses[target_agent_id] - if current_status not in _ACTIVE_STATUSES: + if current_status not in ACTIVE_STATUSES: return json.dumps( { "success": False, diff --git a/tests/test_budget_pause_policy.py b/tests/test_budget_pause_policy.py index 302b23baf..ad7b3af02 100644 --- a/tests/test_budget_pause_policy.py +++ b/tests/test_budget_pause_policy.py @@ -396,6 +396,36 @@ async def test_parked_wait_returns_on_stop_signals() -> None: assert coordinator.statuses["a"] == "budget_paused" +@pytest.mark.asyncio +async def test_park_never_overwrites_a_stop_that_landed_first() -> None: + coordinator = AgentCoordinator() + coordinator.set_budget_policy("pause") + await coordinator.register("a", "strix", parent_id=None) + + parked_epoch = coordinator.resume_epoch + await coordinator.request_stop("a") + assert await coordinator.park_for_budget("a") is False + assert coordinator.statuses["a"] == "stopped" + + await asyncio.wait_for( + coordinator.wait_for_budget_resume("a", parked_epoch=parked_epoch), timeout=1.0 + ) + assert coordinator.statuses["a"] == "stopped" + + +@pytest.mark.asyncio +async def test_parked_children_keep_the_scan_open() -> None: + coordinator = AgentCoordinator() + coordinator.set_budget_policy("pause") + await coordinator.register("root", "strix", parent_id=None) + await coordinator.register("child", "recon", parent_id="root") + assert await coordinator.park_for_budget("child") is True + + active = await coordinator.active_agents_except("root") + assert [a["agent_id"] for a in active] == ["child"] + assert active[0]["status"] == "budget_paused" + + @pytest.mark.asyncio async def test_resume_budget_replaces_the_limit_and_validates_it() -> None: hooks = ReportUsageHooks(model="m", max_budget_usd=10.0, budget_policy="pause")