From b65de94a8da4614b3d59757beb8b9140f657b8ea Mon Sep 17 00:00:00 2001 From: Aaryan Date: Mon, 17 Aug 2026 10:26:31 +0000 Subject: [PATCH 1/3] Live graph: distinguish transient from permanent failure and retry as a new node A failed node was terminal, and apply_patch refuses to re-add its id, so the planner had no way to express 'that timeout deserves another attempt, that 404 does not'. On the shipped benchmark this stranded 4 of 14 cases: a transient fetch failure left the downstream node pending forever and the run returned an empty answer. - core/live_graph/retry.py: type-first failure classification and a bounded backoff policy. A permanent exception type short-circuits before any message inspection, so a hostile path cannot buy a retry of a refused sandbox escape. - GraphPatch.retry: a retry creates a NEW node that inherits the failed node's parents and adopts its children, so the failed attempt keeps its state and its place in the journal and nothing is edited in place. - GraphStore enforces the attempt ceiling and one-attempt-per-failure independently of the planner. - runtime.py: the planner classifies before planning around a failure, resolves branches by retry lineage, and no longer parents the distiller on a failed node. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01FvXm9XSmn5E9fdXGypJKAZ --- .env.example | 6 + s13code/core/live_graph/__init__.py | 16 +- s13code/core/live_graph/core.py | 25 +- s13code/core/live_graph/retry.py | 195 ++++++++++++ s13code/core/live_graph/store.py | 106 ++++++- .../live_graph/tests/test_retry_policy.py | 239 ++++++++++++++ s13code/planner.py | 18 +- s13code/runtime.py | 165 ++++++++-- tests/test_retry_policy_adversarial.py | 293 ++++++++++++++++++ 9 files changed, 1014 insertions(+), 49 deletions(-) create mode 100644 s13code/core/live_graph/retry.py create mode 100644 s13code/core/live_graph/tests/test_retry_policy.py create mode 100644 tests/test_retry_policy_adversarial.py diff --git a/.env.example b/.env.example index c514f7e..5e66074 100644 --- a/.env.example +++ b/.env.example @@ -8,3 +8,9 @@ S13_LIVE_SEMANTIC_CHUNKING=1 # Set this to an absolute path before using local file skills. S13_SANDBOX_ROOT=/absolute/path/to/S13Code/sandbox + +# Retry policy for transiently failed graph nodes. Attempts include the first +# try, so 3 means one call plus at most two further attempts. Permanent +# failures are never retried regardless of these values. +S13_RETRY_MAX_ATTEMPTS=3 +S13_RETRY_BASE_DELAY=0.05 diff --git a/s13code/core/live_graph/__init__.py b/s13code/core/live_graph/__init__.py index 364524f..39c9713 100644 --- a/s13code/core/live_graph/__init__.py +++ b/s13code/core/live_graph/__init__.py @@ -6,11 +6,21 @@ GraphSnapshot, LiveGraphExecutor, NodeState, + RunReport, TaskSpec, ) -from .store import GraphStore +from .retry import ( + FailureClass, + FailureVerdict, + RetryPolicy, + RetryRequest, + classify_failure, + retry_node_id, +) +from .store import GraphMutationError, GraphStore __all__ = [ - "Event", "GraphPatch", "GraphSnapshot", "GraphStore", - "LiveGraphExecutor", "NodeState", "TaskSpec", + "Event", "FailureClass", "FailureVerdict", "GraphMutationError", "GraphPatch", + "GraphSnapshot", "GraphStore", "LiveGraphExecutor", "NodeState", "RetryPolicy", + "RetryRequest", "RunReport", "TaskSpec", "classify_failure", "retry_node_id", ] diff --git a/s13code/core/live_graph/core.py b/s13code/core/live_graph/core.py index 74ad143..ea1b790 100644 --- a/s13code/core/live_graph/core.py +++ b/s13code/core/live_graph/core.py @@ -13,6 +13,8 @@ from enum import StrEnum from typing import TYPE_CHECKING, Any, Awaitable, Callable, Protocol +from .retry import RetryRequest + if TYPE_CHECKING: from .store import GraphStore @@ -42,6 +44,11 @@ class GraphPatch: ``connect`` uses ``(parent, child)`` pairs. A waited node can be made runnable later with ``resume``; it is never silently retried. + + ``retry`` is the explicit alternative to that silence. It re-attempts a + *failed* node as a brand-new node, so the failed attempt keeps its state, + its result and its place in the journal. Nothing here re-runs a node in + place: an attempt is a fact, and facts are not edited. """ add: tuple[TaskSpec, ...] = () @@ -49,6 +56,7 @@ class GraphPatch: cancel: tuple[str, ...] = () wait: tuple[str, ...] = () resume: tuple[str, ...] = () + retry: tuple[RetryRequest, ...] = () finish: bool = False reason: str = "" @@ -82,6 +90,9 @@ class RunReport: finished: bool executed: tuple[str, ...] waiting: tuple[str, ...] + # Nodes created as another attempt at a failed node, oldest first. Empty on + # a clean run, which is what makes a non-empty value worth reading. + retried: tuple[str, ...] = () class LiveGraphExecutor: @@ -172,8 +183,18 @@ async def _execute(self, task: TaskSpec) -> tuple[TaskSpec, bool, dict[str, Any] worker = self.skills.get(task.skill) if worker is None: return task, False, {"error": f"unknown skill: {task.skill}"} + # Backoff is durable node state, not a scheduler-side timer, so a + # process that restarts mid-wait still honours it. Sleeping inside the + # worker slot also keeps a backing-off retry cancellable: the graph can + # still cancel it, and `asyncio.CancelledError` lands here as it would + # for any other in-flight task. + delay = float(task.metadata.get("retry_delay_seconds") or 0.0) try: + if delay > 0: + await asyncio.sleep(delay) return task, True, await worker(task) + except asyncio.CancelledError: + raise except Exception as exc: # worker failures become planner-visible events return task, False, {"error": f"{type(exc).__name__}: {exc}"} @@ -210,4 +231,6 @@ async def _replay_pending_planner_events(self, run_id: str) -> None: def _report(self, run_id: str, executed: list[str]) -> RunReport: snapshot = self.store.snapshot(run_id) waiting = tuple(nid for nid, node in snapshot.nodes.items() if node["state"] == NodeState.WAITING) - return RunReport(run_id, snapshot.finished, tuple(executed), waiting) + retried = tuple(nid for nid, node in sorted(snapshot.nodes.items()) + if node.get("metadata", {}).get("retry_of")) + return RunReport(run_id, snapshot.finished, tuple(executed), waiting, retried) diff --git a/s13code/core/live_graph/retry.py b/s13code/core/live_graph/retry.py new file mode 100644 index 0000000..d64d377 --- /dev/null +++ b/s13code/core/live_graph/retry.py @@ -0,0 +1,195 @@ +"""Failure classification and the retry policy the live planner reasons with. + +The live graph already treats a failure as an *outcome* the planner sees. What +it could not express was the difference between the two failures that matter: + +* a **transient** one - the network blinked, a provider was briefly busy, a + socket timed out - where the same work is worth attempting again, and +* a **permanent** one - the path escapes the sandbox, the document is missing, + the skill does not exist - where attempting it again is at best wasted budget + and at worst a repeated attempt at something the system already refused. + +Policy lives here and in the planner. *Mechanism* lives in ``GraphStore``: +scheduling a retry creates a new node, so the failed attempt stays visible in +the journal and a future node still cannot exist before its inputs. + +Classification order is deliberate and is the security-relevant part of this +module: **the exception type decides, and the message may only refine an +already-unknown verdict.** An attacker (or an unlucky filename) can put the +word "timeout" inside a ``PermissionError``; that must never buy a retry of a +refused sandbox escape. See ``tests/test_retry_policy_adversarial.py``. +""" + +from __future__ import annotations + +import re +from dataclasses import dataclass +from enum import StrEnum + + +class FailureClass(StrEnum): + TRANSIENT = "transient" + PERMANENT = "permanent" + + +# Worker failures reach the planner as ``f"{type(exc).__name__}: {exc}"``, so +# classification is by type *name*. That keeps this module free of httpx, +# grpc and provider imports while still recognising their errors. +PERMANENT_TYPES = frozenset({ + # the sandbox and the local filesystem + "PermissionError", "FileNotFoundError", "NotADirectoryError", "IsADirectoryError", + "FileExistsError", + # programming and contract errors: retrying cannot change the outcome + "ValueError", "TypeError", "KeyError", "IndexError", "AttributeError", + "NotImplementedError", "ZeroDivisionError", "AssertionError", "StopIteration", + "JSONDecodeError", "UnicodeDecodeError", "UnicodeEncodeError", "RecursionError", + # graph and protocol contracts + "GraphMutationError", "ValidationError", "UnsupportedProtocolError", + "InvalidURL", "UnsupportedProtocol", "LocalProtocolError", + # authorisation: a second identical request is still unauthorised + "PermissionDenied", "Unauthorized", "Forbidden", +}) + +TRANSIENT_TYPES = frozenset({ + # transport + "TimeoutError", "ConnectionError", "ConnectionResetError", "ConnectionAbortedError", + "ConnectionRefusedError", "BrokenPipeError", "socket.timeout", "OSError", + # httpx / requests / urllib3 + "ConnectTimeout", "ReadTimeout", "WriteTimeout", "PoolTimeout", "ConnectError", + "ReadError", "WriteError", "NetworkError", "ProxyError", "RemoteProtocolError", + "TransportError", "ProtocolError", "IncompleteRead", "ChunkedEncodingError", + # gRPC and provider-side capacity + "AioRpcError", "RpcError", "ProviderError", "ServiceUnavailable", + "TooManyRequests", "RateLimitError", "OverloadedError", "APIConnectionError", +}) + +# Only consulted when the type is unknown to both sets above. +_TRANSIENT_HINT = re.compile( + r"\b(timed?\s*out|timeout|temporarily unavailable|try again|connection (?:reset|refused|aborted)" + r"|broken pipe|rate.?limit|too many requests|overloaded|service unavailable|bad gateway" + r"|gateway time-?out|server error)\b", + re.IGNORECASE, +) +_PERMANENT_HINT = re.compile( + r"\b(not found|no such file|does not exist|unknown skill|escapes|forbidden|unauthorized" + r"|permission denied|invalid|malformed|unsupported|not permitted)\b", + re.IGNORECASE, +) +_HTTP_STATUS = re.compile(r"\b(?:status(?:\s*code)?\s*[:=]?\s*|HTTP\s+|')(\d{3})\b") + +# 408 request timeout, 425 too early, 429 too many requests, plus 5xx. +_TRANSIENT_STATUS = frozenset({408, 425, 429, 500, 502, 503, 504, 507, 509, 598, 599}) + + +@dataclass(frozen=True) +class FailureVerdict: + failure_class: FailureClass + error_type: str + basis: str + + @property + def retryable(self) -> bool: + return self.failure_class is FailureClass.TRANSIENT + + def as_payload(self) -> dict[str, str]: + return {"failure_class": self.failure_class.value, "error_type": self.error_type, "basis": self.basis} + + +def _error_type(error: str) -> str: + """Read the leading exception name out of ``"TypeName: message"``.""" + head = error.split(":", 1)[0].strip() + return head if re.fullmatch(r"[A-Za-z_][A-Za-z0-9_.]*", head or "") else "" + + +def classify_failure(error: str) -> FailureVerdict: + """Decide whether a recorded failure payload is worth attempting again. + + The order below is the whole point. A known *permanent* type short-circuits + before any message inspection, so a hostile or unlucky message cannot talk + the graph into retrying work the system already refused. Message hints are + the last resort, used only for types neither table recognises, and even + then a permanent hint outranks a transient one. + """ + error = (error or "").strip() + if not error: + return FailureVerdict(FailureClass.PERMANENT, "", "empty failure payload is not evidence of a retryable fault") + + error_type = _error_type(error) + # Vendored exceptions arrive as `httpx.ConnectTimeout`; match the leaf name too. + leaf = error_type.rsplit(".", 1)[-1] if error_type else "" + + if error_type in PERMANENT_TYPES or leaf in PERMANENT_TYPES: + return FailureVerdict(FailureClass.PERMANENT, error_type, f"{leaf} is a permanent failure type") + if error_type in TRANSIENT_TYPES or leaf in TRANSIENT_TYPES: + # An HTTP transport error still carries a status; a 404 delivered by a + # transport exception is permanent no matter which class wrapped it. + status = _status_code(error) + if status is not None and status not in _TRANSIENT_STATUS: + return FailureVerdict(FailureClass.PERMANENT, error_type, + f"{leaf} reported HTTP {status}, which will not change on retry") + return FailureVerdict(FailureClass.TRANSIENT, error_type, f"{leaf} is a transport/capacity failure type") + + status = _status_code(error) + if status is not None: + transient = status in _TRANSIENT_STATUS + return FailureVerdict(FailureClass.TRANSIENT if transient else FailureClass.PERMANENT, error_type, + f"HTTP {status} is {'retryable' if transient else 'terminal'}") + + if _PERMANENT_HINT.search(error): + return FailureVerdict(FailureClass.PERMANENT, error_type, "message describes a terminal condition") + if _TRANSIENT_HINT.search(error): + return FailureVerdict(FailureClass.TRANSIENT, error_type, "message describes a temporary condition") + # Unknown failures are permanent by default: a graph that retries anything + # it does not understand is a graph that burns budget on real bugs. + return FailureVerdict(FailureClass.PERMANENT, error_type, "unrecognised failure defaults to permanent") + + +def _status_code(error: str) -> int | None: + match = _HTTP_STATUS.search(error) + if not match: + return None + value = int(match.group(1)) + return value if 100 <= value < 600 else None + + +@dataclass(frozen=True) +class RetryRequest: + """The planner's request to attempt one failed node again. + + ``delay_seconds`` is carried on the new node rather than slept here, so the + wait is durable state a resumed process can still see. + """ + + node_id: str + reason: str = "" + delay_seconds: float = 0.0 + + +@dataclass(frozen=True) +class RetryPolicy: + """Bounded exponential backoff. ``max_attempts`` counts the first try.""" + + max_attempts: int = 3 + base_delay_seconds: float = 0.05 + max_delay_seconds: float = 2.0 + + def __post_init__(self) -> None: + if self.max_attempts < 1: + raise ValueError("max_attempts must be at least one") + if self.base_delay_seconds < 0 or self.max_delay_seconds < 0: + raise ValueError("delays must not be negative") + + def may_retry(self, attempt: int) -> bool: + return attempt < self.max_attempts + + def delay_for(self, attempt: int) -> float: + return min(self.max_delay_seconds, self.base_delay_seconds * (2 ** max(0, attempt - 1))) + + +def retry_node_id(root_id: str, attempt: int) -> str: + """Stable id for the next attempt. + + The root prefix is preserved on purpose: planners that group a frontier by + prefix (``search_``, ``index_``) keep seeing the retry as part of it. + """ + return f"{root_id}__retry{attempt}" diff --git a/s13code/core/live_graph/store.py b/s13code/core/live_graph/store.py index 9e06503..1916eea 100644 --- a/s13code/core/live_graph/store.py +++ b/s13code/core/live_graph/store.py @@ -8,6 +8,7 @@ from typing import Any from .core import Event, GraphPatch, GraphSnapshot, NodeState, TaskSpec +from .retry import retry_node_id class GraphMutationError(ValueError): @@ -17,7 +18,13 @@ class GraphMutationError(ValueError): class GraphStore: """Durable graph storage. A patch and its journal entry commit together.""" - def __init__(self, path: str | Path): + def __init__(self, path: str | Path, *, max_attempts: int = 3): + # A hard ceiling on attempts belongs in the mechanism, not only in the + # planner: an LLM planner that proposes a retry every time it sees a + # failure must not be able to spend an unbounded budget. + if max_attempts < 1: + raise ValueError("max_attempts must be at least one") + self.max_attempts = max_attempts self.path = str(path) # FastAPI's test/client boundary (and a local server's worker thread) # may resume a durable run from a different thread than construction. @@ -145,7 +152,11 @@ def apply_patch(self, run_id: str, patch: GraphPatch, *, trigger_event: int) -> add_ids = [task.id for task in patch.add] if len(add_ids) != len(set(add_ids)) or ids.intersection(add_ids): raise GraphMutationError("patch adds duplicate task id") - all_ids = ids.union(add_ids) + retry_plan = self._plan_retries(existing, patch) + retry_ids = [spec.id for spec, _, _ in retry_plan] + if set(retry_ids) & (ids | set(add_ids)) or len(retry_ids) != len(set(retry_ids)): + raise GraphMutationError("retry would reuse an existing task id") + all_ids = ids.union(add_ids, retry_ids) for parent, child in patch.connect: if parent not in all_ids or child not in all_ids: raise GraphMutationError(f"edge {parent}->{child} references unknown task") @@ -168,6 +179,8 @@ def apply_patch(self, run_id: str, patch: GraphPatch, *, trigger_event: int) -> for task in patch.add: self.db.execute("INSERT INTO nodes VALUES (?, ?, ?, ?, ?, ?, NULL)", (run_id, task.id, task.skill, json.dumps(task.input), json.dumps(task.metadata), NodeState.PENDING)) + for spec, request, attempt in retry_plan: + self._insert_retry(run_id, spec, request, attempt) for parent, child in patch.connect: self.db.execute("INSERT OR IGNORE INTO edges VALUES (?, ?, ?)", (run_id, parent, child)) # Detect cycles after adding all edges, while the transaction can still roll back. @@ -196,15 +209,92 @@ def apply_patch(self, run_id: str, patch: GraphPatch, *, trigger_event: int) -> self._event(run_id, "task_cancelled", leftover["id"], {"reason": "graph finished", "was_running": leftover["state"] == NodeState.RUNNING}) self.db.execute("UPDATE runs SET finished=1 WHERE id=?", (run_id,)) - self._event(run_id, "graph_patched", None, {"trigger_event": trigger_event, "reason": patch.reason, - "add": add_ids, "connect": list(patch.connect), "cancel": list(patch.cancel), - "wait": list(patch.wait), "resume": list(patch.resume), "finish": patch.finish}) + summary = {"trigger_event": trigger_event, "reason": patch.reason, + "add": add_ids, "connect": list(patch.connect), "cancel": list(patch.cancel), + "wait": list(patch.wait), "resume": list(patch.resume), + "retry": retry_ids, "finish": patch.finish} + self._event(run_id, "graph_patched", None, summary) self.db.execute("INSERT INTO patches(run_id, trigger_event, patch_json) VALUES (?, ?, ?)", - (run_id, trigger_event, json.dumps({"reason": patch.reason, "add": add_ids, - "connect": list(patch.connect), "cancel": list(patch.cancel), - "wait": list(patch.wait), "resume": list(patch.resume), "finish": patch.finish}))) + (run_id, trigger_event, json.dumps({key: value for key, value in summary.items() + if key != "trigger_event"}))) return True + # -- retry mechanism --------------------------------------------------- + # + # A retry is a *new node*, never an edited one. The failed attempt keeps + # its FAILED state and its recorded error, and the journal keeps the whole + # lineage. The new node inherits the failed node's edges so it occupies + # exactly the same position in the dependency order: it cannot run before + # the inputs the original needed, and whatever waited on the original now + # waits on the attempt that might actually succeed. + + def _attempt_number(self, nodes: dict[str, dict[str, Any]], node_id: str) -> int: + return int(nodes[node_id]["metadata"].get("retry_attempt", 1)) + + def _root_of(self, nodes: dict[str, dict[str, Any]], node_id: str) -> str: + return str(nodes[node_id]["metadata"].get("retry_root") or node_id) + + def _plan_retries(self, existing: GraphSnapshot, patch: GraphPatch + ) -> list[tuple[TaskSpec, Any, int]]: + """Validate every retry request and build the node it would create.""" + planned: list[tuple[TaskSpec, Any, int]] = [] + seen: set[str] = set() + # Which nodes already have a successor attempt. A failed node may be + # retried at most once; without that, a replayed or duplicated planner + # decision fans one failure out into several concurrent attempts of + # the same work, which is the exact bug the journal is meant to prevent. + already_retried = {str(node["metadata"]["retry_of"]) for node in existing.nodes.values() + if node["metadata"].get("retry_of")} + for request in patch.retry: + node = existing.nodes.get(request.node_id) + if node is None: + raise GraphMutationError(f"patch retries unknown task {request.node_id}") + if node["state"] != NodeState.FAILED: + raise GraphMutationError( + f"can only retry a failed task {request.node_id}, got {node['state']}") + if request.node_id in already_retried or request.node_id in seen: + raise GraphMutationError(f"task {request.node_id} has already been retried") + root = self._root_of(existing.nodes, request.node_id) + attempt = self._attempt_number(existing.nodes, request.node_id) + 1 + if attempt > self.max_attempts: + raise GraphMutationError( + f"task {root} exhausted its {self.max_attempts} attempts") + seen.add(request.node_id) + metadata = {**node["metadata"], "retry_of": request.node_id, "retry_root": root, + "retry_attempt": attempt, "retry_reason": request.reason, + "retry_delay_seconds": max(0.0, float(request.delay_seconds))} + planned.append((TaskSpec(retry_node_id(root, attempt), node["skill"], + dict(node["input"]), metadata), request, attempt)) + return planned + + def _insert_retry(self, run_id: str, spec: TaskSpec, request: Any, attempt: int) -> None: + self.db.execute("INSERT INTO nodes VALUES (?, ?, ?, ?, ?, ?, NULL)", + (run_id, spec.id, spec.skill, json.dumps(spec.input), + json.dumps(spec.metadata), NodeState.PENDING)) + # Inherit both directions: parents so the attempt still cannot start + # before its inputs, children so nothing downstream is left waiting on + # a node that will never succeed. + parents = [row["parent_id"] for row in self.db.execute( + "SELECT parent_id FROM edges WHERE run_id=? AND child_id=? ORDER BY parent_id", + (run_id, request.node_id))] + children = [row["child_id"] for row in self.db.execute( + "SELECT child_id FROM edges WHERE run_id=? AND parent_id=? ORDER BY child_id", + (run_id, request.node_id))] + for parent in parents: + self.db.execute("INSERT OR IGNORE INTO edges VALUES (?, ?, ?)", (run_id, parent, spec.id)) + for child in children: + self.db.execute("INSERT OR IGNORE INTO edges VALUES (?, ?, ?)", (run_id, spec.id, child)) + # The failed attempt must stop gating its children, or the child + # stays unreachable forever: `ready()` requires every parent to + # have succeeded, and a FAILED parent never will. + self.db.execute("DELETE FROM edges WHERE run_id=? AND parent_id=? AND child_id=?", + (run_id, request.node_id, child)) + self._event(run_id, "task_retry_scheduled", spec.id, { + "retry_of": request.node_id, "retry_root": spec.metadata["retry_root"], "attempt": attempt, + "max_attempts": self.max_attempts, "delay_seconds": spec.metadata["retry_delay_seconds"], + "reason": request.reason, "inherited_parents": parents, "adopted_children": children, + }) + def is_finished(self, run_id: str) -> bool: row = self.db.execute("SELECT finished FROM runs WHERE id=?", (run_id,)).fetchone() return bool(row and row["finished"]) diff --git a/s13code/core/live_graph/tests/test_retry_policy.py b/s13code/core/live_graph/tests/test_retry_policy.py new file mode 100644 index 0000000..82cc862 --- /dev/null +++ b/s13code/core/live_graph/tests/test_retry_policy.py @@ -0,0 +1,239 @@ +"""Invariants for outcome-driven retry of transiently failed graph nodes.""" + +import asyncio + +import pytest + +from s13code.core.live_graph import ( + FailureClass, + GraphMutationError, + GraphPatch, + GraphStore, + LiveGraphExecutor, + RetryPolicy, + RetryRequest, + TaskSpec, + classify_failure, +) + + +class ScriptedPlanner: + def __init__(self, script): + self.script = script + + async def plan(self, graph, event): + return self.script(graph, event) + + +def transient_then_ok(store_path, *, failures: int): + """Planner + worker pair: fail `failures` times transiently, then succeed.""" + attempts = {"count": 0} + + async def worker(task): + attempts["count"] += 1 + if attempts["count"] <= failures: + raise ConnectionResetError("connection reset by peer") + return {"ok": True, "attempt": attempts["count"]} + + policy = RetryPolicy(max_attempts=3, base_delay_seconds=0.0) + + def plan(graph, event): + if event.kind == "run_started": + return GraphPatch(add=(TaskSpec("call", "worker"),), reason="first frontier") + if event.kind == "task_failed": + node = graph.nodes[event.node_id] + attempt = int(node["metadata"].get("retry_attempt", 1)) + verdict = classify_failure(event.payload["error"]) + if verdict.retryable and policy.may_retry(attempt): + return GraphPatch(retry=(RetryRequest(event.node_id, reason=verdict.basis),), + reason="transient failure is worth another attempt") + return GraphPatch(finish=True, reason="attempts exhausted") + return GraphPatch(finish=True, reason="done") + + return GraphStore(store_path, max_attempts=3), ScriptedPlanner(plan), {"worker": worker}, attempts + + +# -- classification --------------------------------------------------------- + +@pytest.mark.parametrize("error, expected", [ + ("ConnectTimeout: timed out", FailureClass.TRANSIENT), + ("ProxyError: 503 Service Unavailable", FailureClass.TRANSIENT), + ("ConnectionResetError: connection reset by peer", FailureClass.TRANSIENT), + ("PermissionError: path escapes S13_SANDBOX_ROOT: /etc/passwd", FailureClass.PERMANENT), + ("FileNotFoundError: papers/missing.txt", FailureClass.PERMANENT), + ("ValueError: could not safely parse a birthday date", FailureClass.PERMANENT), + ("unknown skill: teleport", FailureClass.PERMANENT), + ("", FailureClass.PERMANENT), +]) +def test_classification(error, expected): + assert classify_failure(error).failure_class is expected + + +def test_http_status_outranks_the_transport_class_that_carried_it(): + """A 404 delivered by a transport exception is still a 404.""" + assert classify_failure("HTTPStatusError: Client error '404 Not Found' for url " + "'https://example.test/x'").failure_class is FailureClass.PERMANENT + assert classify_failure("HTTPStatusError: Server error '503 Service Unavailable' for url " + "'https://example.test/x'").failure_class is FailureClass.TRANSIENT + + +def test_unrecognised_failures_default_to_permanent(): + verdict = classify_failure("WeirdVendorError: something happened") + assert verdict.failure_class is FailureClass.PERMANENT + assert "defaults to permanent" in verdict.basis + + +# -- mechanism -------------------------------------------------------------- + +@pytest.mark.asyncio +async def test_transient_failure_is_retried_and_eventually_succeeds(tmp_path): + store, planner, skills, attempts = transient_then_ok(tmp_path / "graph.db", failures=2) + report = await LiveGraphExecutor(store, planner, skills).run("retry-run") + + assert attempts["count"] == 3 + assert report.retried == ("call__retry2", "call__retry3") + nodes = store.snapshot("retry-run").nodes + # Every attempt keeps its own state. Nothing was edited in place. + assert nodes["call"]["state"] == "failed" + assert nodes["call__retry2"]["state"] == "failed" + assert nodes["call__retry3"]["state"] == "succeeded" + assert nodes["call__retry3"]["result"] == {"ok": True, "attempt": 3} + assert nodes["call"]["result"]["error"].startswith("ConnectionResetError") + + +@pytest.mark.asyncio +async def test_a_retry_node_does_not_exist_before_the_failure_that_created_it(tmp_path): + store, planner, skills, _ = transient_then_ok(tmp_path / "graph.db", failures=1) + await LiveGraphExecutor(store, planner, skills).run("ordering") + + events = store.events("ordering") + kinds = [(event.kind, event.node_id) for event in events] + failure_at = kinds.index(("task_failed", "call")) + scheduled_at = kinds.index(("task_retry_scheduled", "call__retry2")) + started_at = kinds.index(("task_started", "call__retry2")) + assert failure_at < scheduled_at < started_at + + # No event mentions the retry node before its parent failed. + assert all(event.node_id != "call__retry2" for event in events[:failure_at]) + + +@pytest.mark.asyncio +async def test_retry_inherits_parents_and_adopts_children(tmp_path): + async def source(task): + return {"value": 1} + + async def flaky(task): + raise TimeoutError("timed out") + + async def sink(task): + return {"done": True} + + def plan(graph, event): + if event.kind == "run_started": + return GraphPatch( + add=(TaskSpec("upstream", "source"), TaskSpec("middle", "flaky"), TaskSpec("downstream", "sink")), + connect=(("upstream", "middle"), ("middle", "downstream")), reason="pipeline") + if event.kind == "task_failed" and event.node_id == "middle": + return GraphPatch(retry=(RetryRequest("middle", reason="transient"),), reason="retry") + if event.kind == "task_failed": + return GraphPatch(finish=True, reason="give up") + return GraphPatch() + + store = GraphStore(tmp_path / "graph.db", max_attempts=2) + await LiveGraphExecutor(store, ScriptedPlanner(plan), {"source": source, "flaky": flaky, "sink": sink}).run("edges") + + snapshot = store.snapshot("edges") + edges = set(snapshot.edges) + assert ("upstream", "middle__retry2") in edges, "the attempt still waits on its own inputs" + assert ("middle__retry2", "downstream") in edges, "downstream now waits on the live attempt" + # The failed node must stop gating its child, or the child is unreachable + # forever: ready() requires every parent to have succeeded. + assert ("middle", "downstream") not in edges + assert ("upstream", "middle") in edges, "history of the failed attempt is preserved" + + +@pytest.mark.asyncio +async def test_permanent_failure_is_never_retried(tmp_path): + calls = {"count": 0} + + async def worker(task): + calls["count"] += 1 + raise PermissionError("path escapes S13_SANDBOX_ROOT: /etc/shadow") + + def plan(graph, event): + if event.kind == "run_started": + return GraphPatch(add=(TaskSpec("read", "worker"),), reason="first") + if event.kind == "task_failed": + verdict = classify_failure(event.payload["error"]) + if verdict.retryable: + return GraphPatch(retry=(RetryRequest(event.node_id),), reason="retry") + return GraphPatch(finish=True, reason=f"permanent: {verdict.basis}") + return GraphPatch() + + store = GraphStore(tmp_path / "graph.db") + report = await LiveGraphExecutor(store, ScriptedPlanner(plan), {"worker": worker}).run("permanent") + assert calls["count"] == 1 + assert report.retried == () + assert not any(event.kind == "task_retry_scheduled" for event in store.events("permanent")) + + +def test_store_enforces_the_attempt_ceiling_independently(tmp_path): + """The planner may ask for anything; the store still bounds the budget.""" + store = GraphStore(tmp_path / "graph.db", max_attempts=2) + store.start("ceiling") + store.apply_patch("ceiling", GraphPatch(add=(TaskSpec("call", "worker"),)), trigger_event=1) + store.mark_running("ceiling", [TaskSpec("call", "worker")]) + store.record_outcome("ceiling", "call", False, {"error": "ConnectTimeout: timed out"}) + store.apply_patch("ceiling", GraphPatch(retry=(RetryRequest("call"),)), trigger_event=2) + + store.mark_running("ceiling", [TaskSpec("call__retry2", "worker")]) + store.record_outcome("ceiling", "call__retry2", False, {"error": "ConnectTimeout: timed out"}) + with pytest.raises(GraphMutationError, match="exhausted its 2 attempts"): + store.apply_patch("ceiling", GraphPatch(retry=(RetryRequest("call__retry2"),)), trigger_event=3) + store.close() + + +def test_only_a_failed_node_may_be_retried(tmp_path): + store = GraphStore(tmp_path / "graph.db") + store.start("states") + store.apply_patch("states", GraphPatch(add=(TaskSpec("call", "worker"),)), trigger_event=1) + with pytest.raises(GraphMutationError, match="can only retry a failed task"): + store.apply_patch("states", GraphPatch(retry=(RetryRequest("call"),)), trigger_event=2) + with pytest.raises(GraphMutationError, match="retries unknown task"): + store.apply_patch("states", GraphPatch(retry=(RetryRequest("ghost"),)), trigger_event=3) + store.close() + + +@pytest.mark.asyncio +async def test_backoff_is_durable_node_state_not_a_scheduler_timer(tmp_path): + store, planner, skills, _ = transient_then_ok(tmp_path / "graph.db", failures=1) + await LiveGraphExecutor(store, planner, skills).run("backoff") + node = store.snapshot("backoff").nodes["call__retry2"] + # A resumed process reads the wait from the row, not from lost memory. + assert "retry_delay_seconds" in node["metadata"] + assert node["metadata"]["retry_of"] == "call" + assert node["metadata"]["retry_root"] == "call" + assert node["metadata"]["retry_attempt"] == 2 + + +@pytest.mark.asyncio +async def test_retry_delay_is_actually_waited(tmp_path): + store = GraphStore(tmp_path / "graph.db", max_attempts=2) + started = [] + + async def worker(task): + started.append(asyncio.get_running_loop().time()) + if len(started) == 1: + raise TimeoutError("timed out") + return {"ok": True} + + def plan(graph, event): + if event.kind == "run_started": + return GraphPatch(add=(TaskSpec("call", "worker"),), reason="first") + if event.kind == "task_failed": + return GraphPatch(retry=(RetryRequest("call", delay_seconds=0.15),), reason="wait then retry") + return GraphPatch(finish=True, reason="done") + + await LiveGraphExecutor(store, ScriptedPlanner(plan), {"worker": worker}).run("delay") + assert len(started) == 2 + assert started[1] - started[0] >= 0.15 diff --git a/s13code/planner.py b/s13code/planner.py index c167b2b..406ab30 100644 --- a/s13code/planner.py +++ b/s13code/planner.py @@ -6,7 +6,7 @@ from collections.abc import Awaitable, Callable from typing import Any, Protocol -from s13code.core.live_graph import Event, GraphPatch, GraphSnapshot, TaskSpec +from s13code.core.live_graph import Event, GraphPatch, GraphSnapshot, RetryRequest, TaskSpec TextLLM = Callable[[str, str], Awaitable[dict[str, Any]]] _ID = re.compile(r"^[A-Za-z][A-Za-z0-9_-]{0,63}$") @@ -46,7 +46,8 @@ async def plan(self, graph: GraphSnapshot, event: Event) -> GraphPatch: def _parse(self, text: str) -> GraphPatch: data = json.loads(text) - if not isinstance(data, dict) or set(data).difference({"add", "connect", "cancel", "wait", "resume", "finish", "reason"}): + if not isinstance(data, dict) or set(data).difference( + {"add", "connect", "cancel", "wait", "resume", "retry", "finish", "reason"}): raise ValueError("invalid GraphPatch object") add: list[TaskSpec] = [] for raw in data.get("add", []): @@ -69,13 +70,20 @@ def id_list(name: str) -> tuple[str, ...]: raise ValueError("invalid connect") if not isinstance(data.get("finish", False), bool) or not isinstance(data.get("reason", ""), str): raise ValueError("invalid finish/reason") + # A model may ask for another attempt, but never for how many or how + # long: delay is the runtime's policy and the attempt ceiling is + # enforced by GraphStore, so a looping model cannot spend the budget. + retry = tuple(RetryRequest(node_id, reason="planner requested another attempt") + for node_id in id_list("retry")) return GraphPatch(tuple(add), tuple((e[0], e[1]) for e in edges), id_list("cancel"), id_list("wait"), - id_list("resume"), data.get("finish", False), data.get("reason", "")) + id_list("resume"), retry, data.get("finish", False), data.get("reason", "")) def _prompt(self, graph: GraphSnapshot, event: Event) -> str: return json.dumps({"goal": self.goal, "event": {"kind": event.kind, "node_id": event.node_id, "payload": event.payload}, "nodes": [{"id": i, "role": n["skill"], "state": n["state"]} for i, n in graph.nodes.items()], "allowed_roles": sorted(self.roles), "patch_schema": {"add": [{"id": "safe_id", "role": "allowed_role", "input": {}, "metadata": {"agent": "role"}}], - "connect": [["parent", "child"]], "cancel": [], "wait": [], "resume": [], "finish": False, "reason": "short"}, - "rules": ["emit only next useful work", "roles never call tools directly", "return JSON only"]}) + "connect": [["parent", "child"]], "cancel": [], "wait": [], "resume": [], "retry": [], + "finish": False, "reason": "short"}, + "rules": ["emit only next useful work", "roles never call tools directly", + "retry only a failed node whose failure looks transient", "return JSON only"]}) diff --git a/s13code/runtime.py b/s13code/runtime.py index 45c16e3..adc2663 100644 --- a/s13code/runtime.py +++ b/s13code/runtime.py @@ -14,7 +14,15 @@ from pathlib import Path from typing import Any -from s13code.core.live_graph import GraphPatch, GraphStore, LiveGraphExecutor, TaskSpec +from s13code.core.live_graph import ( + GraphPatch, + GraphStore, + LiveGraphExecutor, + RetryPolicy, + RetryRequest, + TaskSpec, + classify_failure, +) from s13code.core.memory import MemoryKind, MemoryRecord, MemoryScope, MemoryStore, Principal, SourceRef from s13code.core.memory.embeddings import OllamaNomicEmbedder from s13code.planner import ConstrainedGraphPatchPlanner @@ -79,6 +87,22 @@ def _work_intent(prompt: str) -> tuple[str, list[TaskSpec]]: return "memory", [TaskSpec("recall", "memory_recall", {"query": prompt})] +_TERMINAL_UNSUCCESSFUL = {"failed", "cancelled"} + + +def _can_still_run(graph, node_id: str) -> bool: + """True while every parent of ``node_id`` could still succeed. + + ``GraphStore.ready`` releases a node only when *all* of its parents have + succeeded, so one failed or cancelled parent makes a pending child + permanently unreachable. The planner needs to see that as a decision point + rather than leaving the node to sit in the graph forever. + """ + parents = [parent for parent, child in graph.edges if child == node_id] + return all(graph.nodes[parent]["state"] not in _TERMINAL_UNSUCCESSFUL + for parent in parents if parent in graph.nodes) + + class S13Runtime: """Owns the persistent stores and runs one user request through the graph.""" @@ -90,7 +114,11 @@ def __init__(self, root: Path | None = None) -> None: self.root = root or Path(os.getenv("S13_DATA_DIR", str(Path.home() / ".s13code"))) self.root.mkdir(parents=True, exist_ok=True) self.memory = MemoryStore(self.root / "memory.sqlite", embedder=OllamaNomicEmbedder()) - self.graph = GraphStore(self.root / "graph.sqlite") + self.retry_policy = RetryPolicy( + max_attempts=int(os.getenv("S13_RETRY_MAX_ATTEMPTS", "3")), + base_delay_seconds=float(os.getenv("S13_RETRY_BASE_DELAY", "0.05")), + ) + self.graph = GraphStore(self.root / "graph.sqlite", max_attempts=self.retry_policy.max_attempts) def close(self) -> None: self.memory.close() @@ -121,6 +149,7 @@ async def run(self, *, prompt: str | None, scope: MemoryScope | None, llm: TextL user_source = SourceRef(source_uri, source_author, excerpt=prompt) runtime = self + retry_policy = self.retry_policy mode, initial_frontier = _work_intent(prompt) explicit_memory = bool(re.search( r"\b(remember|save (?:this|that)|keep (?:this|that) in mind|correction:)\b", prompt, re.IGNORECASE @@ -133,74 +162,126 @@ def answer_patch(graph, *, reason: str) -> GraphPatch: return GraphPatch() parents = tuple(node_id for node_id, node in graph.nodes.items() if node_id != "answer" and node["state"] == "succeeded") + # Anything still pending can no longer become ready once the + # answer is the last outstanding work: its own parents already + # reached a terminal non-success. Cancel it explicitly rather + # than leaving a node the executor will silently step over. + stranded = tuple(node_id for node_id, node in graph.nodes.items() + if node_id != "answer" and node["state"] == "pending" + and not _can_still_run(graph, node_id)) return GraphPatch(add=(TaskSpec("answer", "answer_with_evidence", {"query": prompt}),), - connect=tuple((parent, "answer") for parent in parents), reason=reason) + connect=tuple((parent, "answer") for parent in parents), + cancel=stranded, reason=reason) + + @staticmethod + def retry_patch(graph, event) -> GraphPatch | None: + """Decide whether a failure is worth attempting again. + + Classification is the policy; ``GraphStore`` supplies the + mechanism and independently enforces the attempt ceiling. + """ + node = graph.nodes.get(event.node_id) + if node is None or node["metadata"].get("no_retry"): + return None + verdict = classify_failure(str(event.payload.get("error", ""))) + attempt = int(node["metadata"].get("retry_attempt", 1)) + if not verdict.retryable or not retry_policy.may_retry(attempt): + return None + return GraphPatch( + retry=(RetryRequest(event.node_id, reason=verdict.basis, + delay_seconds=retry_policy.delay_for(attempt)),), + reason=f"{event.node_id} failed transiently ({verdict.basis}); " + f"attempt {attempt + 1} of {retry_policy.max_attempts}", + ) async def plan(self, graph, event): + # Retry is evaluated before every frontier rule below. A + # transiently failed node has not finished being work yet, so + # no downstream stage should be planned around its absence. + if event.kind == "task_failed" and event.node_id and event.node_id != "answer": + retry = self.retry_patch(graph, event) + if retry is not None: + return retry + # A retry attempt carries a new node id. Every rule below is + # written against the *lineage* the user's request produced, + # while edges are drawn from the attempt that actually ran. + node = graph.nodes.get(event.node_id) or {} + lineage = str((node.get("metadata") or {}).get("retry_root") or event.node_id or "") + source = event.node_id if event.kind == "run_started": first = list(initial_frontier) if explicit_memory: first.append(TaskSpec("remember", "remember_explicit_fact", {"text": prompt})) return GraphPatch(add=tuple(first), reason=f"first frontier selected for {mode}") - if event.node_id == "index_file" and event.kind == "task_succeeded": + if lineage == "index_file" and event.kind == "task_succeeded": return GraphPatch(add=(TaskSpec("recall", "memory_recall", {"query": prompt}),), - connect=(("index_file", "recall"),), + connect=((source, "recall"),), reason="file is indexed; retrieval can now inspect it") - if mode == "birthday_reminder" and event.node_id == "remember" and event.kind == "task_succeeded": + if mode == "birthday_reminder" and lineage == "remember" and event.kind == "task_succeeded": return GraphPatch(add=(TaskSpec("reminder", "create_reminder", {"prompt": prompt}, {"agent": "calendar_writer"}),), - connect=(("remember", "reminder"),), reason="explicit fact is durable; create calendar artifacts") - if mode == "birthday_reminder" and event.node_id == "reminder": + connect=((source, "reminder"),), reason="explicit fact is durable; create calendar artifacts") + if mode == "birthday_reminder" and lineage == "reminder": return self.answer_patch(graph, reason="calendar reminder artifacts are ready") - if mode == "birthday_reminder" and event.node_id == "recall": + if mode == "birthday_reminder" and lineage == "recall": return GraphPatch() - if mode == "read_file" and event.node_id == "read_file": + if mode == "read_file" and lineage == "read_file": return self.answer_patch(graph, reason="safe sandbox read reached a terminal outcome") - if mode == "index_directory" and event.node_id == "list_directory" and event.kind == "task_succeeded": + if mode == "index_directory" and lineage == "list_directory" and event.kind == "task_succeeded": paths = event.payload.get("paths", []) if not paths: return self.answer_patch(graph, reason="the directory contained no matching files") tasks = tuple(TaskSpec(f"index_{i + 1}", "index_file", {"path": path}) for i, path in enumerate(paths)) return GraphPatch(add=tasks, - connect=tuple(("list_directory", task.id) for task in tasks), + connect=tuple((source, task.id) for task in tasks), reason="directory outcome discovered concrete files; index them in parallel") - if mode == "index_directory" and event.node_id and event.node_id.startswith("index_"): + if mode == "index_directory" and lineage.startswith("index_"): work = [node for node_id, node in graph.nodes.items() if node_id.startswith("index_")] if work and all(node["state"] in {"succeeded", "failed", "cancelled"} for node in work): return self.answer_patch(graph, reason="all discovered documents reached a terminal state") - if mode == "search_fetch" and event.node_id == "search" and event.kind == "task_succeeded": + if mode == "search_fetch" and lineage == "search" and event.kind == "task_succeeded": hits = event.payload.get("hits", [])[:3] if not hits: return self.answer_patch(graph, reason="search returned no URLs; explain the failure") tasks = tuple(TaskSpec(f"fetch_{i + 1}", "fetch_url", {"url": hit["url"]}) for i, hit in enumerate(hits) if hit.get("url")) - return GraphPatch(add=tasks, connect=tuple(("search", task.id) for task in tasks), + return GraphPatch(add=tasks, connect=tuple((source, task.id) for task in tasks), reason="search outcome discovered concrete pages; fetch them in parallel") - if event.node_id == "recall": + if lineage == "recall": if mode == "index_file" and "distill" not in graph.nodes: return GraphPatch(add=(TaskSpec("distill", "distiller", {"query": prompt}, {"agent": "paper_distiller"}),), - connect=(("recall", "distill"),), + connect=((source, "distill"),), reason="retrieved paper evidence is ready for extraction") return self.answer_patch(graph, reason="authorized retrieval completed") - if event.node_id == "distill" and mode != "structured_population": + if lineage == "distill" and mode != "structured_population": return self.answer_patch(graph, reason="specialist synthesis completed") - if event.node_id and event.node_id.startswith(("fetch_", "search_")): + if lineage.startswith(("fetch_", "search_")): work = [node for node_id, node in graph.nodes.items() if node_id.startswith(("fetch_", "search_"))] if work and all(node["state"] in {"succeeded", "failed", "cancelled"} for node in work): if mode in {"fetch", "search_fetch", "parallel_search", "structured_population"} and "distill" not in graph.nodes: - parents = tuple(node_id for node_id in graph.nodes if node_id.startswith("search_")) + # Only *succeeded* research may parent the distiller. + # Connecting it to a failed node made `distill` + # permanently unready and the run returned an empty + # answer; with every attempt spent there is nothing + # to distil, so explain the failure instead. + parents = tuple(node_id for node_id in graph.nodes if node_id.startswith("search_") + and graph.nodes[node_id]["state"] == "succeeded") if not parents: - parents = tuple(node_id for node_id in graph.nodes if node_id.startswith("fetch_")) - return GraphPatch(add=(TaskSpec("distill", "distiller", {"query": prompt}, {"agent": "distiller"}),), - connect=tuple((parent, "distill") for parent in parents), - reason="research evidence landed; specialist synthesis can begin") + parents = tuple(node_id for node_id in graph.nodes if node_id.startswith("fetch_") + and graph.nodes[node_id]["state"] == "succeeded") + if parents: + return GraphPatch(add=(TaskSpec("distill", "distiller", {"query": prompt}, {"agent": "distiller"}),), + connect=tuple((parent, "distill") for parent in parents), + reason="research evidence landed; specialist synthesis can begin") + return self.answer_patch( + graph, reason="every research attempt failed; report the failure rather than stalling") return self.answer_patch(graph, reason="the current research frontier has landed") - if mode == "structured_population" and event.node_id == "distill": + if mode == "structured_population" and lineage == "distill": return GraphPatch(add=(TaskSpec("validate", "coder_validator", {"query": prompt}, {"agent": "structured_validator"}),), - connect=(("distill", "validate"),), reason="validate city structure before answering") - if mode == "structured_population" and event.node_id == "validate": + connect=((source, "validate"),), reason="validate city structure before answering") + if mode == "structured_population" and lineage == "validate": return self.answer_patch(graph, reason="structured population fields validated") if event.kind == "task_failed" and event.node_id != "answer": work = [node for node_id, node in graph.nodes.items() if node_id != "answer"] @@ -244,14 +325,29 @@ async def recall(task: TaskSpec) -> dict[str, Any]: return {"hits": [{"id": hit.id, "kind": hit.kind.value, "text": hit.text, "sources": [source.uri for source in hit.sources]} for hit in hits]} + def _lineage_result(snapshot, root: str) -> dict[str, Any]: + """The successful result for a lineage, whichever attempt produced it.""" + for node_id, node in sorted(snapshot.nodes.items()): + metadata = node.get("metadata") or {} + if (node_id == root or metadata.get("retry_root") == root) \ + and node["state"] == "succeeded" and node.get("result"): + return node["result"] + return {} + async def answer(_: TaskSpec) -> dict[str, Any]: snapshot = runtime.graph.snapshot(run_id) evidence: list[dict[str, Any]] = [] - recall_result = snapshot.nodes.get("recall", {}).get("result") or {} + recall_result = _lineage_result(snapshot, "recall") evidence.extend(recall_result.get("hits", [])) - remembered = snapshot.nodes.get("remember", {}).get("result", {}).get("fact") + remembered = _lineage_result(snapshot, "remember").get("fact") if remembered: evidence = [remembered, *evidence] + # A retried lineage would otherwise contribute the same error once + # per attempt. Report only the last attempt, and say how many there + # were, so the answer never hides that the work was re-tried. + superseded_attempts = {str((node.get("metadata") or {}).get("retry_of")) + for node in snapshot.nodes.values() + if (node.get("metadata") or {}).get("retry_of")} for node_id, node in snapshot.nodes.items(): result = node.get("result") or {} if node["skill"] == "fetch_url" and result.get("text"): @@ -272,9 +368,14 @@ async def answer(_: TaskSpec) -> dict[str, Any]: elif node["skill"] == "create_reminder" and result.get("artifacts"): evidence.append({"text": "Calendar reminders created: " + ", ".join(result["artifacts"]), "sources": result["artifacts"], "kind": "calendar_artifact"}) - elif node["state"] == "failed": - evidence.append({"text": result.get("error", "task failed"), - "sources": [f"graph://{run_id}/{node_id}"], "kind": "failure"}) + elif node["state"] == "failed" and node_id not in superseded_attempts: + metadata = node.get("metadata") or {} + attempts = int(metadata.get("retry_attempt", 1)) + detail = result.get("error", "task failed") + if attempts > 1: + detail = f"{detail} (failed on all {attempts} attempts)" + evidence.append({"text": detail, "sources": [f"graph://{run_id}/{node_id}"], + "kind": "failure"}) # Keep the prompt bounded without hiding which source supplied a claim. bounded, used = [], 0 for item in evidence: diff --git a/tests/test_retry_policy_adversarial.py b/tests/test_retry_policy_adversarial.py new file mode 100644 index 0000000..83ef0ea --- /dev/null +++ b/tests/test_retry_policy_adversarial.py @@ -0,0 +1,293 @@ +"""Part 3: attacks aimed squarely at the retry policy this branch added. + +Each test names the attack, shows why the obvious implementation loses, and +then shows the shipped implementation refusing it. The naive classifier below +is not a straw man - it is the message-matching version this feature started +as, kept here so the regression is executable rather than described. +""" + +from __future__ import annotations + +import asyncio +import re + +import pytest + +import s13code.routes as agent_route +from s13code.core.live_graph import ( + FailureClass, + GraphMutationError, + GraphPatch, + GraphStore, + LiveGraphExecutor, + NodeState, + RetryRequest, + TaskSpec, + classify_failure, +) + +# The pre-fix implementation: look for retry-ish words anywhere in the message. +_NAIVE_TRANSIENT = re.compile( + r"timeout|timed out|connection|reset|refused|unavailable|429|502|503|504", re.IGNORECASE) + + +def naive_classify(error: str) -> FailureClass: + return FailureClass.TRANSIENT if _NAIVE_TRANSIENT.search(error or "") else FailureClass.PERMANENT + + +# -------------------------------------------------------------------------- +# Attack 1 - a permanent failure wearing transient words +# -------------------------------------------------------------------------- +# +# The attacker controls a path, a URL or a filename, and therefore controls +# part of the failure message. Every retry-worthy word can be smuggled into an +# error that is in fact a refused sandbox escape. + +HOSTILE_PATH = "/var/log/connection-reset/504-timeout-service-unavailable.txt" +HOSTILE_ERROR = f"PermissionError: path escapes S13_SANDBOX_ROOT: {HOSTILE_PATH}" + + +def test_before_the_fix_a_sandbox_escape_reads_as_retryable(): + """The failure this branch exists to prevent, still reproducible on demand.""" + assert naive_classify(HOSTILE_ERROR) is FailureClass.TRANSIENT + + +def test_after_the_fix_the_exception_type_decides_and_the_message_cannot_override_it(): + verdict = classify_failure(HOSTILE_ERROR) + assert verdict.failure_class is FailureClass.PERMANENT + assert verdict.error_type == "PermissionError" + assert "permanent failure type" in verdict.basis + + +@pytest.mark.parametrize("error", [ + "FileNotFoundError: papers/504-gateway-timeout-connection-reset.txt", + "ValueError: could not parse 'retry after timeout' from the connection header", + "PermissionError: path escapes S13_SANDBOX_ROOT: /etc/timeout/service-unavailable", + "HTTPStatusError: Client error '404 Not Found' for url 'https://x.test/503-unavailable'", +]) +def test_a_family_of_disguises_is_refused(error): + assert naive_classify(error) is FailureClass.TRANSIENT, "the attack must actually fool the old code" + assert classify_failure(error).failure_class is FailureClass.PERMANENT + + +@pytest.mark.asyncio +async def test_the_graph_does_not_re_attempt_a_refused_sandbox_escape(tmp_path): + """End to end: the disguised failure must be executed exactly once.""" + calls: list[str] = [] + + async def read_file(task): + calls.append(task.input["path"]) + raise PermissionError(f"path escapes S13_SANDBOX_ROOT: {task.input['path']}") + + def plan(graph, event): + if event.kind == "run_started": + return GraphPatch(add=(TaskSpec("read", "read_file", {"path": HOSTILE_PATH}),), reason="first") + if event.kind == "task_failed": + verdict = classify_failure(event.payload["error"]) + if verdict.retryable: + return GraphPatch(retry=(RetryRequest(event.node_id),), reason="retry") + return GraphPatch(finish=True, reason="permanent failure, not retried") + return GraphPatch() + + class Planner: + async def plan(self, graph, event): + return plan(graph, event) + + store = GraphStore(tmp_path / "graph.db") + report = await LiveGraphExecutor(store, Planner(), {"read_file": read_file}).run("escape") + + assert calls == [HOSTILE_PATH], "a refused sandbox path was attempted more than once" + assert report.retried == () + assert not any(event.kind == "task_retry_scheduled" for event in store.events("escape")) + + +# -------------------------------------------------------------------------- +# Attack 2 - a duplicate or replayed planner decision +# -------------------------------------------------------------------------- +# +# The journal is replayed after a crash, or the same decision is delivered +# twice. One failure must still buy exactly one attempt, not one per delivery. + +def _failed_graph(tmp_path): + store = GraphStore(tmp_path / "graph.db", max_attempts=3) + store.start("dupe") + store.apply_patch("dupe", GraphPatch(add=(TaskSpec("call", "worker"),)), trigger_event=1) + store.mark_running("dupe", [TaskSpec("call", "worker")]) + failure = store.record_outcome("dupe", "call", False, {"error": "ConnectTimeout: timed out"}) + return store, failure + + +def test_replaying_the_same_trigger_event_creates_exactly_one_attempt(tmp_path): + store, failure = _failed_graph(tmp_path) + patch = GraphPatch(retry=(RetryRequest("call"),), reason="transient") + + assert store.apply_patch("dupe", patch, trigger_event=failure.sequence) is True + # The crash-recovery path replays the same event; it must be a no-op. + assert store.apply_patch("dupe", patch, trigger_event=failure.sequence) is False + + retries = [nid for nid, node in store.snapshot("dupe").nodes.items() + if node["metadata"].get("retry_of")] + assert retries == ["call__retry2"] + scheduled = [e for e in store.events("dupe") if e.kind == "task_retry_scheduled"] + assert len(scheduled) == 1 + store.close() + + +def test_a_second_distinct_decision_cannot_fan_one_failure_into_two_attempts(tmp_path): + """A duplicate arriving under a *new* event id is the harder case.""" + store, failure = _failed_graph(tmp_path) + store.apply_patch("dupe", GraphPatch(retry=(RetryRequest("call"),)), trigger_event=failure.sequence) + + with pytest.raises(GraphMutationError, match="has already been retried"): + store.apply_patch("dupe", GraphPatch(retry=(RetryRequest("call"),)), trigger_event=failure.sequence + 99) + + assert sum(1 for node in store.snapshot("dupe").nodes.values() + if node["metadata"].get("retry_of")) == 1 + store.close() + + +def test_one_patch_cannot_retry_the_same_node_twice(tmp_path): + store, failure = _failed_graph(tmp_path) + with pytest.raises(GraphMutationError, match="has already been retried"): + store.apply_patch("dupe", GraphPatch(retry=(RetryRequest("call"), RetryRequest("call"))), + trigger_event=failure.sequence) + # The failed transaction left no half-built attempt behind. + assert list(store.snapshot("dupe").nodes) == ["call"] + store.close() + + +@pytest.mark.asyncio +async def test_crash_recovery_replay_does_not_duplicate_a_scheduled_attempt(tmp_path): + """Kill the process between the failure and its patch, then resume.""" + store = GraphStore(tmp_path / "graph.db", max_attempts=3) + store.start("crash") + store.apply_patch("crash", GraphPatch(add=(TaskSpec("call", "worker"),)), trigger_event=1) + store.mark_running("crash", [TaskSpec("call", "worker")]) + store.record_outcome("crash", "call", False, {"error": "ConnectTimeout: timed out"}) + store.close() # the patch for that failure was never committed + + reopened = GraphStore(tmp_path / "graph.db", max_attempts=3) + seen: list[int] = [] + + class Planner: + async def plan(self, graph, event): + seen.append(event.sequence) + if event.kind == "task_failed": + return GraphPatch(retry=(RetryRequest(event.node_id),), reason="replayed decision") + return GraphPatch(finish=True, reason="done") + + async def worker(task): + return {"ok": True} + + await LiveGraphExecutor(reopened, Planner(), {"worker": worker}).run("crash", resume=True) + nodes = reopened.snapshot("crash").nodes + assert sorted(nodes) == ["call", "call__retry2"] + assert nodes["call__retry2"]["state"] == "succeeded" + reopened.close() + + +# -------------------------------------------------------------------------- +# Attack 3 - a retry result arriving after cancellation +# -------------------------------------------------------------------------- + +@pytest.mark.asyncio +async def test_a_cancelled_retry_leaks_no_late_result_into_the_graph(tmp_path): + """A sibling finishes the graph while an attempt is still backing off.""" + leaked: list[str] = [] + sibling_may_finish = asyncio.Event() + + async def flaky(task): + if task.metadata.get("retry_of"): + leaked.append(task.id) # must never run: the graph cancelled it + return {"stale": True} + raise TimeoutError("timed out") + + async def sibling(task): + await sibling_may_finish.wait() + return {"answer": "good enough without the flaky branch"} + + class Planner: + async def plan(self, graph, event): + if event.kind == "run_started": + return GraphPatch(add=(TaskSpec("flaky", "flaky"), TaskSpec("sibling", "sibling")), + reason="race two approaches") + if event.kind == "task_failed" and event.node_id == "flaky": + # A long backoff: the sibling will win while this one waits. + sibling_may_finish.set() + return GraphPatch(retry=(RetryRequest("flaky", delay_seconds=30),), reason="retry later") + if event.node_id == "sibling": + return GraphPatch(finish=True, reason="sibling answered; abandon the flaky branch") + return GraphPatch() + + store = GraphStore(tmp_path / "graph.db", max_attempts=3) + report = await asyncio.wait_for( + LiveGraphExecutor(store, Planner(), {"flaky": flaky, "sibling": sibling}).run("race"), timeout=10) + + assert leaked == [], "a cancelled attempt produced work after the graph moved on" + nodes = store.snapshot("race").nodes + assert nodes["flaky__retry2"]["state"] == NodeState.CANCELLED + assert nodes["flaky__retry2"]["result"] is None + assert report.finished + kinds = {(event.kind, event.node_id) for event in store.events("race")} + assert ("task_succeeded", "flaky__retry2") not in kinds + assert ("task_cancelled", "flaky__retry2") in kinds + + +def test_the_store_refuses_an_outcome_for_a_cancelled_attempt(tmp_path): + """The durable backstop, independent of any executor-side race guard.""" + store = GraphStore(tmp_path / "graph.db", max_attempts=3) + store.start("late") + store.apply_patch("late", GraphPatch(add=(TaskSpec("call", "worker"),)), trigger_event=1) + store.mark_running("late", [TaskSpec("call", "worker")]) + store.record_outcome("late", "call", False, {"error": "ConnectTimeout: timed out"}) + store.apply_patch("late", GraphPatch(retry=(RetryRequest("call"),)), trigger_event=2) + store.mark_running("late", [TaskSpec("call__retry2", "worker")]) + store.apply_patch("late", GraphPatch(cancel=("call__retry2",), reason="no longer needed"), trigger_event=3) + + with pytest.raises(GraphMutationError, match="not running"): + store.record_outcome("late", "call__retry2", True, {"stale": "result"}) + assert store.snapshot("late").nodes["call__retry2"]["result"] is None + store.close() + + +# -------------------------------------------------------------------------- +# The user-visible regression the floor exposed +# -------------------------------------------------------------------------- + +def test_an_exhausted_research_branch_answers_instead_of_returning_nothing(app_client, monkeypatch): + """Before this branch, a transient fetch failure returned an empty answer. + + The floor's `shannon` case ended with `fetch_1: failed`, `distill: pending` + and `answer: ''` - `ready()` never releases a child whose parent failed, so + the executor simply ran out of work. + """ + async def fake_gateway(_app, prompt: str, _system: str): + assert "Authorized memory evidence" in prompt + return {"text": "I could not fetch the page: the connection timed out on every attempt.", + "provider": "synthetic", "model": "synthetic"} + + monkeypatch.setattr(agent_route, "gateway_text_llm", fake_gateway) + + attempts = {"count": 0} + + async def always_times_out(url: str, **_kwargs): + attempts["count"] += 1 + raise TimeoutError("timed out") + + monkeypatch.setattr("s13code.runtime.fetch_url", always_times_out) + + response = app_client.post("/v1/agent/runs", json={ + "tenant_id": "synthetic-course", "project_id": "retry", "user_id": "student-synthetic", + "prompt": "Fetch https://example.test/page and tell me what it says."}) + assert response.status_code == 200 + body = response.json() + + assert body["status"] == "completed" + assert body["answer"], "the run must not return an empty answer" + assert attempts["count"] == 3, "the transient failure should be attempted up to the policy ceiling" + + states = {nid: node["state"] for nid, node in body["graph"]["nodes"].items()} + assert states["answer"] == "succeeded" + assert "pending" not in states.values(), "no node may be stranded in pending" + kinds = [event["kind"] for event in body["events"]] + assert kinds.count("task_retry_scheduled") == 2 From ea8bac33a9c39e3a70f9d395410984399a9d37ba Mon Sep 17 00:00:00 2001 From: Aaryan Date: Mon, 17 Aug 2026 10:39:52 +0000 Subject: [PATCH 2/3] Report failure instead of stalling, and document the retry extension - answer_with_evidence checks node failure before the per-skill evidence arms. A failed researcher previously matched the researcher arm, produced no evidence at all, and the run answered 'no authorized memory matched' instead of saying the research had failed. - scripts/repro_retry.py: end-to-end proof against a local origin that is 503 twice then 200, and a second that is 404 forever. - scripts/offline_gateway_model.py: deterministic Ollama-compatible model tier so the proof reproduces with no GPU and no provider key. glc_v3 and S13Code are unmodified above it. - README: the seven required items for the extension. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01FvXm9XSmn5E9fdXGypJKAZ --- .gitignore | 1 + README.md | 267 +++++++++++++ s13code/core/live_graph/retry.py | 9 +- .../live_graph/tests/test_retry_policy.py | 5 + s13code/runtime.py | 25 +- scripts/offline_gateway_model.py | 353 ++++++++++++++++++ scripts/repro_retry.py | 136 +++++++ 7 files changed, 786 insertions(+), 10 deletions(-) create mode 100644 scripts/offline_gateway_model.py create mode 100644 scripts/repro_retry.py diff --git a/.gitignore b/.gitignore index b112038..d222d2d 100644 --- a/.gitignore +++ b/.gitignore @@ -14,4 +14,5 @@ htmlcov/ benchmark.json benchmark.md a2a-proof.json +retry-proof.json .DS_Store diff --git a/README.md b/README.md index 7cf7872..9aa06a2 100644 --- a/README.md +++ b/README.md @@ -124,6 +124,273 @@ Add one subsection to this README in the same pull request. It must contain: Do not commit `.env`, credentials, personal memory, generated databases, unrestricted local paths, benchmark output containing private data, or provider responses containing secrets. Use synthetic identities in every proof. +## Extension: a live graph that tells a timeout apart from a 404 + +### 1. The user-visible capability + +A single flaky network call no longer costs the user their answer. Before this +change a failed node was terminal, and because `apply_patch` refuses to re-add +an existing task id the planner had no way to say "that timeout deserves +another attempt, that 404 does not". On the shipped non-browser benchmark that +gap was not theoretical: **4 of the 14 cases returned a completely empty +answer**, because the planner attached the next stage to a node that had +failed, and `GraphStore.ready()` only releases a node once *every* parent has +succeeded. The child sat in `pending` forever, the executor ran out of work, +and the run returned silently with nothing. The live graph now classifies +every failure before it plans around it: a transient failure is re-attempted +as a **new node** with bounded exponential backoff, a permanent one is never +re-attempted, and either way the user gets either a grounded answer or an +explicit explanation of what failed. Across the same benchmark the result goes +from **10/14 completed with 4 empty answers and 4 stranded nodes** to +**14/14 completed, 0 empty answers, 0 stranded nodes**. + +Classification is type-first by design, and that ordering is the +security-relevant part: a known permanent exception type short-circuits before +any message inspection, so an attacker-controlled path or URL cannot smuggle +the word "timeout" into a `PermissionError` and buy a retry of a refused +sandbox escape. A retry never edits a node in place — it creates a successor +that inherits the failed node's parents and adopts its children — so every +attempt keeps its own state, its own recorded error and its own place in the +journal. + +### 2. The exact API request + +```bash +curl -s http://127.0.0.1:8113/v1/agent/runs \ + -H 'Content-Type: application/json' \ + -d '{ + "tenant_id": "synthetic-course", + "project_id": "retry-proof", + "user_id": "student-synthetic", + "agent_id": "assistant", + "prompt": "Fetch http://127.0.0.1:8231/report and tell me what it says." + }' +``` + +`http://127.0.0.1:8231/report` is a local origin started by +`scripts/repro_retry.py` that answers `503 Service Unavailable` to its first +two requests and serves a real page on the third. The permanent counterpart +uses the same request against `http://127.0.0.1:8231/always-missing`, which +answers `404` forever. + +### 3. The graph and the ordered event trace + +Transient origin, run `run-910f4ce055b0`: + +``` +544 run_started - +545 graph_patched - first frontier selected for fetch +546 task_started fetch_1 +547 task_failed fetch_1 HTTPStatusError: Server error '503 Service Unavailable' +548 task_retry_scheduled fetch_1__retry2 attempt 2/3 after 0.05s — HTTP 503 is retryable +549 graph_patched - fetch_1 failed transiently (HTTP 503 is retryable); attempt 2 of 3 +550 task_started fetch_1__retry2 +551 task_failed fetch_1__retry2 HTTPStatusError: Server error '503 Service Unavailable' +552 task_retry_scheduled fetch_1__retry3 attempt 3/3 after 0.1s — HTTP 503 is retryable +553 graph_patched - fetch_1__retry2 failed transiently (HTTP 503 is retryable); attempt 3 of 3 +554 task_started fetch_1__retry3 +555 task_succeeded fetch_1__retry3 +556 graph_patched - research evidence landed; specialist synthesis can begin +557 task_started distill +558 task_succeeded distill +559 graph_patched - specialist synthesis completed +560 task_started answer +561 task_succeeded answer +562 graph_patched - grounded answer produced +``` + +Final edges: `[["distill","answer"],["fetch_1__retry3","answer"],["fetch_1__retry3","distill"]]` + +Two invariants are readable directly from that trace. **No future node exists +before its inputs**: `fetch_1__retry2` appears for the first time at sequence +548, *after* the failure at 547, and `distill` is created at 556 only once a +fetch attempt has actually succeeded — it is wired to `fetch_1__retry3`, never +to the two attempts that failed. And **the failed attempts are still there**: +`fetch_1` and `fetch_1__retry2` remain `failed` with their recorded errors; +nothing was rewritten to make the run look clean. + +Permanent origin, run `run-c103337da291` — the whole graph, with no retry +event anywhere in it: + +``` +563 run_started - +564 graph_patched - first frontier selected for fetch +565 task_started fetch_1 +566 task_failed fetch_1 HTTPStatusError: Client error '404 Not Found' +567 graph_patched - every research attempt failed; report the failure rather than stalling +568 task_started answer +569 task_succeeded answer +570 graph_patched - grounded answer produced +``` + +### 4. The actual final result + +Transient origin — the origin's own counter confirms it received exactly three +requests for `/report`: + +``` +Based on the authorized evidence: +- Quarterly reliability report. [source: graph://run-910f4ce055b0/distill] +- The retry budget absorbed three upstream incidents this quarter. [source: graph://run-910f4ce055b0/distill] +- Quarterly reliability report. [source: http://127.0.0.1:8231/report] +- The retry budget absorbed three upstream incidents this quarter. [source: http://127.0.0.1:8231/report] + +Sources consulted: graph://run-910f4ce055b0/distill, http://127.0.0.1:8231/report. +``` + +Permanent origin — one attempt, no retry, and the failure is reported rather +than swallowed: + +``` +The authorized evidence did not contain a passage matching this request. + +One step did not succeed: The fetch_url step failed: HTTPStatusError: Client error +'404 Not Found' for url 'http://127.0.0.1:8231/always-missing' +[source: graph://run-c103337da291/fetch_1] + +No usable source was available. +``` + +### 5. Evidence and provider/agent assignments + +| node | agent | skill | state | attempt | provider/model | +|---|---|---|---|---|---| +| `fetch_1` | fetch_url | `fetch_url` | `failed` | 1 | — (local skill) | +| `fetch_1__retry2` | fetch_url | `fetch_url` | `failed` | 2 | — (local skill) | +| `fetch_1__retry3` | fetch_url | `fetch_url` | `succeeded` | 3 | — (local skill) | +| `distill` | distiller | `distiller` | `succeeded` | 1 | `ollama` / `s13-offline-extractive:1.0` | +| `answer` | answer_with_evidence | `answer_with_evidence` | `succeeded` | 1 | `ollama` / `s13-offline-extractive:1.0` | + +Evidence in the final answer is `kind: web_page` from +`http://127.0.0.1:8231/report` (the successful third attempt) plus +`kind: role_output` from `graph://run-910f4ce055b0/distill`. Both model calls +crossed real HTTP to `glc_v3` on `127.0.0.1:8111`, which selected the `ollama` +provider; `S13Code` never saw a credential. Fetch attempts are local skills and +carry no provider. + +The provider above is an **offline deterministic model**, not a frontier one. +`scripts/offline_gateway_model.py` serves Ollama's `/api/chat` and `/api/embed` +with an extractive engine so that this section reproduces byte-for-byte with no +GPU and no provider key — every layer above the model tier (`glc_v3` routing, +policy and accounting; the graph, journal, memory and A2A in `S13Code`) is the +real one. Point `OLLAMA_URL` at a real Ollama, or give `glc_v3` a provider key, +and the identical commands produce the identical graph with better prose. + +Retry is policy in the planner and mechanism in the store. `GraphStore` +independently enforces the attempt ceiling and the one-attempt-per-failure +rule, so a planner — including an LLM planner emitting `retry` in a +`GraphPatch` — cannot loop the budget away. + +### 6. The adversarial failure and its fix + +**The attack.** A permanent failure wearing transient words. The user controls +part of a path, therefore part of the error message, and every retry-worthy +keyword can be smuggled into a failure that is in fact a *refused sandbox +escape*: + +``` +PermissionError: path escapes S13_SANDBOX_ROOT: /var/log/connection-reset/504-timeout-service-unavailable.txt +``` + +**The failure before the fix.** The first cut of this feature classified by +searching the message for retry-ish words. That implementation is kept in +`tests/test_retry_policy_adversarial.py` as `naive_classify`, so the regression +is executable rather than described: + +```python +def test_before_the_fix_a_sandbox_escape_reads_as_retryable(): + assert naive_classify(HOSTILE_ERROR) is FailureClass.TRANSIENT +``` + +It returns `TRANSIENT`, so the graph would re-attempt a path the sandbox had +already refused — up to the ceiling, three times, spending budget on an access +the system exists to deny. + +**The fix.** `classify_failure` decides on the **exception type first**. A type +in `PERMANENT_TYPES` short-circuits before any message inspection, so no +message can promote it. Message hints are consulted only for types neither +table recognises, a permanent hint outranks a transient one, and an +unrecognised failure defaults to permanent — a graph that retries what it does +not understand is a graph that burns budget on real bugs. + +```python +def test_after_the_fix_the_exception_type_decides_and_the_message_cannot_override_it(): + verdict = classify_failure(HOSTILE_ERROR) + assert verdict.failure_class is FailureClass.PERMANENT + assert verdict.error_type == "PermissionError" +``` + +`test_the_graph_does_not_re_attempt_a_refused_sandbox_escape` runs the same +attack end to end through the executor and asserts the hostile path is +attempted **exactly once**. A parametrised family of disguises is checked the +same way: each must genuinely fool `naive_classify` *and* be refused by +`classify_failure`. A related case falls out of the same rule — `ProxyError` is +a transport type, but `ProxyError: 403 Forbidden` is permanent, because a proxy +that answers 403 will answer 403 again. + +Two further attacks are covered in the same file. A **duplicate or replayed +planner decision** must not fan one failure into several concurrent attempts: +replaying the same `trigger_event` is a no-op via the `patches` table, a +duplicate arriving under a *new* event id is rejected with +`task ... has already been retried`, and a crash between a failure and its +patch replays into exactly one attempt on resume. A **result arriving after +cancellation** must not enter the graph: when a sibling finishes the run while +an attempt is still backing off, the attempt is cancelled during its sleep, its +worker never executes, no `task_succeeded` is journalled for it, and +`record_outcome` refuses outright for a node that is no longer `running`. + +### 7. Reproducing this from a fresh checkout + +Unzip or clone `glc_v3`, `S13Code` and `S13Proof` beside one another. No +provider key and no GPU are required for any step below. + +```bash +# 1. tests and lint (44 upstream + 32 added) +cd S13Code +uv sync +uv run ruff check . +uv run pytest -q # 76 passed + +cd ../S13Proof && uv sync && uv run pytest -q +uv run python run_a2a_proof.py --output a2a-proof.json + +# 2. the offline model tier, so glc_v3's real OllamaProvider has something to call +cd ../S13Code +uv run python scripts/offline_gateway_model.py --port 11434 & + +# 3. glc_v3, unmodified, routing to that local model +cd ../glc_v3 && uv sync +OLLAMA_MODEL=s13-offline-extractive:1.0 OLLAMA_URL=http://127.0.0.1:11434 \ +LLM_ORDER=ollama uv run uvicorn glc.main:app --host 127.0.0.1 --port 8111 & + +# 4. S13Code +cd ../S13Code +export GLC_BASE_URL=http://127.0.0.1:8111 +export S13_GATEWAY_PROVIDER=ollama +export S13_SANDBOX_ROOT="$PWD/sandbox" +export S13_CHUNK_MODEL=s13-offline-extractive:1.0 +uv run uvicorn s13code.main:app --host 127.0.0.1 --port 8113 & + +curl http://127.0.0.1:8113/healthz && curl http://127.0.0.1:8113/readyz + +# 5. the retry proof: both scenarios, graph, ordered trace, agents, answers +uv run python scripts/repro_retry.py --output retry-proof.json + +# 6. the whole non-browser benchmark +cd ../S13Proof +uv run python run_benchmark.py --base-url http://127.0.0.1:8113 --output benchmark.md +``` + +To see the *before* state, run step 6 against a checkout of the commit +preceding this branch: `shannon`, `tokyo_weather`, `populations` and +`structured_growth` come back `failed` with an empty answer and one node +stranded in `pending`. On this branch all fourteen cases complete with a +non-empty answer. + +Steps 5 and 6 write `retry-proof.json`, `benchmark.json` and `benchmark.md`; +all three are already covered by `.gitignore` and are not committed. Every +identity used above is synthetic. + ## License MIT. See `LICENSE`. diff --git a/s13code/core/live_graph/retry.py b/s13code/core/live_graph/retry.py index d64d377..e2bd364 100644 --- a/s13code/core/live_graph/retry.py +++ b/s13code/core/live_graph/retry.py @@ -75,7 +75,14 @@ class FailureClass(StrEnum): r"|permission denied|invalid|malformed|unsupported|not permitted)\b", re.IGNORECASE, ) -_HTTP_STATUS = re.compile(r"\b(?:status(?:\s*code)?\s*[:=]?\s*|HTTP\s+|')(\d{3})\b") +# Reads the status out of the shapes libraries actually produce, e.g. httpx's +# ``Server error '503 Service Unavailable' for url ...`` and ``HTTP 429``. +# The leading context is required: a bare three-digit run inside a URL path is +# not a status, which is what stops ``/503-service-unavailable`` in an attacker +# controlled URL from re-classifying a 404. +_HTTP_STATUS = re.compile( + r"(?:status(?:\s*code)?\s*[:=]?\s*|\bHTTP[/ ]\S*?\s*|['\"(]" + r"|\A[A-Za-z_][A-Za-z0-9_.]*:\s+)(\d{3})\b", re.IGNORECASE) # 408 request timeout, 425 too early, 429 too many requests, plus 5xx. _TRANSIENT_STATUS = frozenset({408, 425, 429, 500, 502, 503, 504, 507, 509, 598, 599}) diff --git a/s13code/core/live_graph/tests/test_retry_policy.py b/s13code/core/live_graph/tests/test_retry_policy.py index 82cc862..f9b2d0f 100644 --- a/s13code/core/live_graph/tests/test_retry_policy.py +++ b/s13code/core/live_graph/tests/test_retry_policy.py @@ -75,6 +75,11 @@ def test_http_status_outranks_the_transport_class_that_carried_it(): "'https://example.test/x'").failure_class is FailureClass.PERMANENT assert classify_failure("HTTPStatusError: Server error '503 Service Unavailable' for url " "'https://example.test/x'").failure_class is FailureClass.TRANSIENT + # ProxyError is a transport type, but a proxy that answers 403 will answer + # 403 again. The status wins over the class that delivered it. + forbidden = classify_failure("ProxyError: 403 Forbidden") + assert forbidden.failure_class is FailureClass.PERMANENT + assert "403" in forbidden.basis def test_unrecognised_failures_default_to_permanent(): diff --git a/s13code/runtime.py b/s13code/runtime.py index adc2663..6385f48 100644 --- a/s13code/runtime.py +++ b/s13code/runtime.py @@ -350,7 +350,22 @@ async def answer(_: TaskSpec) -> dict[str, Any]: if (node.get("metadata") or {}).get("retry_of")} for node_id, node in snapshot.nodes.items(): result = node.get("result") or {} - if node["skill"] == "fetch_url" and result.get("text"): + # Failure is checked first. When it sat at the end of this + # chain an earlier skill arm claimed the node and produced no + # evidence at all, so a run whose every researcher failed + # answered "no authorized memory matched" instead of saying + # that the research had failed. + if node["state"] == "failed": + if node_id in superseded_attempts: + continue + metadata = node.get("metadata") or {} + attempts = int(metadata.get("retry_attempt", 1)) + detail = result.get("error", "task failed") + if attempts > 1: + detail = f"{detail} (failed on all {attempts} attempts)" + evidence.append({"text": f"The {node['skill']} step failed: {detail}", + "sources": [f"graph://{run_id}/{node_id}"], "kind": "failure"}) + elif node["skill"] == "fetch_url" and result.get("text"): evidence.append({"text": result["text"][:12_000], "sources": [result["url"]], "kind": "web_page"}) elif node["skill"] == "web_search": for hit in result.get("hits", []): @@ -368,14 +383,6 @@ async def answer(_: TaskSpec) -> dict[str, Any]: elif node["skill"] == "create_reminder" and result.get("artifacts"): evidence.append({"text": "Calendar reminders created: " + ", ".join(result["artifacts"]), "sources": result["artifacts"], "kind": "calendar_artifact"}) - elif node["state"] == "failed" and node_id not in superseded_attempts: - metadata = node.get("metadata") or {} - attempts = int(metadata.get("retry_attempt", 1)) - detail = result.get("error", "task failed") - if attempts > 1: - detail = f"{detail} (failed on all {attempts} attempts)" - evidence.append({"text": detail, "sources": [f"graph://{run_id}/{node_id}"], - "kind": "failure"}) # Keep the prompt bounded without hiding which source supplied a claim. bounded, used = [], 0 for item in evidence: diff --git a/scripts/offline_gateway_model.py b/scripts/offline_gateway_model.py new file mode 100644 index 0000000..769c5d6 --- /dev/null +++ b/scripts/offline_gateway_model.py @@ -0,0 +1,353 @@ +"""Deterministic, dependency-light stand-in for a local Ollama instance. + +Why this exists +--------------- +`glc_v3` reaches its local model tier through ``OllamaProvider``, which speaks +Ollama's ``/api/chat``. `S13Code` reaches its embedder through Ollama's +``/api/embed``. Both are ordinary HTTP. + +This module serves those two endpoints with a deterministic extractive engine +instead of neural weights. Everything above the model tier - `glc_v3` routing, +policy, audit and cost accounting; `S13Code`'s live graph, journal, scoped +memory, semantic chunking and A2A - runs unmodified and for real. + +The point is byte-reproducible traces: a reader with no GPU and no provider key +replays a run and gets the same graph, the same event ordering and the same +final text. Swap ``OLLAMA_URL`` back to a real Ollama (or set a provider key in +`glc_v3`) and nothing above this file changes. + +Run: python offline_model_server.py --port 11434 +""" + +from __future__ import annotations + +import argparse +import json +import math +import re +from collections import Counter +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from typing import Any + +MODEL_NAME = "s13-offline-extractive:1.0" +EMBED_DIMENSIONS = 256 + +_SENTENCE = re.compile(r"(?<=[.!?])\s+(?=[A-Z(\[])|\n+") +_WORD = re.compile(r"[a-z0-9]+") +_STOP = frozenset("""a an and are as at be been but by for from has have how i if in into is it its of on or +that the their them then there these they this to was were what when where which who why will with you your +about across after all also any because before between both can could do does each did having here more most +not only other our over same should so some such than through under up very we would across say says said +tell give me my""".split()) + + +# --------------------------------------------------------------------------- +# deterministic embeddings +# --------------------------------------------------------------------------- + +def embed(text: str) -> list[float]: + """Hashed bag of unigrams and bigrams, L2 normalised. + + Bigrams matter: pure unigram bags rank "chain of thought" and "thought + chain" identically, which makes the retrieval proofs uninformative. + """ + tokens = [token for token in _WORD.findall(text.lower()) if token not in _STOP] + grams = Counter(tokens) + grams.update(f"{left}_{right}" for left, right in zip(tokens, tokens[1:])) + values = [0.0] * EMBED_DIMENSIONS + for gram, count in grams.items(): + slot = sum((index + 1) * ord(char) for index, char in enumerate(gram)) % EMBED_DIMENSIONS + # Sub-linear term frequency; a word repeated 40 times is not 40x the signal. + values[slot] += 1.0 + math.log(count) + norm = math.sqrt(sum(value * value for value in values)) + return [value / norm for value in values] if norm else values + + +# --------------------------------------------------------------------------- +# extractive generation +# --------------------------------------------------------------------------- + +def _terms(text: str) -> set[str]: + return {token for token in _WORD.findall(text.lower()) if token not in _STOP and len(token) > 2} + + +def _sentences(text: str) -> list[str]: + parts = [part.strip() for part in _SENTENCE.split(text) if part and part.strip()] + return [part for part in parts if len(part.split()) >= 3] + + +def _rank(sentences: list[str], query: set[str], limit: int, *, pad: bool = True) -> list[str]: + """Greedy relevance with redundancy suppression (a small MMR). + + ``pad`` backfills with lead sentences when too few score above zero. Some + real requests - "tell me what it says" - carry no content terms at all, and + lead-based extraction is a better answer than silence. + """ + scored = [] + for position, sentence in enumerate(sentences): + words = _terms(sentence) + if not words: + continue + overlap = len(words & query) / math.sqrt(len(words)) + # Earlier sentences in a passage carry the claim; later ones qualify it. + scored.append((overlap - position * 1e-4, sentence, words)) + scored.sort(key=lambda item: -item[0]) + chosen: list[tuple[str, set[str]]] = [] + for score, sentence, words in scored: + if score <= 0: + break + if any(len(words & taken) / max(1, len(words | taken)) > 0.6 for _, taken in chosen): + continue + chosen.append((sentence, words)) + if len(chosen) >= limit: + break + if pad: + taken = {sentence for sentence, _ in chosen} + for sentence in sentences: + if len(chosen) >= limit: + break + if sentence not in taken: + chosen.append((sentence, _terms(sentence))) + taken.add(sentence) + # Restore document order: an extract reads as prose, not as a ranking. + order = {sentence: index for index, sentence in enumerate(sentences)} + return sorted((sentence for sentence, _ in chosen), key=lambda item: order.get(item, 0)) + + +_ITEM_START = re.compile(r"(?m)^-\s*\[kind:\s*([^\]]+)\]\s*") + + +def _parse_evidence(block: str) -> list[dict[str, str]]: + """Read the answer worker's rendered evidence back into records. + + An item is NOT one line: a document chunk keeps its original newlines, so + the record runs from one ``- [kind: …]`` marker to the next. The trailing + ``[source: …]`` is the last such bracket in that span - chunk bodies of + scraped pages contain Markdown links with their own brackets. + """ + starts = list(_ITEM_START.finditer(block)) + items = [] + for index, match in enumerate(starts): + end = starts[index + 1].start() if index + 1 < len(starts) else len(block) + body = block[match.end():end].rstrip() + source_match = None + for candidate in re.finditer(r"\[source:\s*([^\]]*)\]", body): + source_match = candidate + if source_match is None: + continue + items.append({"kind": match.group(1).strip(), + "text": body[:source_match.start()].strip(), + "source": source_match.group(1).strip()}) + return items + + +def _answer_from_evidence(prompt: str) -> str: + request_match = re.search(r"User request:\n(.*?)\n\nAuthorized memory evidence:\n(.*)", prompt, re.S) + if not request_match: + return _generic(prompt) + question, evidence_block = request_match.group(1).strip(), request_match.group(2) + items = _parse_evidence(evidence_block) + if not items or "(No authorized durable memory matched" in evidence_block: + return ("I could not answer this from authorized durable memory: no scoped record matched the " + "request. Nothing in the supplied evidence supports a claim here.") + + query = _terms(question) + failures = [item for item in items if item["kind"] == "failure"] + usable = [item for item in items if item["kind"] != "failure"] + + lines: list[str] = [] + for item in sorted(usable, key=lambda entry: -len(_terms(entry["text"]) & query))[:5]: + picked = _rank(_sentences(item["text"]), query, 2) or [item["text"][:300]] + for sentence in picked: + # Role outputs are already bulleted; do not bullet them twice. + clean = sentence.lstrip("-• \t").rstrip(".") + if clean: + lines.append(f"- {clean}. [source: {item['source']}]") + if len(lines) >= 6: + break + + parts: list[str] = [] + if lines: + parts.append("Based on the authorized evidence:") + parts.extend(lines) + else: + parts.append("The authorized evidence did not contain a passage matching this request.") + if failures: + parts.append("") + for failure in failures: + parts.append(f"One step did not succeed: {failure['text']} [source: {failure['source']}]") + parts.append("") + parts.append(f"Sources consulted: {', '.join(sorted({item['source'] for item in usable}))}." + if usable else "No usable source was available.") + return "\n".join(parts) + + +def _segment(prompt: str) -> str: + """Markdown segmenter reply: a verbatim suffix, or nothing. + + Lexical cohesion, not headings: score each paragraph gap by how little + vocabulary the two sides share, and cut at the weakest seam if it is weak + enough. The reply is always sliced out of the block, never rewritten, so + Rohan V2's suffix check is a real test rather than a formality. + """ + block_match = re.search(r"---\n(.*)\n---", prompt, re.S) + if not block_match: + return "" + block = block_match.group(1) + offsets = [match.end() for match in re.finditer(r"\n\s*\n", block)] + if not offsets: + return "" + best_offset, best_score = None, 0.0 + for offset in offsets: + left, right = _terms(block[:offset]), _terms(block[offset:]) + if len(left) < 12 or len(right) < 12: + continue + overlap = len(left & right) / len(left | right) + score = 1.0 - overlap + if score > best_score: + best_offset, best_score = offset, score + # 0.86 keeps ordinary topic drift inside one chunk; only a genuine subject + # change clears it. Tuned on the sandbox/papers fixtures. + if best_offset is None or best_score < 0.86: + return "" + return block[best_offset:].strip() + + +def _harvest(value: Any, out: list[str]) -> None: + """Collect human-readable prose out of a nested worker result. + + Role workers are handed the whole upstream result dict. Flattening it to + JSON and ranking that produces sentences full of field names; walking it + for the fields that actually hold text does not. + """ + if isinstance(value, str): + if len(value.split()) >= 4: + out.append(value) + elif isinstance(value, dict): + for key in ("text", "answer", "snippet", "title", "content"): + if isinstance(value.get(key), str): + _harvest(value[key], out) + for key, item in value.items(): + if key not in {"text", "answer", "snippet", "title", "content"} and isinstance(item, (dict, list)): + _harvest(item, out) + elif isinstance(value, list): + for item in value: + _harvest(item, out) + + +def _role_reply(prompt: str, system: str) -> str: + try: + payload = json.loads(prompt) + except (ValueError, TypeError): + return _generic(prompt) + task = payload.get("task") or {} + question = str(task.get("query") or task.get("question") or "") + upstream = payload.get("upstream_evidence") or payload.get("hits") or {} + harvested: list[str] = [] + _harvest(upstream, harvested) + sentences = _sentences("\n".join(harvested)) + # "tell me what it says" has no content terms to match on - its query words + # are all URL noise. Falling back to the lead sentences beats reporting no + # findings when upstream evidence plainly exists. + fragments = _rank(sentences, _terms(question), 4) or sentences[:3] + if "coder_validator" in system: + if not fragments: + return ("VALIDATION FAILED: no upstream values were supplied, so no comparison is possible. " + "Do not report a fastest-growing result.") + return ("VALIDATION: the supplied upstream values were checked for a shared geographic definition, a " + "comparable year and an explicit growth rate.\n" + "\n".join(f"- {item}" for item in fragments)) + if not fragments: + return "No upstream evidence was supplied to this role, so it produced no findings." + return "\n".join(f"- {item}" for item in fragments) + + +def _generic(prompt: str) -> str: + picked = _rank(_sentences(prompt), _terms(prompt), 3) + return "\n".join(picked) if picked else prompt.strip()[:400] or "Acknowledged." + + +def generate(messages: list[dict[str, Any]]) -> str: + system = " ".join(str(message.get("content", "")) for message in messages if message.get("role") == "system") + user = "\n".join(str(message.get("content", "")) for message in messages if message.get("role") == "user") + if "markdown document segmenter" in user.lower() or "markdown document segmenter" in system.lower(): + return _segment(user) + if "Authorized memory evidence:" in user: + return _answer_from_evidence(user) + if "role in a constrained graph" in system or "researcher role" in system or "retriever role" in system: + return _role_reply(user, system) + if "Return only valid GraphPatch JSON" in system: + # The constrained planner must fall back deterministically rather than + # receive a patch this engine is not qualified to invent. + return "not a graph patch" + return _generic(user) + + +# --------------------------------------------------------------------------- +# HTTP surface +# --------------------------------------------------------------------------- + +class Handler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def log_message(self, *_args: Any) -> None: # keep proof output readable + return + + def _send(self, payload: dict[str, Any], status: int = 200) -> None: + body = json.dumps(payload).encode() + self.send_response(status) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def do_GET(self) -> None: # noqa: N802 - stdlib naming + if self.path.startswith("/api/tags"): + self._send({"models": [{"name": MODEL_NAME, "model": MODEL_NAME}, + {"name": "nomic-embed-text", "model": "nomic-embed-text"}]}) + elif self.path in {"/", "/healthz"}: + self._send({"ok": True, "engine": "offline-extractive"}) + else: + self._send({"error": "not found"}, 404) + + def do_POST(self) -> None: # noqa: N802 - stdlib naming + length = int(self.headers.get("Content-Length", "0")) + try: + body = json.loads(self.rfile.read(length) or b"{}") + except ValueError: + self._send({"error": "invalid json"}, 400) + return + path = self.path.split("?")[0] + if path == "/api/chat": + text = generate(body.get("messages", [])) + self._send({"model": body.get("model", MODEL_NAME), "created_at": "1970-01-01T00:00:00Z", + "message": {"role": "assistant", "content": text}, "done": True, + "done_reason": "stop", + "prompt_eval_count": sum(len(str(m.get("content", "")).split()) + for m in body.get("messages", [])), + "eval_count": len(text.split())}) + elif path == "/api/embed": + raw = body.get("input", "") + batch = raw if isinstance(raw, list) else [raw] + self._send({"model": body.get("model", "nomic-embed-text"), + "embeddings": [embed(str(item)) for item in batch]}) + elif path == "/api/embeddings": + self._send({"embedding": embed(str(body.get("prompt", "")))}) + elif path == "/api/generate": + # `release()` posts here to unload a model; there is nothing to unload. + self._send({"model": body.get("model", MODEL_NAME), "response": "", "done": True}) + else: + self._send({"error": f"unsupported path {path}"}, 404) + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--host", default="127.0.0.1") + parser.add_argument("--port", type=int, default=11434) + args = parser.parse_args() + server = ThreadingHTTPServer((args.host, args.port), Handler) + print(f"offline model server on http://{args.host}:{args.port} (model {MODEL_NAME})", flush=True) + server.serve_forever() + + +if __name__ == "__main__": + main() diff --git a/scripts/repro_retry.py b/scripts/repro_retry.py new file mode 100644 index 0000000..5c29547 --- /dev/null +++ b/scripts/repro_retry.py @@ -0,0 +1,136 @@ +"""End-to-end proof of the transient-vs-permanent retry policy. + +Starts a local origin that is deliberately flaky, asks the running `S13Code` +service to fetch it, and prints the graph, the ordered event trace, the +provider/agent assignments and the final answer for two scenarios: + +1. a **transient** origin - 503 for the first two requests, then a real page. + The graph must retry and finish with a grounded answer. +2. a **permanent** origin - 404 forever. The graph must *not* retry, and must + still answer by explaining the failure. + +Nothing here is mocked inside `S13Code`: the requests cross real HTTP, the +graph is the durable one, and the answer comes back through `glc_v3`. + + uv run python scripts/repro_retry.py --base-url http://127.0.0.1:8113 +""" + +from __future__ import annotations + +import argparse +import json +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path + +import httpx + +PAGE = (b"

Quarterly reliability report

" + b"

The retry budget absorbed three upstream incidents this quarter. " + b"Mean time to recovery was four minutes. No incident required manual intervention.

" + b"") + +# Synthetic identities only; nothing here refers to a real tenant or person. +SCOPE = {"tenant_id": "synthetic-course", "project_id": "retry-proof", + "user_id": "student-synthetic", "agent_id": "assistant"} + + +class FlakyHandler(BaseHTTPRequestHandler): + """503 for the first `fail_times` requests to /report, then the page.""" + + fail_times = 2 + seen = 0 + lock = threading.Lock() + + protocol_version = "HTTP/1.1" + + def log_message(self, *_args): + return + + def do_GET(self): # noqa: N802 - stdlib naming + if self.path == "/always-missing": + self.send_response(404) + self.send_header("Content-Length", "0") + self.end_headers() + return + with FlakyHandler.lock: + FlakyHandler.seen += 1 + attempt = FlakyHandler.seen + if attempt <= FlakyHandler.fail_times: + self.send_response(503, "Service Unavailable") + self.send_header("Retry-After", "1") + self.send_header("Content-Length", "0") + self.end_headers() + return + self.send_response(200) + self.send_header("Content-Type", "text/html") + self.send_header("Content-Length", str(len(PAGE))) + self.end_headers() + self.wfile.write(PAGE) + + +def render(title: str, body: dict) -> str: + graph = body.get("graph", {}) + nodes = graph.get("nodes", {}) + lines = [f"## {title}", "", + f"status: `{body.get('status')}` · run: `{body.get('run_id')}`", "", + "| node | agent | skill | state | attempt | provider/model |", + "|---|---|---|---|---|---|"] + for node_id, node in sorted(nodes.items()): + metadata = node.get("metadata") or {} + result = node.get("result") or {} + attempt = metadata.get("retry_attempt", 1) + lines.append(f"| `{node_id}` | {metadata.get('agent', node.get('skill'))} | `{node.get('skill')}` " + f"| `{node.get('state')}` | {attempt} | " + f"{result.get('provider') or '-'}/{result.get('model') or '-'} |") + lines += ["", f"edges: `{json.dumps(graph.get('edges', []))}`", "", "ordered event trace:", ""] + for event in body.get("events", []): + payload = event.get("payload", {}) + detail = "" + if event["kind"] == "task_retry_scheduled": + detail = (f" attempt {payload.get('attempt')}/{payload.get('max_attempts')}" + f" after {payload.get('delay_seconds')}s — {payload.get('reason')}") + elif event["kind"] == "task_failed": + detail = f" {payload.get('error', '')}" + elif event["kind"] == "graph_patched": + detail = f" {payload.get('reason', '')}" + lines.append(f"{event['sequence']:>3} {event['kind']:<22} {event.get('node_id') or '-':<24}{detail}") + lines += ["", "final answer:", "", body.get("answer") or "*(empty)*", ""] + return "\n".join(lines) + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--base-url", default="http://127.0.0.1:8113") + parser.add_argument("--port", type=int, default=8231) + parser.add_argument("--output", type=Path, default=Path("retry-proof.json")) + args = parser.parse_args() + + origin = ThreadingHTTPServer(("127.0.0.1", args.port), FlakyHandler) + thread = threading.Thread(target=origin.serve_forever, daemon=True) + thread.start() + base = f"http://127.0.0.1:{args.port}" + results = {} + try: + with httpx.Client(timeout=300) as client: + transient = client.post(f"{args.base_url}/v1/agent/runs", json={ + **SCOPE, "prompt": f"Fetch {base}/report and tell me what it says."}).json() + results["transient"] = transient + print(render("Transient origin: 503, 503, then 200", transient)) + print(f"\norigin received {FlakyHandler.seen} request(s) for /report\n") + + permanent = client.post(f"{args.base_url}/v1/agent/runs", json={ + **SCOPE, "project_id": "retry-proof-permanent", + "prompt": f"Fetch {base}/always-missing and tell me what it says."}).json() + results["permanent"] = permanent + print(render("Permanent origin: 404 forever", permanent)) + + results["origin_requests_for_report"] = FlakyHandler.seen + args.output.write_text(json.dumps(results, indent=2)) + print(f"\nwrote {args.output}") + finally: + origin.shutdown() + + +if __name__ == "__main__": + main() From a4df962caafa4946d16d7572d0095556da31e8d3 Mon Sep 17 00:00:00 2001 From: Aaryan Date: Tue, 18 Aug 2026 12:09:08 +0000 Subject: [PATCH 3/3] Fix adversarial test's hidden dependency on a live Ollama embedder test_an_exhausted_research_branch_answers_instead_of_returning_nothing uses app_client, which boots a real S13Runtime defaulting to OllamaNomicEmbedder. Every memory write in that test's run - the inbound prompt, the answer episode - tries to embed over real HTTP to localhost:11434. It passed in the original sandbox only because an unrelated local server happened to be running on that port; in a clean environment with nothing on 11434 it fails with urllib.error.URLError: Connection refused. Fix: swap in DeterministicEmbedder before the run, exactly as every other app_client test in this suite already does. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01FvXm9XSmn5E9fdXGypJKAZ --- tests/test_retry_policy_adversarial.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/tests/test_retry_policy_adversarial.py b/tests/test_retry_policy_adversarial.py index 83ef0ea..f1ebe27 100644 --- a/tests/test_retry_policy_adversarial.py +++ b/tests/test_retry_policy_adversarial.py @@ -25,6 +25,7 @@ TaskSpec, classify_failure, ) +from s13code.core.memory.embeddings import DeterministicEmbedder # The pre-fix implementation: look for retry-ish words anywhere in the message. _NAIVE_TRANSIENT = re.compile( @@ -261,6 +262,12 @@ def test_an_exhausted_research_branch_answers_instead_of_returning_nothing(app_c and `answer: ''` - `ready()` never releases a child whose parent failed, so the executor simply ran out of work. """ + # Every memory write in this run - the inbound prompt, the answer episode - + # embeds its text. The default embedder calls a live Ollama; swap in the + # deterministic one so this test needs no local model server, exactly like + # every other app_client test in this suite. + app_client.app.state.s13_runtime.memory.embedder = DeterministicEmbedder(256) + async def fake_gateway(_app, prompt: str, _system: str): assert "Authorized memory evidence" in prompt return {"text": "I could not fetch the page: the connection timed out on every attempt.",