[F098] orchestrator: keep waiting record through a re-park during probe-success resume

resolve_wait deleted the waiting record (in-memory + durable) BEFORE calling
spawn_agent. A re-park in the window between the probe-success clear and the
spawn — the provider's rate limit lifts then immediately re-limits, or a second
provider limit lands — bails spawn with an OFFLINE instance (the parked-provider
short-circuit). Deleting the record first orphaned the agent: with no record
the probe-resume loop can never revive it and the spawn gate bails every tick,
so the agent is lost until the operator intervenes.

Fix: spawn first, then tear down the record only once a container actually
launched (instance.state == ACTIVE). On an OFFLINE bail the record stays so the
next probe-success re-attempts the resume. On a spawn EXCEPTION the record is
torn down + re-raised so the probe loop doesn't keep re-resuming a task that
moved to a different state (e.g. readiness refused -> task auto-blocked) —
matching the pre-fix behavior where the record was deleted before the spawn.
This commit is contained in:
Renn F
2026-06-28 20:56:04 +02:00
parent a4a5756f67
commit 0506176796
2 changed files with 168 additions and 9 deletions
+32 -9
View File
@@ -5380,8 +5380,6 @@ class AgentOrchestrator:
return None
record = self._waiting_records[agent_id]
del self._waiting_records[agent_id]
await self._delete_waiting_record(agent_id)
# Generate resume prompt
resume_prompt = self._generate_resume_prompt(record, resolution)
@@ -5391,13 +5389,38 @@ class AgentOrchestrator:
prior = self._instances.get(agent_id)
prior_git_context = prior.config.git_context if prior and prior.config else None
# Respawn
return await self.spawn_agent(
agent_id=agent_id,
initial_prompt=resume_prompt,
task_id=record.task_id,
git_context=prior_git_context,
)
# Respawn FIRST, then tear down the record only once a container actually
# launched. The old order deleted the record (in-memory + durable) before
# the spawn: a re-park during the resume window — the provider's rate limit
# lifts then immediately re-limits, or a second provider limit lands —
# bails spawn with an OFFLINE instance (the parked-provider short-circuit),
# and deleting the record first orphaned the agent. With no record the
# probe-resume loop can never revive it and the spawn gate bails every
# tick, so the agent is lost until the operator intervenes. Keeping the
# record through a bail lets the next probe-success re-attempt the resume.
try:
instance = await self.spawn_agent(
agent_id=agent_id,
initial_prompt=resume_prompt,
task_id=record.task_id,
git_context=prior_git_context,
)
except Exception:
# Spawn failed (e.g. readiness refused → task auto-blocked). Tear
# down the record so the probe loop doesn't keep re-resuming a task
# that has moved to a different state; the blocked-task path takes
# over. This matches the pre-fix behavior where the record was
# deleted before the spawn attempt.
del self._waiting_records[agent_id]
await self._delete_waiting_record(agent_id)
raise
if instance is None or instance.state == AgentState.OFFLINE:
# Spawn bailed without launching (provider re-parked). Keep the record
# so the probe-resume loop re-attempts on the next clear.
return instance
del self._waiting_records[agent_id]
await self._delete_waiting_record(agent_id)
return instance
def _generate_resume_prompt(
self,
@@ -0,0 +1,136 @@
"""F098: a re-park during probe-success resume must not orphan the agent.
``resolve_wait`` deletes the waiting record (in-memory + durable) and then calls
``spawn_agent`` to respawn the parked agent. If the provider re-parks in the
window between the probe-success clear and the spawn (the rate limit lifts then
immediately re-limits, or a second provider limit lands), ``spawn_agent`` bails
with an OFFLINE instance — the F095 parked-provider short-circuit. The old order
deleted the record BEFORE the spawn, so a bail orphaned the agent: no record
means the probe-resume loop can never revive it and the spawn gate bails every
tick. The record must stay until a container actually launches.
"""
from __future__ import annotations
from datetime import UTC, datetime
from typing import Any
from unittest.mock import AsyncMock
from uuid import uuid4
import pytest
from roboco.models.runtime import AgentInstance, WaitingRecord
from roboco.runtime.orchestrator import (
AgentOrchestrator,
AgentReadinessError,
AgentState,
)
def _make_orchestrator() -> AgentOrchestrator:
orch = AgentOrchestrator.__new__(AgentOrchestrator)
orch._waiting_records = {}
orch._instances = {}
return orch
def _record(agent_id: str = "be-dev-1", provider: str = "anthropic") -> WaitingRecord:
return WaitingRecord(
agent_id=agent_id,
task_id=str(uuid4()),
waiting_for="rate_limit_lifted",
waiting_since=datetime.now(UTC),
context={"provider": provider},
)
def _instance(state: AgentState) -> AgentInstance:
cfg = type("C", (), {"provider_type": "anthropic", "model": "opus"})()
return AgentInstance(agent_id="be-dev-1", state=state, config=cfg)
class TestResolveWaitReparkKeepsRecord:
async def test_record_kept_when_spawn_bails_offline(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A re-park bails spawn with an OFFLINE instance — the record must stay
so the next probe-success re-attempts the resume (not orphan the agent)."""
orch = _make_orchestrator()
orch._waiting_records = {"be-dev-1": _record()}
offline = _instance(AgentState.OFFLINE)
monkeypatch.setattr(orch, "spawn_agent", AsyncMock(return_value=offline))
delete_mock = AsyncMock()
monkeypatch.setattr(orch, "_delete_waiting_record", delete_mock)
result = await orch.resolve_wait(
"be-dev-1", {"reason": "rate_limit_lifted", "provider": "anthropic"}
)
assert result is offline
# Record kept — the probe-resume loop can re-attempt on the next clear.
assert "be-dev-1" in orch._waiting_records
delete_mock.assert_not_awaited()
async def test_record_deleted_when_spawn_launches(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A successful launch (ACTIVE) tears down the record as before."""
orch = _make_orchestrator()
orch._waiting_records = {"be-dev-1": _record()}
active = _instance(AgentState.ACTIVE)
monkeypatch.setattr(orch, "spawn_agent", AsyncMock(return_value=active))
delete_mock = AsyncMock()
monkeypatch.setattr(orch, "_delete_waiting_record", delete_mock)
result = await orch.resolve_wait(
"be-dev-1", {"reason": "rate_limit_lifted", "provider": "anthropic"}
)
assert result is active
assert "be-dev-1" not in orch._waiting_records
delete_mock.assert_awaited_once_with("be-dev-1")
async def test_record_deleted_when_spawn_raises(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A spawn failure (e.g. readiness refused → task auto-blocked) tears
down the record so the probe loop doesn't keep re-resuming a task that
moved to a different state. Matches the pre-fix behavior (the record
was deleted before the spawn attempt)."""
orch = _make_orchestrator()
orch._waiting_records = {"be-dev-1": _record()}
async def _boom(*_a: Any, **_k: Any) -> AgentInstance:
raise AgentReadinessError("not ready")
monkeypatch.setattr(orch, "spawn_agent", _boom)
delete_mock = AsyncMock()
monkeypatch.setattr(orch, "_delete_waiting_record", delete_mock)
with pytest.raises(AgentReadinessError):
await orch.resolve_wait(
"be-dev-1", {"reason": "rate_limit_lifted", "provider": "anthropic"}
)
assert "be-dev-1" not in orch._waiting_records
delete_mock.assert_awaited_once_with("be-dev-1")
async def test_record_kept_across_repeated_offline_bails(
self, monkeypatch: pytest.MonkeyPatch
) -> None:
"""Two probe-success resumes that both bail (provider flaps) must both
find the record still present — the agent is never orphaned across the
flap, and the durable delete is never called."""
orch = _make_orchestrator()
orch._waiting_records = {"be-dev-1": _record()}
offline = _instance(AgentState.OFFLINE)
monkeypatch.setattr(orch, "spawn_agent", AsyncMock(return_value=offline))
delete_mock = AsyncMock()
monkeypatch.setattr(orch, "_delete_waiting_record", delete_mock)
for _ in range(3):
await orch.resolve_wait(
"be-dev-1", {"reason": "rate_limit_lifted", "provider": "anthropic"}
)
assert "be-dev-1" in orch._waiting_records
delete_mock.assert_not_awaited()