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/.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/__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..e2bd364 --- /dev/null +++ b/s13code/core/live_graph/retry.py @@ -0,0 +1,202 @@ +"""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, +) +# 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}) + + +@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..f9b2d0f --- /dev/null +++ b/s13code/core/live_graph/tests/test_retry_policy.py @@ -0,0 +1,244 @@ +"""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 + # 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(): + 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..6385f48 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,17 +325,47 @@ 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"): + # 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", []): @@ -272,9 +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": - evidence.append({"text": result.get("error", "task failed"), - "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() diff --git a/tests/test_retry_policy_adversarial.py b/tests/test_retry_policy_adversarial.py new file mode 100644 index 0000000..f1ebe27 --- /dev/null +++ b/tests/test_retry_policy_adversarial.py @@ -0,0 +1,300 @@ +"""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, +) +from s13code.core.memory.embeddings import DeterministicEmbedder + +# 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. + """ + # 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.", + "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