mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-30 21:44:50 +02:00
[eric] agents: transport heals respawn the CLI on the same transcript instead of rebuilding it (ENG-382)
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01En8dRGsJPLrJCQBEkTH4Mp
This commit is contained in:
co-authored by
Claude Fable 5
parent
a4be8c469f
commit
42e4f9116b
@@ -235,7 +235,8 @@ class AgentManager(SessionLifecycle, SessionHistory, SessionPersistence, Messagi
|
||||
turn = TurnState()
|
||||
p_stderr_buffer: List[str] = []
|
||||
# Read BEFORE build_agent_options consumes these flags: a fresh-session/fork request must force the persistent client to respawn (same branch id would otherwise fingerprint-match a client still holding the old transcript).
|
||||
p_force_respawn = bool(session.needs_fresh_session or session.needs_fork or fork_session)
|
||||
p_force_respawn = bool(session.needs_fresh_session or session.needs_fork or fork_session or session.needs_respawn)
|
||||
session.needs_respawn = False
|
||||
try:
|
||||
logger.info(f"[SPAWN-PHASE] run-loop start session={session_id[:8]} t={time.monotonic():.3f}")
|
||||
(options, options_kwargs, prompt_content, p_stderr_buffer,
|
||||
|
||||
@@ -132,6 +132,8 @@ class AgentSession(BaseModel):
|
||||
needs_fork: bool = False
|
||||
# Stronger than needs_fork: drop resume= and replay history into a fresh sdk_session_id; fork_session alone won't re-read mcp_servers.
|
||||
needs_fresh_session: bool = False
|
||||
# A new CLI process that RESUMES the same transcript (dead transport, stale token, core sidecar never connected); unlike needs_fresh_session nothing is rebuilt, so no history is ever re-authored as text (ENG-382).
|
||||
needs_respawn: bool = False
|
||||
# Auto-continue: agent loop dispatches a hidden turn at end-of-loop using pending_continuation_prompt. Race-free vs background tasks.
|
||||
pending_continuation: bool = False
|
||||
pending_continuation_prompt: Optional[str] = None
|
||||
@@ -147,10 +149,10 @@ class AgentSession(BaseModel):
|
||||
empty_finish_progress_mark: int = 0
|
||||
# One honest "stopped without a report" line per exhausted nudge budget; resets with the budget.
|
||||
empty_finish_surfaced: bool = False
|
||||
# What a fresh CLI session may carry as history: "full" (asks, the model's own replies as gists, tool trail), "minimal" (asks + tool calls, zero model text), "none". Ratchets DOWN when a provider policy filter blocks a recap-bearing turn and never back up: Anthropic's anti-distillation classifier reads a replay of the model's own outputs as "duplicating model outputs" (Alex, 57 blocks in 4 days, every one at spawn).
|
||||
history_prefix_mode: Literal["full", "minimal", "none"] = "full"
|
||||
# What a fresh CLI session may carry as history: "minimal" (the user's asks, the tool trail, a model-written summary of the dropped span; never the model's own replies verbatim) or "none". Ratchets to "none" when a provider policy filter blocks a recap-bearing turn and never back up: on the subscription lane Anthropic's anti-distillation classifier blocked 192 of our recap turns in 14 days (0 on API keys), reading replayed model text as "duplicating model outputs".
|
||||
history_prefix_mode: Literal["minimal", "none"] = "minimal"
|
||||
# What the LAST spawned turn actually carried, so a block can tell a recap-caused refusal from a plain one.
|
||||
history_prefix_sent: Literal["full", "minimal", "none"] = "none"
|
||||
history_prefix_sent: Literal["minimal", "none"] = "none"
|
||||
# Consecutive dirty deaths this session was MID-TURN for; the crash auto-resume breaker (hermes #30719 pairing: auto-resume must never outrun its circuit breaker).
|
||||
crash_interrupt_count: int = 0
|
||||
# Outage rounds spent on this ask: the in-turn ladder covers only 335s, and the work is checkpointed, so a longer drop is waited out rather than ending the task.
|
||||
|
||||
@@ -59,9 +59,9 @@ def arm_reconnect_resume(session: AgentSession, retry_after_s: Optional[int] = N
|
||||
|
||||
session.reconnect_attempts = attempts + 1
|
||||
session.awaiting_reconnect = True
|
||||
# Only a dead transport leaves the CLI holding a corpse; a 429 is a healthy pipe carrying a NO, and respawning for that spends a process to be told the same thing.
|
||||
# Only a dead transport leaves the CLI holding a corpse; a 429 is a healthy pipe carrying a NO, and respawning for that spends a process to be told the same thing. The new process resumes the same transcript.
|
||||
if connection_lost:
|
||||
session.needs_fresh_session = True
|
||||
session.needs_respawn = True
|
||||
session.pending_continuation = True
|
||||
session.pending_continuation_prompt = RECONNECT_PROMPT
|
||||
session.pending_continuation_delay_s = delay
|
||||
|
||||
@@ -21,8 +21,9 @@ AUTH_RETRY_PROMPT = (
|
||||
|
||||
@typechecked
|
||||
def try_auth_self_heal(session: AgentSession, delay_s: int = 0) -> bool:
|
||||
"""Queue the one hidden retry on a fresh CLI. False = budget spent or a continuation is
|
||||
already pending, and the caller should show the honest banner instead.
|
||||
"""Queue the one hidden retry on a new CLI process that resumes the same transcript (the
|
||||
process is what holds the stale token; the conversation is fine). False = budget spent or a
|
||||
continuation is already pending, and the caller should show the honest banner instead.
|
||||
|
||||
delay_s: codex tokens ROTATE on a 1-2 minute cadence; an instant retry lands inside the same
|
||||
rotation window, burns the one-shot budget, and the user then gets a banner for a condition
|
||||
@@ -31,7 +32,7 @@ def try_auth_self_heal(session: AgentSession, delay_s: int = 0) -> bool:
|
||||
if session.auth_retry_used or session.pending_continuation:
|
||||
return False
|
||||
session.auth_retry_used = True
|
||||
session.needs_fresh_session = True
|
||||
session.needs_respawn = True
|
||||
session.pending_continuation = True
|
||||
session.pending_continuation_prompt = AUTH_RETRY_PROMPT
|
||||
session.pending_continuation_delay_s = max(0, delay_s)
|
||||
|
||||
@@ -55,19 +55,18 @@ def core_mcp_failed_to_connect(init_payload: object) -> bool:
|
||||
|
||||
@typechecked
|
||||
def note_core_mcp_health(session: object, session_id: str, init_payload: object) -> bool:
|
||||
"""Arm a fresh-session rebuild when the core server did not connect. The next turn respawns the
|
||||
CLI (the machinery ENG-258 already uses for unclassified failures), which is the only cure:
|
||||
MCP registration happens once at connect, so a toolless session stays toolless forever.
|
||||
Returns whether it armed."""
|
||||
"""Arm a CLI respawn when the core server did not connect. MCP registration happens once at
|
||||
connect, so a toolless session stays toolless forever; the new process resumes the same
|
||||
transcript (nothing about the conversation changed, so nothing is rebuilt). Returns whether it armed."""
|
||||
if not core_mcp_failed_to_connect(init_payload):
|
||||
return False
|
||||
status = core_mcp_status(init_payload)
|
||||
logger.warning(
|
||||
f"Agent {session_id}: core MCP server reported '{status}', not connected; "
|
||||
"arming a fresh CLI session so the agent is not left without its tools"
|
||||
"arming a CLI respawn so the agent is not left without its tools"
|
||||
)
|
||||
try:
|
||||
session.needs_fresh_session = True # type: ignore[attr-defined]
|
||||
session.needs_respawn = True # type: ignore[attr-defined]
|
||||
except Exception:
|
||||
return False
|
||||
try:
|
||||
|
||||
@@ -57,7 +57,7 @@ async def test_first_token_expiry_heals_silently(monkeypatch):
|
||||
assert not any(m.role == "assistant" for m in session.messages)
|
||||
assert "agent:auth_error" not in events
|
||||
assert session.auth_retry_used is True
|
||||
assert session.needs_fresh_session is True, "the fresh CLI is what drops the stale token"
|
||||
assert session.needs_respawn is True, "a new CLI process is what drops the stale token; the transcript stays"
|
||||
assert session.pending_continuation is True and session.pending_continuation_prompt
|
||||
assert turn.number == 1, "healing must not skip the turn bookkeeping"
|
||||
|
||||
|
||||
@@ -27,7 +27,7 @@ def test_delay_lands_on_the_session():
|
||||
s = p_session()
|
||||
assert try_auth_self_heal(s, delay_s=75) is True
|
||||
assert s.pending_continuation_delay_s == 75
|
||||
assert s.pending_continuation is True and s.needs_fresh_session is True
|
||||
assert s.pending_continuation is True and s.needs_respawn is True
|
||||
|
||||
|
||||
def test_budget_is_one_per_ask():
|
||||
|
||||
@@ -71,19 +71,19 @@ def test_a_malformed_entry_is_ignored_rather_than_read_as_broken():
|
||||
|
||||
|
||||
class P_Session:
|
||||
needs_fresh_session = False
|
||||
needs_respawn = False
|
||||
|
||||
|
||||
def test_a_failed_connect_arms_a_fresh_cli_session():
|
||||
s = P_Session()
|
||||
assert note_core_mcp_health(s, "sess-1", p_init([{"name": "openswarm-core", "status": "failed"}])) is True
|
||||
assert s.needs_fresh_session is True, "a toolless session stays toolless until the CLI respawns"
|
||||
assert s.needs_respawn is True, "a toolless session stays toolless until the CLI respawns"
|
||||
|
||||
|
||||
def test_a_healthy_connect_changes_nothing():
|
||||
s = P_Session()
|
||||
assert note_core_mcp_health(s, "sess-2", p_init([{"name": "openswarm-core", "status": "connected"}])) is False
|
||||
assert s.needs_fresh_session is False
|
||||
assert s.needs_respawn is False
|
||||
|
||||
|
||||
def test_the_turn_runner_consults_this_on_init():
|
||||
|
||||
@@ -122,14 +122,14 @@ def test_auth_failure_self_heals_once_then_stops_respawning(monkeypatch):
|
||||
# rule that still matters is that it happens ONCE; a credential that fails twice is genuinely
|
||||
# dead, and respawning forever would just hide it behind an endless retry.
|
||||
session, _ = p_drive_error(monkeypatch, Exception("401 invalid authentication credentials"))
|
||||
assert session.needs_fresh_session is True, "one rebuild is the heal"
|
||||
assert session.needs_respawn is True, "one respawn on the same transcript is the heal"
|
||||
assert session.auth_retry_used is True
|
||||
assert not [m for m in session.messages if m.role == "system"], "no card on the first expiry"
|
||||
|
||||
session.needs_fresh_session = False
|
||||
session.needs_respawn = False
|
||||
session.pending_continuation = False
|
||||
p_drive_error(monkeypatch, Exception("401 invalid authentication credentials"), session=session)
|
||||
assert session.needs_fresh_session is False, "the budget is spent; stop respawning"
|
||||
assert session.needs_respawn is False, "the budget is spent; stop respawning"
|
||||
assert [m for m in session.messages if m.role == "system"], "the second failure is honest"
|
||||
|
||||
|
||||
|
||||
@@ -46,7 +46,7 @@ def test_an_outage_parks_the_turn_instead_of_ending_it(monkeypatch):
|
||||
assert session.pending_continuation is True, "the work is queued to continue"
|
||||
assert session.pending_continuation_delay_s == RECONNECT_BACKOFFS[0]
|
||||
assert session.awaiting_reconnect is True
|
||||
assert session.needs_fresh_session is True, "the CLI died with the outage; resume on a fresh one"
|
||||
assert session.needs_respawn is True, "the CLI died with the outage; a new process resumes the transcript"
|
||||
assert not [m for m in session.messages if m.role == "system"], "no card: nothing is over yet"
|
||||
assert "agent:reconnect_wait" in [e for e, _ in events]
|
||||
assert "agent:rate_limited" not in [e for e, _ in events]
|
||||
@@ -204,11 +204,11 @@ def test_a_dead_socket_respawns_the_cli_but_a_429_does_not(monkeypatch):
|
||||
respawning for that spends a whole process to be told the same thing (caught by the existing
|
||||
test_rate_limit_does_not_respawn_the_cli when this shipped ungated)."""
|
||||
dead, _ = p_drive(monkeypatch, ConnectionError("Connection reset by peer"))
|
||||
assert dead.needs_fresh_session is True
|
||||
assert dead.needs_respawn is True
|
||||
assert dead.awaiting_reconnect is True
|
||||
|
||||
throttled, _ = p_drive(monkeypatch, Exception("429 rate_limit_error: overloaded"))
|
||||
assert throttled.needs_fresh_session is False, "a refusal is not a broken pipe"
|
||||
assert throttled.needs_respawn is False, "a refusal is not a broken pipe"
|
||||
assert throttled.awaiting_reconnect is True, "but it is still worth waiting out"
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user