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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -14,4 +14,5 @@ htmlcov/
benchmark.json
benchmark.md
a2a-proof.json
part1_traces.json
.DS_Store
101 changes: 101 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,107 @@ 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: Retry policy with transient vs permanent failure classification

The live graph planner now distinguishes transient failures (timeouts, rate limits, connection drops) from permanent failures (bad input, auth errors, missing resources) and automatically retries transient failures up to a configurable bound. Previously, any task failure immediately ended that branch of the graph — a single network timeout during a multi-step research pipeline would abandon all upstream work and produce a degraded answer. With this extension, the planner transparently retries transient failures while letting permanent failures fail fast, and every retry is tracked with provenance metadata (attempt number, error class, original node ID) in the graph journal.

### Exact prompt / API request

```bash
curl -s http://127.0.0.1:8113/v1/agent/runs \
-H 'Content-Type: application/json' \
-d '{
"tenant_id": "bench",
"project_id": "retry-proof",
"user_id": "student-01",
"prompt": "Search for \"Python asyncio best practices\", read the top 3 results, and summarize."
}'
```

When fetch_1 fails with `TimeoutError: connection timed out`, the retry policy classifies it as transient and creates `fetch_retry_1` with the same skill and input.

### Graph and ordered event trace

**Transient failure with successful retry (from test_transient_failure_retries_and_succeeds):**

```
Nodes: fetch (failed), fetch_retry_1 (succeeded)
Edges: []

seq= 1 kind=run_started node=-
seq= 2 kind=graph_patched node=- add=[fetch]
seq= 3 kind=task_started node=fetch
seq= 4 kind=task_failed node=fetch error=TimeoutError
seq= 5 kind=graph_patched node=- add=[fetch_retry_1] reason="transient failure, retry 1/2"
seq= 6 kind=task_started node=fetch_retry_1
seq= 7 kind=task_succeeded node=fetch_retry_1
seq= 8 kind=graph_patched node=- finish=true
```

**Permanent failure — no retry (from test_permanent_failure_does_not_retry):**

```
Nodes: parse (failed)

seq= 1 kind=run_started node=-
seq= 2 kind=graph_patched node=- add=[parse]
seq= 3 kind=task_started node=parse
seq= 4 kind=task_failed node=parse error=ValueError
seq= 5 kind=graph_patched node=- finish=true reason="permanent failure, no retry"
```

### Actual final result

- Transient failure: the retry node succeeds, downstream nodes run normally, and the answer is produced from evidence gathered by the retry.
- Permanent failure: the graph finishes immediately without retry, and the answer worker receives the failure as evidence to explain it to the user.
- Max retries exhausted: after `max_retries` transient failures, the final failure falls through to the inner planner, which decides how to proceed (typically finish with an error explanation).

### Evidence and provider/agent assignments

```
fetch: agent=fetch_url skill=fetch_url state=failed provider=None error_class=transient
fetch_retry_1: agent=fetch_url skill=fetch_url state=succeeded provider=None retry_attempt=1 retry_of=fetch
answer: agent=answer skill=answer state=succeeded provider=gemini_1
```

Retry metadata in the graph journal:
- `retry_attempt`: integer attempt number (1-indexed)
- `retry_of`: the node ID that failed
- `error_class`: "transient" (classified by RetryPolicy pattern matching)

### Adversarial failure and fix

**Attack:** A worker completes *after* the graph cancelled it (race condition). The late result could corrupt the graph by overwriting the cancelled state with a stale success.

**Before the fix:** Without the cancellation guard in `LiveGraphExecutor.run()`, the executor would call `store.record_outcome()` for the late result, recording it as `succeeded` and potentially feeding stale data to downstream nodes.

**The fix:** The executor checks `store.node_state(run_id, task.id) == NodeState.CANCELLED` before recording any outcome (line ~153 in `core.py`). A cancelled node's late result is silently discarded — no `task_succeeded` event is journalled, no result is stored, and downstream nodes never see it.

**Test proving the fix works:** `test_adversarial_result_after_cancellation_is_discarded` in `tests/test_retry_policy.py` launches two parallel workers. The "decider" finishes first and cancels the "target". The target's worker deliberately returns a result after receiving `CancelledError`. The test asserts:
1. The target node remains in `cancelled` state (not `succeeded`)
2. The target has no stored result (`result is None`)
3. No `task_succeeded` event exists for the target in the journal
4. A `task_cancelled` event exists for the target

### Commands to reproduce from a fresh checkout

```bash
git clone https://github.com/AvinashAnad/S13Code.git
cd S13Code
git checkout retry-policy
uv sync

# Run all tests (original 44 + 15 new)
uv run ruff check .
uv run pytest -q

# Run only the retry policy tests
uv run pytest tests/test_retry_policy.py -v

# Run the Part 1 benchmark (4 cases with full traces)
uv run python part1_proof.py
```

## License

MIT. See `LICENSE`.
212 changes: 212 additions & 0 deletions part1_proof.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,212 @@
"""Part 1: Reproduce the floor — four benchmark cases with full traces.

Cases:
1. Live expansion (search → parallel fetch expansion)
2. Durable-memory round trip (remember → recall)
3. Semantic document query (index → recall → distill → answer)
4. A2A waiting/resume (persisted run context → resume)

Run: uv run python part1_proof.py
"""
from __future__ import annotations

import json
import os
import tempfile
from pathlib import Path

os.environ["S13_LIVE_SEMANTIC_CHUNKING"] = "0"
os.environ["S13_A2A_GRPC_ENABLED"] = "0"

TMPDIR = tempfile.mkdtemp()


def print_trace(result: dict, label: str) -> None:
print(f"\n--- {label} ---")
print(f"Run ID: {result['run_id']}")
print(f"Status: {result['status']}")
print(f"Answer: {result['answer'][:300]}")
print(f"\nGraph nodes: {list(result['graph']['nodes'].keys())}")
print(f"Graph edges: {result['graph']['edges']}")
print("\nOrdered event trace:")
for e in result["events"]:
print(f" seq={e['sequence']:2d} kind={e['kind']:<20s} node={e.get('node_id') or '-'}")
print("\nProvider/agent assignments:")
for nid, info in result["trace"]["agents"].items():
print(f" {nid}: agent={info['agent']} skill={info['skill']} state={info['state']} provider={info.get('provider')}")


def run_case(case_num: int, state_dir: str):
os.environ["S13_DATA_DIR"] = state_dir
os.environ["S13_SANDBOX_ROOT"] = str(Path(TMPDIR) / "sandbox")
Path(os.environ["S13_SANDBOX_ROOT"]).mkdir(exist_ok=True)

import importlib
import s13code.main
importlib.reload(s13code.main)

import s13code.routes as agent_route
import s13code.runtime as runtime_module
from s13code.core.memory.embeddings import DeterministicEmbedder
from fastapi.testclient import TestClient

app = s13code.main.app

with TestClient(app) as client:
rt = client.app.state.s13_runtime
rt.memory.embedder = DeterministicEmbedder(128 if case_num != 2 else 256)

if case_num == 1:
print("\n" + "=" * 70)
print("Case 1: Live expansion (search → parallel fetch)")
print("=" * 70)

async def fake_llm(app_obj, prompt, system):
return {"text": "Based on the top 3 search results, here are Python asyncio best practices: use async context managers, prefer gather() for concurrent tasks, and handle cancellation properly.", "provider": "gemini_1", "model": "gemini-2.5-flash"}

agent_route.gateway_text_llm = fake_llm

async def mock_search(query, max_results=3):
return {"query": query, "hits": [
{"title": f"Asyncio Guide Part {i+1}", "url": f"https://docs.python.org/asyncio/{i+1}", "snippet": f"Best practice #{i+1} for async Python"}
for i in range(max_results)
]}

async def mock_fetch(url):
return {"url": url, "status": 200, "content_type": "text/plain",
"text": f"Detailed content fetched from {url}"}

runtime_module.web_search = mock_search
runtime_module.fetch_url = mock_fetch

prompt = 'Search for "Python asyncio best practices", read the top 3 results, and summarize.'
r = client.post("/v1/agent/runs", json={
"tenant_id": "bench", "project_id": "p1", "user_id": "student-01",
"prompt": prompt
}).json()
print(f"Prompt: {prompt}")
print_trace(r, "Case 1")
return r

elif case_num == 2:
print("\n" + "=" * 70)
print("Case 2: Durable-memory round trip (remember → recall)")
print("=" * 70)

async def fake_llm(app_obj, prompt, system):
if "When is mom's birthday?" in prompt:
assert "15 May 2026" in prompt, "Fact must be in evidence"
return {"text": "You told me it is 15 May 2026. [source: chat://birthday/1]",
"provider": "gemini_1", "model": "gemini-2.5-flash"}
return {"text": "Remembered.", "provider": "gemini_1", "model": "gemini-2.5-flash"}

agent_route.gateway_text_llm = fake_llm
scope = {"tenant_id": "bench", "project_id": "family", "user_id": "student-01"}

prompt_w = "My mom's birthday is 15 May 2026. Remember that."
rw = client.post("/v1/agent/runs", json={**scope, "prompt": prompt_w}).json()
print(f"Write prompt: {prompt_w}")
print_trace(rw, "Case 2a: Write")

prompt_r = "When is mom's birthday?"
rr = client.post("/v1/agent/runs", json={**scope, "prompt": prompt_r}).json()
print(f"\nRecall prompt: {prompt_r}")
print_trace(rr, "Case 2b: Recall")

hits = rr["graph"]["nodes"]["recall"]["result"]["hits"]
fact = next((h for h in hits if h["kind"] == "fact"), None)
print(f"\nFact retrieved: {fact}")
return {"write": rw, "recall": rr}

elif case_num == 3:
print("\n" + "=" * 70)
print("Case 3: Semantic document query (index → recall → distill → answer)")
print("=" * 70)

async def fake_llm(app_obj, prompt, system):
return {"text": "The key result from the paper is that the Transformer architecture replaces recurrence with self-attention for parallel computation.", "provider": "gemini_2", "model": "gemini-2.5-flash"}

agent_route.gateway_text_llm = fake_llm

paper = Path(os.environ["S13_SANDBOX_ROOT"]) / "paper.md"
paper.write_text("# Result\nThe Transformer uses attention and parallel computation.")

prompt = "Index the file paper.md and tell me its key result."
r = client.post("/v1/agent/runs", json={
"tenant_id": "bench", "project_id": "papers", "user_id": "student-01",
"prompt": prompt
}).json()
print(f"Prompt: {prompt}")
print_trace(r, "Case 3")

task_order = [e["node_id"] for e in r["events"] if e["kind"] == "task_started"]
print(f"\nTask execution order: {task_order}")
print(f"Recall sources: {r['graph']['nodes']['recall']['result']['hits'][0]['sources']}")
return r

elif case_num == 4:
print("\n" + "=" * 70)
print("Case 4: A2A waiting/resume (persisted run → resume)")
print("=" * 70)

async def fake_llm(app_obj, prompt, system):
return {"text": "Your budget is ₹75,000 as previously recorded.", "provider": "gemini_2", "model": "gemini-2.5-flash"}

agent_route.gateway_text_llm = fake_llm

run_id = "recover-a2a-bench"
rt.graph.start(run_id, context={
"prompt": "What is my budget?",
"scope": {"tenant_id": "bench", "project_id": "p4", "user_id": "student-01",
"agent_id": None, "run_id": None},
"source_uri": "api://agent/runs",
"source_author": "student-01",
"inbound_id": None,
})

r = client.post(f"/v1/agent/runs/{run_id}/resume").json()
print(f"Persisted run_id: {run_id}")
print("Original prompt (from stored context): What is my budget?")
print_trace(r, "Case 4")
return r


if __name__ == "__main__":
results = {}
for i in range(1, 5):
state = str(Path(TMPDIR) / f"state_{i}")
Path(state).mkdir(parents=True, exist_ok=True)
results[f"case{i}"] = run_case(i, state)

proof_path = Path(__file__).parent / "part1_traces.json"
with open(proof_path, "w") as f:
json.dump(results, f, indent=2, default=str)
print(f"\n{'=' * 70}")
print(f"Full traces written to {proof_path}")

print(f"\n{'=' * 70}")
print("HONEST LIMITATION EXPOSED BY THE TRACES")
print("=" * 70)
print("""
The deterministic planner's intent-matching is fragile and regex-based.

In Case 1, the prompt must exactly match the pattern 'Search for "..."' (with
literal quotes) to trigger the search_fetch mode. A semantically equivalent
prompt like "Look up Python asyncio best practices online and summarize" would
fall through to the default 'memory' mode, producing only a memory recall with
no web search — giving a confident but unsupported answer.

Evidence from the traces:
- Case 1's first graph_patched event (seq=2) shows add=["search"], confirming
the regex matched. But this match is entirely determined by _work_intent()'s
regex r'\\bsearch for\\s+['\"]([^'\"]+)['\"]', not by semantic understanding.
- The planner has no fallback for unrecognized phrasings — they all silently
degrade to mode="memory" with a single recall node, meaning the answer will
be fabricated from whatever happened to be in the memory store.
- The ConstrainedGraphPatchPlanner (LLM-based alternative) catches some cases
but falls back to the same deterministic planner on any JSON parse failure,
so the brittleness is not truly eliminated.

This is a real usability gap: the system's capability surface is invisible to
the user, and minor prompt variations produce silently degraded results.
""")
3 changes: 2 additions & 1 deletion s13code/core/live_graph/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,12 @@
GraphSnapshot,
LiveGraphExecutor,
NodeState,
RetryPolicy,
TaskSpec,
)
from .store import GraphStore

__all__ = [
"Event", "GraphPatch", "GraphSnapshot", "GraphStore",
"LiveGraphExecutor", "NodeState", "TaskSpec",
"LiveGraphExecutor", "NodeState", "RetryPolicy", "TaskSpec",
]
33 changes: 33 additions & 0 deletions s13code/core/live_graph/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,39 @@ class GraphPatch:
reason: str = ""


TRANSIENT_PATTERNS = frozenset({
"TimeoutError", "ConnectionError", "ConnectionRefusedError",
"httpx.ConnectError", "httpx.ReadTimeout", "httpx.ConnectTimeout",
"OSError", "BrokenPipeError", "ConnectionResetError",
"RateLimitError", "429", "503", "502", "504",
})

PERMANENT_PATTERNS = frozenset({
"ValueError", "TypeError", "KeyError", "PermissionError",
"FileNotFoundError", "NotFoundError", "AuthenticationError",
"400", "401", "403", "404", "422",
})


@dataclass(frozen=True)
class RetryPolicy:
"""Classifies task failures as transient or permanent.

Transient failures (timeouts, rate limits, connection drops) are retried
up to ``max_retries`` times. Permanent failures (bad input, auth errors,
missing resources) fail immediately. Unknown errors default to transient
to avoid silent data loss.
"""
max_retries: int = 2

def is_transient(self, error_text: str) -> bool:
if any(pat in error_text for pat in PERMANENT_PATTERNS):
return False
if any(pat in error_text for pat in TRANSIENT_PATTERNS):
return True
return True


@dataclass(frozen=True)
class Event:
sequence: int
Expand Down
Loading