mirror of
https://github.com/usestrix/strix.git
synced 2026-10-01 02:03:55 +00:00
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.
This commit is contained in:
parent
d355838ea0
commit
463b149bdb
3 changed files with 75 additions and 28 deletions
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue