Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions supervisor/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -380,6 +380,7 @@ def cmd_run_foreground(args):
mode=config.pause_handling_mode,
max_auto_interventions=config.max_auto_interventions,
),
runtime_recovery_policy=config.runtime_recovery_policy(),
)

print(f"[DEBUG MODE] Foreground controller — for debugging only")
Expand Down
110 changes: 104 additions & 6 deletions supervisor/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,15 @@
"judge_temperature", "judge_max_tokens", "worker_trust_level",
"notification_channels", "pause_handling_mode", "max_auto_interventions",
"poll_interval_sec", "read_lines",
"runtime_recovery_enabled", "runtime_recovery_profile",
"runtime_recovery_reset_timezone", "runtime_recovery_reset_grace_seconds",
"runtime_recovery_transient_delay_seconds",
"runtime_recovery_rate_limit_fallback_delay_seconds",
"runtime_recovery_quiet_windows", "runtime_recovery_max_attempts_per_run",
"provider_retry_enabled", "provider_retry_reset_timezone",
"provider_retry_reset_grace_seconds", "provider_retry_transient_delay_seconds",
"provider_retry_rate_limit_fallback_delay_seconds",
"provider_retry_skip_windows", "provider_retry_max_attempts_per_run",
"explainer_model", "explainer_temperature", "explainer_max_tokens",
"deep_explainer_model", "deep_explainer_temperature", "deep_explainer_max_tokens",
"clarification_escalation_confidence",
Expand Down Expand Up @@ -78,13 +87,20 @@ def coerce_config_value(key: str, value: str):
ftype = known[key].type
if value.lower() in ("null", "none", "~"):
return None
if ftype in ("float", float):
if _field_accepts(ftype, "float", float):
return float(value)
if ftype in ("int", int):
if _field_accepts(ftype, "int", int):
return int(value)
if _field_accepts(ftype, "bool", bool):
return value.lower() in ("1", "true", "yes", "on")
return value


def _field_accepts(ftype, type_name: str, pytype) -> bool:
text = str(ftype)
return ftype in (type_name, pytype) or type_name in text


@dataclass
class RuntimeConfig:
# -- Execution Surface --
Expand Down Expand Up @@ -139,6 +155,25 @@ class RuntimeConfig:
branch_confidence_threshold: float = 0.75
default_agent_timeout_sec: int = 300

# -- Runtime recovery --
runtime_recovery_enabled: bool | None = None
runtime_recovery_profile: str | None = None
runtime_recovery_reset_timezone: str | None = None
runtime_recovery_reset_grace_seconds: int | None = None
runtime_recovery_transient_delay_seconds: int | None = None
runtime_recovery_rate_limit_fallback_delay_seconds: int | None = None
runtime_recovery_quiet_windows: list[str] | None = None
runtime_recovery_max_attempts_per_run: int | None = None

# Back-compat aliases for configs written during the initial provider-retry rollout.
provider_retry_enabled: bool = True
provider_retry_reset_timezone: str = "Asia/Shanghai"
provider_retry_reset_grace_seconds: int = 60
provider_retry_transient_delay_seconds: int = 60
provider_retry_rate_limit_fallback_delay_seconds: int = 300
provider_retry_skip_windows: list[str] = field(default_factory=list)
provider_retry_max_attempts_per_run: int = 3

# -- Notifications --
notification_channels: list[dict] = field(default_factory=lambda: [
{"kind": "tmux_display"},
Expand Down Expand Up @@ -171,10 +206,12 @@ def from_env(cls, prefix: str = "SUPERVISOR_") -> "RuntimeConfig":
if field_name not in known:
continue
ftype = known[field_name].type
if ftype in ("float", float):
if _field_accepts(ftype, "float", float):
data[field_name] = float(val)
elif ftype in ("int", int):
elif _field_accepts(ftype, "int", int):
data[field_name] = int(val)
elif _field_accepts(ftype, "bool", bool):
data[field_name] = val.lower() in ("1", "true", "yes", "on")
else:
data[field_name] = val
return cls(**data)
Expand Down Expand Up @@ -216,10 +253,12 @@ def load(cls, config_path: str | Path | None = None) -> "RuntimeConfig":
if field_name not in known:
continue
ftype = known[field_name].type
if ftype in ("float", float):
if _field_accepts(ftype, "float", float):
setattr(base, field_name, float(val))
elif ftype in ("int", int):
elif _field_accepts(ftype, "int", int):
setattr(base, field_name, int(val))
elif _field_accepts(ftype, "bool", bool):
setattr(base, field_name, val.lower() in ("1", "true", "yes", "on"))
else:
setattr(base, field_name, val)
return base
Expand All @@ -233,6 +272,52 @@ def effective_target(self) -> str:
"""Resolve the effective surface target (surface_target > pane_target)."""
return self.surface_target or self.pane_target

def runtime_recovery_policy(self):
from supervisor.runtime_recovery import RuntimeRecoveryPolicy, policy_with_profile

defaults = RuntimeConfig()
windows = (
self.runtime_recovery_quiet_windows
if self.runtime_recovery_quiet_windows is not None
else self.provider_retry_skip_windows
)
if isinstance(windows, str):
windows = [item.strip() for item in windows.split(",") if item.strip()]
return policy_with_profile(RuntimeRecoveryPolicy(
enabled=(
self.runtime_recovery_enabled
if self.runtime_recovery_enabled is not None
else self.provider_retry_enabled
),
profile=self.runtime_recovery_profile or "",
reset_timezone=(
self.runtime_recovery_reset_timezone
if self.runtime_recovery_reset_timezone is not None
else self.provider_retry_reset_timezone
),
reset_grace_seconds=(
self.provider_retry_reset_grace_seconds
if self.runtime_recovery_reset_grace_seconds is None
else self.runtime_recovery_reset_grace_seconds
),
transient_delay_seconds=(
self.provider_retry_transient_delay_seconds
if self.runtime_recovery_transient_delay_seconds is None
else self.runtime_recovery_transient_delay_seconds
),
rate_limit_fallback_delay_seconds=(
self.provider_retry_rate_limit_fallback_delay_seconds
if self.runtime_recovery_rate_limit_fallback_delay_seconds is None
else self.runtime_recovery_rate_limit_fallback_delay_seconds
),
quiet_windows=tuple(windows or ()),
max_attempts_per_run=(
self.provider_retry_max_attempts_per_run
if self.runtime_recovery_max_attempts_per_run is None
else self.runtime_recovery_max_attempts_per_run
),
))

def default_config_yaml(self) -> str:
"""Render a commented YAML template suitable for ``init``."""
return (
Expand All @@ -253,6 +338,19 @@ def default_config_yaml(self) -> str:
"# trust: low | standard | high (high = minimal supervision)\n"
f"worker_trust_level: \"{self.worker_trust_level}\"\n"
"\n"
"# Runtime recovery: orthogonal to workflow steps; retries provider/connectivity failures.\n"
f"runtime_recovery_enabled: {str(self.runtime_recovery_enabled if self.runtime_recovery_enabled is not None else True).lower()}\n"
"# Optional preset profile. Example: \"glm\" reserves UTC+8 12:00-18:00 and uses 5h fallback.\n"
f"runtime_recovery_profile: \"{self.runtime_recovery_profile or ''}\"\n"
"# Naive reset timestamps in Chinese Claude/GLM output are usually UTC+8.\n"
f"runtime_recovery_reset_timezone: \"{self.runtime_recovery_reset_timezone or self.provider_retry_reset_timezone}\"\n"
f"runtime_recovery_reset_grace_seconds: {self.runtime_recovery_reset_grace_seconds if self.runtime_recovery_reset_grace_seconds is not None else self.provider_retry_reset_grace_seconds}\n"
f"runtime_recovery_transient_delay_seconds: {self.runtime_recovery_transient_delay_seconds if self.runtime_recovery_transient_delay_seconds is not None else self.provider_retry_transient_delay_seconds}\n"
f"runtime_recovery_rate_limit_fallback_delay_seconds: {self.runtime_recovery_rate_limit_fallback_delay_seconds if self.runtime_recovery_rate_limit_fallback_delay_seconds is not None else self.provider_retry_rate_limit_fallback_delay_seconds}\n"
"# Optional quiet windows, e.g. [\"UTC+8 12:00-18:00\"]\n"
"runtime_recovery_quiet_windows: []\n"
f"runtime_recovery_max_attempts_per_run: {self.runtime_recovery_max_attempts_per_run if self.runtime_recovery_max_attempts_per_run is not None else self.provider_retry_max_attempts_per_run}\n"
"\n"
"# LLM judge (set to null for stub/offline mode)\n"
"# Examples: anthropic/claude-haiku-4-5-20251001, openai/gpt-4o-mini\n"
f"judge_model: null\n"
Expand Down
1 change: 1 addition & 0 deletions supervisor/daemon/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -479,6 +479,7 @@ def _run_worker(self, entry: RunEntry, spec, state) -> None:
mode=self.config.pause_handling_mode,
max_auto_interventions=self.config.max_auto_interventions,
),
runtime_recovery_policy=self.config.runtime_recovery_policy(),
)
loop.run_sidecar(
spec, state, terminal,
Expand Down
106 changes: 105 additions & 1 deletion supervisor/loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import logging
import signal
import time
from datetime import datetime, timezone

from supervisor.domain.enums import DeliveryState, TopState, DecisionType
from supervisor.domain.models import (
Expand All @@ -26,6 +27,12 @@
from supervisor.interventions import AutoInterventionManager
from supervisor.notifications import NotificationEvent, NotificationManager
from supervisor.pause_summary import PAUSE_CLASSES, latest_human_escalation, summarize_state
from supervisor.runtime_recovery import (
RuntimeRecoveryObservation,
RuntimeRecoveryPolicy,
detect_runtime_recovery,
seconds_until_recovery,
)
from supervisor.progress import write_progress
from supervisor.protocol.reason_code import (
ESC_AUTHORIZATION_REQUIRED,
Expand Down Expand Up @@ -82,7 +89,8 @@ def __init__(self, store, judge_model: str | None = None,
judge_temperature: float = 0.1, judge_max_tokens: int = 512,
worker_profile: WorkerProfile | None = None,
notification_manager: NotificationManager | None = None,
auto_intervention_manager: AutoInterventionManager | None = None):
auto_intervention_manager: AutoInterventionManager | None = None,
runtime_recovery_policy: RuntimeRecoveryPolicy | None = None):
self.store = store
self.judge_client = JudgeClient(
model=judge_model,
Expand All @@ -98,6 +106,7 @@ def __init__(self, store, judge_model: str | None = None,
self.worker_profile = worker_profile or WorkerProfile()
self.notification_manager = notification_manager or NotificationManager()
self.auto_intervention_manager = auto_intervention_manager or AutoInterventionManager(mode="notify_only")
self.runtime_recovery_policy = runtime_recovery_policy or RuntimeRecoveryPolicy()
# Set while a sidecar loop is active; consulted by helpers that need
# to cooperate with daemon stop_event / SIGTERM.
self._interrupted_ref = None
Expand Down Expand Up @@ -904,6 +913,11 @@ def _run_sidecar_inner(
effective_idle_timeout_sec = ZERO_POLL_IDLE_TIMEOUT_SEC
pending_text = None
last_activity_at = time.monotonic()
recovery_attempts: dict[str, int] = {}
recovery_total_attempts = 0
scheduled_recovery: RuntimeRecoveryObservation | None = None
scheduled_recovery_signature = ""
exhausted_recovery_signatures: set[str] = set()
# Delivery ack: transient monotonic time of last injection (not persisted)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
delivery_ack_deadline = 0.0 # 0 = not awaiting ack
DELIVERY_ACK_TIMEOUT = 60 # seconds
Expand Down Expand Up @@ -1041,6 +1055,79 @@ def _run_sidecar_inner(
# 2. Parse checkpoint with identity
checkpoints = adapter.parse_checkpoints(text, run_id=state.run_id, surface_id=surface_id)
if not checkpoints:
wall_clock_now = datetime.now(timezone.utc)
recovery = detect_runtime_recovery(
text,
now=wall_clock_now,
policy=self.runtime_recovery_policy,
)
if recovery is not None:
if recovery_total_attempts < self.runtime_recovery_policy.max_attempts_per_run:
if scheduled_recovery is None or scheduled_recovery_signature != recovery.signature:
scheduled_recovery = recovery
scheduled_recovery_signature = recovery.signature
self.store.append_session_event(
Comment thread
coderabbitai[bot] marked this conversation as resolved.
state.run_id,
"runtime_recovery_scheduled",
{
"kind": recovery.kind,
"reason": recovery.reason,
"retry_at": recovery.retry_at.isoformat(),
"attempt": recovery_total_attempts + 1,
"max_attempts": self.runtime_recovery_policy.max_attempts_per_run,
},
)
self.store.save(state)
elif recovery.signature not in exhausted_recovery_signatures:
exhausted_recovery_signatures.add(recovery.signature)
self.store.append_session_event(
state.run_id,
"runtime_recovery_exhausted",
{
"kind": recovery.kind,
"reason": recovery.reason,
"attempts": recovery_total_attempts,
"max_attempts": self.runtime_recovery_policy.max_attempts_per_run,
},
)
self.store.save(state)

if (
scheduled_recovery is not None
and seconds_until_recovery(scheduled_recovery.retry_at, now=wall_clock_now) <= 0
):
attempts = recovery_attempts.get(scheduled_recovery.signature, 0)
instruction = self._build_runtime_recovery_instruction(
state, scheduled_recovery
)
state.last_injected_node_id = state.current_node_id
state.last_injected_attempt = state.current_attempt
state.last_injection_seq = state.checkpoint_seq
self.store.save(state)
if not self._inject_or_pause(state, terminal, instruction, spec=spec):
return
recovery_attempts[scheduled_recovery.signature] = attempts + 1
recovery_total_attempts += 1
self.store.append_session_event(
state.run_id,
"runtime_recovery_injected",
{
"kind": scheduled_recovery.kind,
"reason": scheduled_recovery.reason,
"attempt": recovery_total_attempts,
"max_attempts": self.runtime_recovery_policy.max_attempts_per_run,
},
)
scheduled_recovery = None
scheduled_recovery_signature = ""
delivery_ack_deadline = time.monotonic() + DELIVERY_ACK_TIMEOUT
time.sleep(effective_poll_interval)
continue

if scheduled_recovery is not None:
time.sleep(effective_poll_interval)
continue

if effective_idle_timeout_sec and effective_idle_timeout_sec > 0:
idle_for = now - last_activity_at
if idle_for >= effective_idle_timeout_sec:
Expand Down Expand Up @@ -1165,6 +1252,8 @@ def _run_sidecar_inner(
if (state.checkpoint_seq > state.last_injection_seq
and state.delivery_state not in (DeliveryState.IDLE, DeliveryState.STARTED_PROCESSING)):
self._set_delivery_state(state, DeliveryState.STARTED_PROCESSING, reason="checkpoint received")
scheduled_recovery = None
scheduled_recovery_signature = ""
logger.info("checkpoint: %s (id=%s)", checkpoint.summary, checkpoint.checkpoint_id)
if checkpoint.status in {"working", "step_done", "workflow_done"}:
self._reset_recovery_tracking(state, clear_escalations=False)
Expand Down Expand Up @@ -1347,6 +1436,21 @@ def _get_cwd(self, terminal, state=None) -> str | None:
return state.workspace_root
return None

def _build_runtime_recovery_instruction(self, state, recovery: RuntimeRecoveryObservation) -> HandoffInstruction:
content = (
"retry\n\n"
"Runtime recovery detected a provider/connectivity failure outside the task logic. "
f"Retry the last interrupted action and continue current_node={state.current_node_id}. "
f"Observed {recovery.kind}: {recovery.reason}"
)
return HandoffInstruction.make(
content=content,
node_id=state.current_node_id,
current_attempt=state.current_attempt,
triggered_by_decision_id="",
trigger_type="runtime_recovery",
)

def _wait_for_injection_window(self, state, terminal, *, instruction_id: str) -> tuple[bool, str]:
readiness_fn = getattr(terminal, "injection_readiness", None)
if not callable(readiness_fn):
Expand Down
Loading
Loading