Files
interactive-story/backend/tools/m11_long_run.py
T
JesseMarkowitzandClaude Opus 5 f8d401029f Stop a turn locking out its own memory bank, and let the long run notice
The first M01 trial with the memory bank on was 26 turns on a GPU host. It
accepted every turn and reported "complete". It also wrote two memories and
no summary, and logged 180 `database is locked` errors, while derived status
still read `idle`.

The cause was a single uncommitted UPDATE. Retrieval bumped each used
memory's counter before the model call, and the turn commits only after the
reply has streamed. SQLite has one writer, so the turn held the write lock for
the whole reply. Every post-turn memory, summary and status write in that
window waited out the five-second timeout and failed. Recording the failure
needed a write as well, and without a rollback first it raised
PendingRollbackError. The loss therefore reached the log and never reached
the status the Insights panel reads, which F08 forbids. The draco run never
hit this because the bank was off there.

- `retrieve_memories` now only reads. `record_use` writes the counters in the
  turn's single commit, so a turn that never lands counts nothing.
- The post-turn task's outer handler rolls back before it records a failure.

The harness could not have caught any of this. It read three prompt sections
under names the builder does not use: `memories` (really `used_memories`),
`story_history` (really `history`/`recent_history`), and a `knowledge` prefix
that matched the fixed instruction section instead of the imported passages.
Memory tokens read 0 whatever the prompt held, and the in-history and
in-memories recall checks could never come out true. The labels are now
constants, pinned by a test against a prompt the real builder assembled.

The harness also stops at the first sign of failed post-turn work. It checks
/derived and new server.log lines after every turn, keeps its log position
across --resume, and waits for background work to settle before its final
checks. A run with no memories or no summaries now ends "failed", not
"complete".

Both new application tests fail on fec46f6: the lock probe sees
`database is locked`, and memory status stays `idle`. The full backend suite
passes (1392 passed, 17 skipped). A 26-turn re-run against the same host had
0 lock errors, wrote 7 memories and 2 summaries, and used them in the prompt
from turn 8.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_0136VBTMUKWYeU6G9HgbDbND
2026-09-13 20:21:21 -04:00

1158 lines
53 KiB
Python

"""M11 M01-M04: a real 100-turn campaign, against a real narrator, over HTTP.
python -m tools.m11_long_run --turns 100 --out <dir>
Run from `backend/`. Reads `AIDND_TEST_ENDPOINT`, `AIDND_TEST_MODEL` and
`AIDND_TEST_EMBED_MODEL`, all three required.
## Why this is a script that spawns servers rather than a test
M01's pass condition is not "100 requests succeeded". It is **100+ accepted turns
with no continuity, state, history, authority, lineage or recovery corruption**,
across genuine application restarts, with a fact planted at the beginning
recoverable at the end through memory rather than through the transcript.
Three of those words decide the shape of this harness:
*Accepted* — a turn counts when the application committed it, so every turn is
checked for a committed action and a state document, not for an HTTP 200.
*Restarts* — M02 says a new test client, a reconnected browser and a reopened
session do **not** count. So the storyteller runs as a real `uvicorn` process,
started with the command `DEVELOPMENT.md` documents, and is killed and restarted
at planned points. Everything that survives crosses as bytes on disk.
*Recoverable* — the planted clue has to be pushed out of the recent-history
window and then retrieved, so the run measures the window at intervals and the
recall check at the end asks the application what it would actually send.
M01's step list also asks for **"summary/memory activation"**, and both are
per-campaign switches that default to off. An earlier version of this harness
never turned them on, so a hundred turns ran with an empty memory bank and no
summaries: six of M01's seven clauses were exercised and the seventh was
reported by silence. Worse for M04, whose whole question is whether a fact
planted at turn one is still reachable at turn a hundred — with the bank off it
was reachable through narrative state alone, and the retrieval path M6 built was
never asked. `setup` now turns both on and proves it, and `measure` records how
many memories and summaries exist so that "the bank stayed empty" is a number in
the evidence rather than an absence nobody looked for.
Turning both switches on did not make them work. The first run with them on
was a 26-turn trial on a GPU host. It accepted every turn and reported
"complete" with two memories, no summary and 180 `database is locked` errors in
`server.log`. Two defects hid that result. The application lost its own
post-turn writes and then could not record the loss. The harness read three
prompt sections under names the builder does not use, so it measured 0 memory
tokens whatever was in the prompt. The harness now checks derived status and
`server.log` after every turn and stops at the first sign of failed background
work. A run with no memories or no summaries at the end is reported as
`failed`, not `complete`.
## What it records
A JSON line per turn (`timeline.jsonl`) carrying the context measurements M03
wants, the recall evidence for M04 (`recall.json`), the exported campaign
(`bundle.json`) and `summary.json`. Everything is written as it happens, so a
run that dies at turn 80 still leaves 80 turns of evidence rather than nothing.
## Resuming
A hundred turns is hours of wall clock, and the first release attempt lost one
at turn 97 to a host crash. `timeline.jsonl` survived that; the run did not,
because starting the harness again began a new campaign at turn 1.
So a run now checkpoints `resume.json` beside its evidence — after the prologue,
after every scheduled operation, and after every turn — and `--resume` picks the
campaign back up where it stopped: same adventure row, same accepted count, same
place in the beat cycle, and the scheduled operations that already fired are not
fired again. The file is written under a temporary name and renamed, because the
failure it exists to survive is the host dying mid-write.
`resume.json` is operational state, never evidence. `timeline.jsonl` remains the
append-only record and nothing here rewrites it. A finished run deletes its
`resume.json`, which makes the file's presence mean exactly one thing: there is
an unfinished run in this directory. The harness refuses to start a fresh
campaign in a directory that already holds one, because two campaigns
interleaved in one timeline are worse evidence than none.
## Timeouts, and why they are an option rather than a constant
How long a turn takes belongs to the inference host, not to the application. On
the reference host a turn cost 229-291 seconds at the recommended window, which
fits comfortably inside the 600-second model timeout this harness used to
hard-code. A slower host does not, and the consequence was not a slow run: a
turn that overran the timeout raised out of the loop and ended the run with a
traceback and no summary.
`--turn-timeout` sets what the application will wait for one narrator reply, and
the harness waits longer still, so that the application's own error arrives
inside the stream rather than being cut off at the socket. The default is
deliberately generous; measure your host before lowering it.
`--max-consecutive-failures` ends a run that has stopped producing turns, with
its evidence and a summary written and `--resume` still able to continue it,
instead of spinning against a narrator that is not answering.
"""
from __future__ import annotations
import argparse
import json
import os
import socket
import subprocess
import sys
import time
import urllib.error
import urllib.request
from datetime import datetime
from pathlib import Path
HERE = Path(__file__).resolve().parent
BACKEND = HERE.parent
ENDPOINT = os.environ.get("AIDND_TEST_ENDPOINT", "")
MODEL = os.environ.get("AIDND_TEST_MODEL", "")
EMBED_MODEL = os.environ.get("AIDND_TEST_EMBED_MODEL", "")
#: What the application will wait for one narrator reply, in seconds. The
#: settings schema bounds this at 30..3600 (`app/schemas.py`), and this default
#: sits high inside that range deliberately: an overrun turn is a lost turn, and
#: over a hundred of them the cost of waiting is far below the cost of a restart.
DEFAULT_TURN_TIMEOUT = 1800
#: How much longer the harness waits than the application does. The application
#: has to be the thing that times out, because it reports the failure inside the
#: stream the harness is reading; a harness that gave up first would record a
#: socket error and throw away what the application was about to say.
HARNESS_TIMEOUT_MARGIN = 300
#: Turns in a row not accepted before the run stops and writes what it has. One
#: rejected turn is ordinary — `failed_call` causes one on purpose. A run of
#: them means the narrator is gone, and every further attempt costs a timeout.
DEFAULT_MAX_CONSECUTIVE_FAILURES = 5
#: Operational state, not evidence. Its presence means an unfinished run.
RESUME_FILE = "resume.json"
#: The prompt sections this harness reads, spelled the way
#: `app/context/builder.py` and `app/knowledge/classes.py` spell them. A wrong
#: label raises no error. `.get` returns 0 and a membership check returns
#: False, so each misspelling becomes a check that can never pass. Three of
#: them did: `memories`, `story_history`, and a `knowledge` prefix that matched
#: the fixed instruction section instead of the retrieved passages.
#: `test_context_memory.py` pins these names against a prompt the real builder
#: assembled. They are copied rather than imported, because importing `app`
#: here would build a database engine in the harness process.
HISTORY_LABELS = ("history", "recent_history")
SUMMARY_LABEL = "story_summary"
MEMORIES_LABEL = "used_memories"
STATE_LABEL = "narrative_state"
CANON_LABEL = "campaign_canon"
IMPORTED_KNOWLEDGE_LABELS = (
"imported_canon_always", "imported_canon", "imported_reference",
"imported_inspiration",
)
#: Lines `app/derived.py` and `app/memorybank.py` log when post-turn work fails.
#: The second one means the failure could not be written to derived status at
#: all. Status alone therefore cannot prove the work is healthy: in the first
#: 26-turn GPU trial, 20 failures were logged this way and the status still read
#: `idle`.
LOG_FAILURE_MARKERS = (
"work failed for adventure",
"could not record derived-work failure",
)
#: The planted clue. Distinctive enough that its presence anywhere is
#: unambiguous, and phrased as something a story would actually establish.
CLUE = "the silver key opens the crypt beneath the Old Abbey"
CLUE_SENTINEL = "SILVER-KEY-CRYPT-OLD-ABBEY"
#: The clue as accepted state, which is the half of M04 that does not depend on
#: the narrator remembering anything.
#:
#: The content goes in `value`. `add_fact` requires `predicate` and accepts
#: `subject`, `object`, `value` and `fact_id` — and nothing else
#: (`app/narrative/events.py` SPECS). An earlier version of this harness put the
#: clue in a `detail` key, which that event does not define: the correction was
#: accepted, the fact was created, and the clue text went nowhere. The stored
#: fact said only that Aldric knows of the abbey, so `_recall`'s
#: `fact_still_in_state` could not answer True however well the application
#: behaved. A check that can only fail is worse than no check, and this is the
#: second harness defect of that shape M11 has found.
CLUE_FACT = {
"type": "add_fact",
"subject": "aldric",
"predicate": "knows",
"object": "abbey",
"value": f"{CLUE} ({CLUE_SENTINEL})",
"fact_id": "silver-key-opens-crypt",
}
CANON = [
"The dead do not return. No rite, relic or bargain has ever returned anyone.",
"The abbey crypt has been sealed since the founding.",
"Aldric is the protagonist and the one the reader plays.",
]
CANON_MD = """# Westhaven
## The Old Abbey
The abbey above Westhaven has stood since the founding. Its crypt is sealed.
## What cannot happen here
The dead do not return. No rite, relic or bargain in Westhaven has ever
returned anyone from death, and none ever will.
"""
REFERENCE_MD = """# Roads and weather of the Fen
The fen road floods between the autumn rains and the first hard frost. Traders
take the ridge track instead, which adds a day.
"""
INSPIRATION_MD = """# Tone notes
Rain on slate. Lamplight through smoke. People who say less than they mean.
"""
#: The beats the campaign plays through, cycled. Written so the story keeps
#: moving and keeps giving the state extractor something to do, rather than a
#: hundred repetitions of one sentence.
BEATS = [
"I ask Mara what she has heard about the abbey.",
"I walk down to the waterfront and watch the boats.",
"I ask the ferryman about the fen road.",
"I look through my pack for anything useful.",
"I go back to the tavern and sit by the fire.",
"I ask Mara whether Edrin has been seen.",
"I take the ridge track north out of town.",
"I stop at the shrine on the ridge and look back at Westhaven.",
"I talk to the trader waiting out the rain.",
"I check the sky and decide whether to press on.",
]
def free_port() -> int:
with socket.socket() as s:
s.bind(("127.0.0.1", 0))
return s.getsockname()[1]
class Storyteller:
"""The real application, started the way `DEVELOPMENT.md` says to start it."""
def __init__(self, db_path: Path, log: Path, *,
turn_timeout: int = DEFAULT_TURN_TIMEOUT):
self.db_path = db_path
self.port = free_port()
self.log_path = log
self.proc = None
self.starts = 0
self.turn_timeout = turn_timeout
self.stream_timeout = turn_timeout + HARNESS_TIMEOUT_MARGIN
def start(self) -> None:
self.starts += 1
handle = open(self.log_path, "ab")
self.proc = subprocess.Popen(
[str(BACKEND / ".venv/bin/uvicorn"), "app.main:app",
"--host", "127.0.0.1", "--port", str(self.port)],
cwd=str(BACKEND), stdout=handle, stderr=subprocess.STDOUT,
env={**os.environ, "AIDND_DB_PATH": str(self.db_path),
"AIDND_DATABASE_URL": "", "DATABASE_URL": ""},
)
deadline = time.monotonic() + 90
while time.monotonic() < deadline:
if self.proc.poll() is not None:
raise SystemExit(f"server exited early; see {self.log_path}")
try:
self.call("GET", "/settings")
return
except (urllib.error.URLError, ConnectionError, OSError):
time.sleep(0.1)
raise SystemExit(f"server never became ready; see {self.log_path}")
def stop(self) -> None:
if self.proc and self.proc.poll() is None:
self.proc.terminate()
try:
self.proc.wait(timeout=20)
except subprocess.TimeoutExpired:
self.proc.kill()
self.proc.wait(timeout=20)
def listening(self) -> bool:
try:
self.call("GET", "/settings")
return True
except Exception:
return False
def restart(self) -> None:
"""A genuine OS process boundary, proved gone before it is replaced."""
self.stop()
assert not self.listening(), "the old process is still answering"
self.port = free_port()
self.start()
def call(self, method: str, path: str, payload=None, timeout=600):
data = json.dumps(payload).encode() if payload is not None else None
request = urllib.request.Request(
f"http://127.0.0.1:{self.port}/api{path}", data=data, method=method,
headers={"Content-Type": "application/json"} if data else {},
)
with urllib.request.urlopen(request, timeout=timeout) as response:
body = response.read().decode()
return json.loads(body) if body else None
def stream(self, path: str, payload, timeout=None) -> list[dict]:
"""A turn. The reply is SSE, and a failed turn is an event, not a status.
`app/sse.py`: "A failed turn is still an HTTP 200 response, because the
error is reported inside the stream the client is already reading." A
harness that read the status code would call every failure a success —
which is precisely the class of harness defect M8's review warned about.
"""
if timeout is None:
timeout = self.stream_timeout
request = urllib.request.Request(
f"http://127.0.0.1:{self.port}/api{path}",
data=json.dumps(payload).encode(), method="POST",
headers={"Content-Type": "application/json"},
)
events: list[dict] = []
with urllib.request.urlopen(request, timeout=timeout) as response:
for raw in response:
line = raw.decode(errors="replace").strip()
if line.startswith("data:"):
try:
events.append(json.loads(line[5:].strip()))
except json.JSONDecodeError:
pass
return events
class Run:
"""One long campaign, and everything measured about it."""
def __init__(self, server: Storyteller, out: Path, *, turns_target: int,
turn_timeout: int = DEFAULT_TURN_TIMEOUT):
self.server = server
self.out = out
self.timeline = (out / "timeline.jsonl").open("a")
self.adv = 0
self.accepted = 0
self.events: list[dict] = []
self.turns_target = turns_target
self.turn_timeout = turn_timeout
#: Where the beat cycle stands. Carried across a resume, so a continued
#: campaign keeps moving rather than replaying its first ten beats.
self.beat = 0
#: Plan keys whose scheduled operation has already fired, keyed by turn
#: number rather than by step name: two of the steps are called `retry`
#: and two `restart`, so a name does not identify one.
self.completed_steps: set[int] = set()
self.resumed = False
self.elapsed_before = 0.0
self.session_started = time.monotonic()
#: How far `background_failures` has read `server.log`. Carried across a
#: resume, so a failure is reported once, not again every session.
self.log_offset = 0
# ------------------------------------------------------------ recording
def note(self, kind: str, **fields) -> None:
entry = {"at": datetime.now().isoformat(timespec="seconds"),
"kind": kind, "accepted_turns": self.accepted, **fields}
self.events.append(entry)
self.timeline.write(json.dumps(entry, default=str) + "\n")
self.timeline.flush()
# ------------------------------------------------------------- resuming
def elapsed(self) -> int:
"""Seconds of run time, across every session this campaign has had."""
return round(self.elapsed_before + (time.monotonic() - self.session_started))
def save_resume(self) -> None:
"""Checkpoint enough to pick this campaign up again, atomically."""
payload = {
"adventure": self.adv,
"accepted": self.accepted,
"beat": self.beat,
"completed_steps": sorted(self.completed_steps),
"server_starts": self.server.starts,
"elapsed_seconds": self.elapsed(),
"turns_target": self.turns_target,
"log_offset": self.log_offset,
"written": datetime.now().isoformat(timespec="seconds"),
}
tmp = self.out / (RESUME_FILE + ".tmp")
tmp.write_text(json.dumps(payload, indent=2))
tmp.replace(self.out / RESUME_FILE)
def adopt(self, prior: dict) -> None:
"""Take on the state a previous session checkpointed.
A checkpoint can be at most one turn behind the database, because a turn
is committed by the application before this file is written. Erring that
way costs one extra turn on a campaign that wants *at least* a hundred,
which is the harmless direction.
"""
self.adv = prior["adventure"]
self.accepted = prior["accepted"]
self.beat = prior.get("beat", 0)
self.completed_steps = set(prior.get("completed_steps") or [])
self.elapsed_before = prior.get("elapsed_seconds", 0)
self.log_offset = prior.get("log_offset", 0)
self.resumed = True
def reattach(self) -> None:
"""Prove the campaign is still there, and restate what a run needs.
Settings are re-applied rather than trusted. They live in the database,
and `failed_call` deliberately points the model at a name the server
does not serve before putting it back; a host that died inside that
window left the campaign configured to fail every turn it is given.
"""
page = self.server.call("GET", f"/adventures/{self.adv}/actions?limit=1")
self.apply_settings()
self.note("resumed", adventure=self.adv, accepted=self.accepted,
beat=self.beat, actions_in_db=page["total"],
completed_steps=sorted(self.completed_steps),
process_starts=self.server.starts,
elapsed_before_seconds=round(self.elapsed_before))
# ------------------------------------------------------------- campaign
def apply_settings(self) -> dict:
settings = self.server.call("PUT", "/settings", {
"endpoint_url": ENDPOINT, "model": MODEL,
"embedding_model": EMBED_MODEL, "context_token_budget": 16384,
"max_output_tokens": 500,
"model_timeout_seconds": self.turn_timeout,
"memory_top_k": 4,
})
self.note("settings", model=settings["model"],
budget=settings["context_token_budget"],
model_timeout_seconds=settings["model_timeout_seconds"])
return settings
def setup(self) -> None:
self.apply_settings()
created = self.server.call("POST", "/adventures", {
"title": "Continuity Test (M11 long run)",
"opening": (
"Rain over Westhaven. Aldric sits in the Crooked Lantern with a "
"silver key in his pocket and no-one to give it to."
),
"canon_rules": CANON,
"persona_name": "Aldric",
"narration_length": "brief",
})
self.adv = created["id"]
self.note("campaign", id=self.adv)
# M01 asks for "summary/memory activation". Both are per-campaign and
# both default to False (`models.Adventure`), so a campaign created and
# played without touching them never writes a memory or a summary.
# Enabled here, and read back rather than assumed: a PATCH that silently
# did nothing would leave the same hole this closes.
self.server.call("PATCH", f"/adventures/{self.adv}",
{"memory_bank_enabled": True, "auto_summarize": True})
back = self.server.call("GET", f"/adventures/{self.adv}")
self.note("memory_activated",
memory_bank_enabled=back.get("memory_bank_enabled"),
auto_summarize=back.get("auto_summarize"))
if not (back.get("memory_bank_enabled") and back.get("auto_summarize")):
raise SystemExit(
"the campaign did not accept memory-bank and auto-summarize; "
"M01's summary/memory clause cannot be measured from this run.")
for name, body, kind in (("canon.md", CANON_MD, "canon"),
("reference.md", REFERENCE_MD, "reference"),
("inspiration.md", INSPIRATION_MD, "inspiration")):
self.upload(name, body, kind)
# The cast and the opening scene, as accepted state rather than prose.
self.correct([
{"type": "create_entity", "entity": "aldric",
"entity_type": "character", "name": "Aldric"},
{"type": "create_entity", "entity": "mara",
"entity_type": "character", "name": "Mara"},
{"type": "create_entity", "entity": "edrin",
"entity_type": "character", "name": "Edrin"},
{"type": "create_entity", "entity": "tavern",
"entity_type": "location", "name": "The Crooked Lantern"},
{"type": "create_entity", "entity": "abbey",
"entity_type": "location", "name": "The Old Abbey"},
{"type": "create_entity", "entity": "silver_key",
"entity_type": "item", "name": "the silver key"},
{"type": "set_possession", "item": "silver_key", "owner": "aldric"},
{"type": "set_scene",
"summary": "Aldric and Mara in the Crooked Lantern, rain outside.",
"location": "tavern", "present": ["aldric", "mara"]},
], note="the opening cast")
def upload(self, name: str, body: str, classification: str) -> None:
"""Multipart by hand: the harness speaks HTTP, not the test client."""
boundary = "----m11longrun"
parts = (
f"--{boundary}\r\nContent-Disposition: form-data; name=\"classification\""
f"\r\n\r\n{classification}\r\n"
f"--{boundary}\r\nContent-Disposition: form-data; name=\"file\"; "
f"filename=\"{name}\"\r\nContent-Type: text/markdown\r\n\r\n{body}\r\n"
f"--{boundary}--\r\n"
).encode()
request = urllib.request.Request(
f"http://127.0.0.1:{self.server.port}/api/adventures/{self.adv}/knowledge",
data=parts, method="POST",
headers={"Content-Type": f"multipart/form-data; boundary={boundary}"},
)
with urllib.request.urlopen(request, timeout=120) as response:
body_out = json.loads(response.read().decode())
self.note("knowledge", file=name, classification=classification,
id=body_out["id"])
def correct(self, events, note="") -> None:
self.server.call("POST", f"/adventures/{self.adv}/state/corrections",
{"events": events, "note": note})
self.note("state_correction", events=len(events), note=note)
# ----------------------------------------------------------------- play
def turn(self, text: str, *, kind="do") -> dict:
started = time.monotonic()
before = self.count_actions()
try:
events = self.server.stream(f"/adventures/{self.adv}/actions",
{"type": kind, "text": text})
except urllib.error.HTTPError as exc:
self.note("turn_failed", text=text, status=exc.code,
detail=exc.read().decode()[:300])
return {"accepted": False}
except Exception as exc: # noqa: BLE001
# A socket timeout, or a connection dropped mid-stream. This used to
# propagate out of the loop and end the run with a traceback and no
# summary, which on a slow host is the likeliest way to lose one.
self.note("turn_failed", text=text,
detail=f"{type(exc).__name__}: {exc}"[:300])
return {"accepted": False}
errors = [e for e in events if e.get("type") == "error"]
if errors:
self.note("turn_error", text=text, detail=errors[0].get("detail", "")[:300])
return {"accepted": False, "error": errors[0].get("detail", "")}
after = self.count_actions()
if after <= before:
self.note("turn_not_accepted", text=text)
return {"accepted": False}
self.accepted += 1
seconds = time.monotonic() - started
sample = self.measure()
self.note("turn", text=text, seconds=round(seconds, 1), **sample)
return {"accepted": True, "seconds": seconds, **sample}
def count_actions(self) -> int:
return self.server.call("GET", f"/adventures/{self.adv}/actions?limit=1")["total"]
def measure(self) -> dict:
"""M03's numbers, read from the prompt the app would send right now."""
report = self.server.call("GET", f"/adventures/{self.adv}/context")
tokens = report["tokens"]
sections = {s["label"]: s["tokens"] for s in report["sections"]}
window = report.get("window") or {}
return {
"total_actions": report["history"]["total"],
"history_included": report["history"]["included"],
# Where the history window was cut, and the step it takes when it
# moves. A `floor_depth` that is the same on two consecutive turns
# is the prompt's prefix having been preserved, which is the whole
# of what trimming the window in blocks buys; a run that recorded
# neither could not say whether it engaged, held or stepped, and a
# continuity finding could not be attributed. `.get` because a
# campaign resumed against an older build has neither.
"history_floor_depth": report["history"].get("floor_depth"),
"history_trim_block": report["history"].get("trim_block"),
# M01's summary/memory clause, as a count on every turn. A bank that
# stays at zero is then visible in the evidence while the run is
# still going, instead of being discovered afterwards.
"memories_in_bank": self.bank_size(),
"prompt_tokens": tokens["total"],
"budget": tokens["budget"],
"configured_budget": tokens.get("configured_budget"),
"output_reserve": tokens["output_reserve"],
"protected": tokens["protected"],
"available_for_history": tokens["available_for_history"],
"summary_tokens": sections.get(SUMMARY_LABEL, 0),
"memory_tokens": sections.get(MEMORIES_LABEL, 0),
"knowledge_tokens": sum(
sections.get(label, 0) for label in IMPORTED_KNOWLEDGE_LABELS),
"state_tokens": sections.get(STATE_LABEL, 0),
"canon_tokens": sections.get(CANON_LABEL, 0),
"window_verified": window.get("verified"),
"window_tokens": window.get("tokens"),
# Whether the *planted clue* is still visible anywhere in the
# assembled prompt. Named for what it measures: an earlier version
# called this `canon_present`, which it never was — the campaign
# canon's presence is `canon_tokens`, which is non-zero on every
# turn. This one going to zero is M04's precondition: the clue has
# left the recent-history window and can only come back through
# memory, summary or state.
"clue_in_prompt": CLUE_SENTINEL in json.dumps(report["sections"]),
}
def bank_size(self) -> int:
"""How many memories the bank holds. Never raises: this is measurement,
and a turn is not worth failing over a count."""
try:
return len(self.server.call(
"GET", f"/adventures/{self.adv}/memories") or [])
except Exception: # noqa: BLE001
return -1
def summary_count(self) -> int:
"""How many summaries the campaign has written, on any branch. Never
raises, and -1 means unreadable, as for `bank_size`."""
try:
derived = self.server.call("GET", f"/adventures/{self.adv}/derived") or {}
return len(derived.get("summaries") or [])
except Exception: # noqa: BLE001
return -1
def background_failures(self) -> list[str]:
"""Every sign since the last call that post-turn work failed.
Two sources, because neither is enough alone. Derived status is what the
application says. The server log also catches a failure the application
could not record, which is the failure that made the first GPU trial
report `idle` over 180 `database is locked` errors. Log lines are read
from where the last call stopped, so each failure is reported once.
"""
found: list[str] = []
try:
derived = self.server.call("GET", f"/adventures/{self.adv}/derived") or {}
for row in derived.get("status") or []:
if row.get("status") == "failed":
found.append(f"{row.get('kind')}: {(row.get('detail') or '')[:200]}")
except Exception as exc: # noqa: BLE001 - an unreadable status is noted, not a failure
self.note("derived_unreadable", error=f"{type(exc).__name__}: {exc}"[:200])
log = getattr(self.server, "log_path", None)
if log is not None and Path(log).exists():
with open(log, "rb") as handle:
handle.seek(self.log_offset)
fresh = handle.read()
self.log_offset = handle.tell()
for line in fresh.decode(errors="replace").splitlines():
if any(marker in line for marker in LOG_FAILURE_MARKERS):
found.append(f"server.log: {line.strip()[:200]}")
return found
def settle(self, *, quiet_seconds: int = 30, limit_seconds: int = 900) -> None:
"""Waits for post-turn work to stop changing derived status.
Used before the final checks. The last turn's memory and summary passes
run after that turn returns, and a check taken before they finish can
miss their failure or their output. No endpoint reports running work, so
"settled" means derived status unchanged for `quiet_seconds`."""
deadline = time.monotonic() + limit_seconds
last, since = None, time.monotonic()
while time.monotonic() < deadline:
try:
now = json.dumps(self.server.call(
"GET", f"/adventures/{self.adv}/derived"), sort_keys=True)
except Exception: # noqa: BLE001
now = None
if now is not None and now == last:
if time.monotonic() - since >= quiet_seconds:
return
else:
last, since = now, time.monotonic()
time.sleep(2)
self.note("settle_timeout", limit_seconds=limit_seconds)
def state(self) -> dict:
return self.server.call("GET", f"/adventures/{self.adv}/state")
def head(self) -> tuple:
page = self.server.call("GET", f"/adventures/{self.adv}/actions?limit=1")
return page["total"], page["can_undo"], page["can_redo"]
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--turns", type=int, default=100)
parser.add_argument("--out", required=True)
parser.add_argument(
"--resume", action="store_true",
help="continue the unfinished run in --out rather than starting one")
parser.add_argument(
"--turn-timeout", type=int, default=DEFAULT_TURN_TIMEOUT,
help=("seconds the application waits for one narrator reply "
f"(30-3600, default {DEFAULT_TURN_TIMEOUT})"))
parser.add_argument(
"--max-consecutive-failures", type=int,
default=DEFAULT_MAX_CONSECUTIVE_FAILURES,
help="stop and write the evidence after this many unaccepted turns")
args = parser.parse_args()
if not (ENDPOINT and MODEL and EMBED_MODEL):
# The embedding model is not optional. Without one the summary pass
# still writes memories, but `memorybank.retrieve` answers "No embedding
# model configured" and returns none — so the bank fills and M04's
# memory path is never asked, which is the hole `setup` exists to close.
print("set AIDND_TEST_ENDPOINT, AIDND_TEST_MODEL and AIDND_TEST_EMBED_MODEL")
return 2
if not 30 <= args.turn_timeout <= 3600:
# The settings schema's own bound, checked here so a mistyped timeout
# fails in the first second rather than on the first PUT.
print("--turn-timeout must be between 30 and 3600 seconds")
return 2
out = Path(args.out)
out.mkdir(parents=True, exist_ok=True)
db_path = out / "campaign.db"
resume_path = out / RESUME_FILE
prior = _resume_state(resume_path, out / "timeline.jsonl", args.resume)
if isinstance(prior, str):
print(prior)
return 2
server = Storyteller(db_path, out / "server.log",
turn_timeout=args.turn_timeout)
if prior:
server.starts = prior.get("server_starts", 1)
server.start()
run = Run(server, out, turns_target=args.turns,
turn_timeout=args.turn_timeout)
if prior:
run.adopt(prior)
try:
if run.resumed:
run.reattach()
else:
run.setup()
# ---- The planted clue, at the very beginning. ----
run.turn(f"I tell Mara quietly that {CLUE} — {CLUE_SENTINEL}.")
run.correct([CLUE_FACT], note="the planted clue, as accepted state")
# Proved here, at turn one, where it costs a single request. M04
# asks whether the clue is still recoverable a hundred turns later,
# and that question is meaningless if it was never stored — so a
# run whose clue did not land is stopped rather than spending hours
# measuring nothing. See CLUE_FACT for how this went wrong before.
planted = any(CLUE_SENTINEL in json.dumps(fact) for fact
in run.state()["document"].get("facts") or [])
run.note("clue_planted", sentinel=CLUE_SENTINEL,
verified_in_state=planted)
if not planted:
raise SystemExit(
"the planted clue is not in accepted state, so M04 cannot "
"be measured from this run. Stopping before the campaign "
"starts rather than reporting a recall failure later.")
# The first checkpoint, and the point from which --resume works: the
# campaign exists and its clue is planted.
run.save_resume()
# ---- The long middle. ----
plan = _schedule(args.turns)
consecutive_failures = 0
aborted = None
while run.accepted < args.turns:
at = run.accepted + 1
step = plan.get(at)
if step is not None and at not in run.completed_steps:
try:
_do_step(run, server, step)
except Exception as exc: # noqa: BLE001
# A step that fails is a finding, not a reason to lose the
# other ninety turns. It is recorded loudly and the campaign
# goes on, because an abandoned run proves nothing at all.
run.note("step_failed", step=step, at=at,
error=f"{type(exc).__name__}: {exc}"[:300])
# Fired, however it went. Before this the step was keyed only by
# the accepted count, so a refused turn afterwards ran the whole
# operation again — and a restart or a retry performed twice is
# not the test the schedule describes.
run.completed_steps.add(at)
run.save_resume()
result = run.turn(BEATS[run.beat % len(BEATS)])
run.beat += 1
if result.get("accepted"):
consecutive_failures = 0
else:
consecutive_failures += 1
if consecutive_failures >= args.max_consecutive_failures:
aborted = (
f"{consecutive_failures} turns in a row were not accepted; "
"the narrator is not answering. Stopping with the evidence "
"written, so --resume can carry this campaign on."
)
run.note("run_aborted", reason=aborted,
accepted_turns=run.accepted)
run.save_resume()
break
# Checked on every turn, not only at the end. A memory or summary
# pass that fails is lost M01 evidence from that turn onward. A run
# that carries on would report "complete" over a bank that stopped
# filling, which is what the first GPU trial did.
failures = run.background_failures()
if failures:
aborted = (
f"post-turn memory/summary work failed ({len(failures)} "
"signs, first: " + failures[0] + "). Stopping with the "
"evidence written; see server.log."
)
run.note("run_aborted", reason=aborted, failures=failures[:20],
accepted_turns=run.accepted)
run.save_resume()
break
run.save_resume()
# ---- M04: the recall check, with controls. ----
# Skipped on an aborted run: it asks the narrator a question, and the
# reason the run stopped is that the narrator does not answer.
recall = None
if aborted is None:
run.note("recall_begin")
recall = _recall(run)
(out / "recall.json").write_text(json.dumps(recall, indent=2))
# ---- Export whatever exists, for the recovery evidence. ----
# Attempted even for an aborted run: the recovery check and the storage
# numbers are worth having at whatever turn count was reached.
bundle_bytes = None
try:
bundle = server.call("GET", f"/adventures/{run.adv}/export")
(out / "bundle.json").write_text(json.dumps(bundle))
bundle_bytes = len((out / "bundle.json").read_bytes())
run.note("exported", bytes=bundle_bytes)
except Exception as exc: # noqa: BLE001
run.note("export_failed", error=f"{type(exc).__name__}: {exc}"[:300])
# Every turn was accepted, but that does not complete M01. Its
# summary/memory clause needs a bank that filled and a summary that was
# written, and the recall turn's own post-turn work has to have
# finished without failing. Otherwise the run is "failed", not
# "complete".
failed_reason = None
if aborted is None:
run.settle()
late = run.background_failures()
if late:
failed_reason = (f"post-turn work failed after the last turn "
f"({len(late)} signs, first: {late[0]})")
else:
failed_reason = _activation_shortfall(run)
if failed_reason:
run.note("run_failed", reason=failed_reason)
summary = {
"status": ("aborted" if aborted
else "failed" if failed_reason else "complete"),
"aborted_reason": aborted,
"failed_reason": failed_reason,
"summaries": _or_none(run.summary_count),
"accepted_turns": run.accepted,
"turns_requested": args.turns,
"restarts": server.starts - 1,
"process_starts": server.starts,
"resumed": run.resumed,
"elapsed_seconds": run.elapsed(),
"turn_timeout_seconds": args.turn_timeout,
"recall": recall,
"final_state": _or_none(lambda: run.state()["document"]),
"final_measurement": _or_none(run.measure),
"db_bytes": db_path.stat().st_size,
"bundle_bytes": bundle_bytes,
# Whether the seventh clause of M01's step list actually happened.
"memories_in_bank": _or_none(run.bank_size),
}
(out / "summary.json").write_text(json.dumps(summary, indent=2, default=str))
print(json.dumps({k: v for k, v in summary.items()
if k not in ("final_state",)}, indent=2, default=str)[:2000])
if aborted is None:
# A finished run has nothing to resume, and the file's absence is
# what lets a later run use this directory.
resume_path.unlink(missing_ok=True)
return 1 if failed_reason else 0
return 1
finally:
server.stop()
run.timeline.close()
def _activation_shortfall(run) -> str | None:
"""Why M01's summary/memory clause was not exercised, or None if it was.
A count of -1 means the count could not be read. It counts as a shortfall,
because an unreadable bank does not show that the bank filled."""
memories, summaries = run.bank_size(), run.summary_count()
missing = []
if memories <= 0:
missing.append(f"memories_in_bank={memories}")
if summaries <= 0:
missing.append(f"summaries={summaries}")
if not missing:
return None
return ("M01's summary/memory clause was not exercised: "
+ ", ".join(missing))
def _or_none(read):
"""A summary field worth having when it can be read, and worth skipping when
it cannot. An aborted run still reports the fields that do answer."""
try:
return read()
except Exception as exc: # noqa: BLE001
return {"unavailable": f"{type(exc).__name__}: {exc}"[:200]}
def _resume_state(resume_path: Path, timeline_path: Path, resuming: bool):
"""The previous session's checkpoint, or a message saying why there is none.
Returns the parsed checkpoint, `None` for a clean start, or a string to
print before exiting. The refusals matter as much as the resume: two
campaigns interleaved in one `timeline.jsonl` and one `campaign.db` are
worse evidence than none, and that is what starting fresh on top of an
existing run produces.
A non-empty `timeline.jsonl` is what says a previous attempt got far enough
to record something, and it is the harness's own artifact rather than the
application's. Testing it rather than `campaign.db` means a run that died
before it recorded anything — a server that never came up, a missing
endpoint — leaves the directory usable, because nothing was written into it
to collide with.
"""
if resuming:
if not resume_path.exists():
return (f"--resume: there is no {resume_path.name} in "
f"{resume_path.parent}. Either the run never reached its "
"first checkpoint, or this is the wrong directory.")
try:
prior = json.loads(resume_path.read_text())
except (OSError, json.JSONDecodeError) as exc:
return f"--resume: cannot read {resume_path}: {exc}"
if not prior.get("adventure"):
return f"--resume: {resume_path} names no campaign."
return prior
if resume_path.exists():
return (f"{resume_path} exists, so this directory holds an unfinished "
"run. Pass --resume to carry it on, or choose a new --out.")
if timeline_path.exists() and timeline_path.stat().st_size > 0:
return (f"{timeline_path} already has entries but there is no "
f"{resume_path.name}, so a previous run recorded something here "
"and then died before its first checkpoint. Choose a new --out: "
"starting here would put a second campaign in the same timeline "
"and the same database.")
return None
def _schedule(turns: int) -> dict:
"""Where each required history operation happens. Spread, not clustered."""
unit = max(1, turns // 13)
return {
unit * 1: "save_point_1",
unit * 2: "restart",
unit * 3: "undo_redo",
unit * 4: "retry",
unit * 5: "save_point_2",
unit * 6: "restart_with_retained_history",
unit * 7: "retry",
unit * 8: "undo_diverge",
unit * 9: "take_selection",
unit * 10: "failed_call",
unit * 11: "restore_save_point",
unit * 12: "restart",
}
def _do_step(run: Run, server: Storyteller, step: str) -> None:
adv = run.adv
if step == "save_point_1":
point = server.call("POST", f"/adventures/{adv}/checkpoints",
{"name": "Before the ridge", "note": "planted clue is behind us"})
run.note("save_point", id=point["id"], name=point["name"])
elif step == "save_point_2":
point = server.call("POST", f"/adventures/{adv}/checkpoints",
{"name": "On the ridge", "note": ""})
run.note("save_point", id=point["id"], name=point["name"])
elif step in ("restart", "restart_with_retained_history"):
if step == "restart_with_retained_history":
server.call("POST", f"/adventures/{adv}/undo")
run.note("undo", why="leave retained history across the restart")
before = _snapshot(run)
server.restart()
after = _snapshot(run)
run.note("restart", number=server.starts - 1,
identical=before == after,
before=before, after=after)
elif step == "undo_redo":
total_before, _, _ = run.head()
server.call("POST", f"/adventures/{adv}/undo")
after_undo = run.head()
server.call("POST", f"/adventures/{adv}/redo")
after_redo = run.head()
run.note("undo_redo", before=total_before, after_undo=after_undo[0],
after_redo=after_redo[0], restored=after_redo[0] == total_before)
elif step == "retry":
# Retry regenerates a turn, so it streams like one. The first version of
# this harness called it as JSON and died on the SSE body — found by the
# shakeout run rather than fifty turns into the release campaign, which
# is what the shakeout was for.
events = server.stream(f"/adventures/{adv}/retry", {})
errors = [e for e in events if e.get("type") == "error"]
page = server.call("GET", f"/adventures/{adv}/actions?limit=3")
takes = max((a.get("take_count") or 1) for a in page["actions"])
run.note("retry", ok=not errors, takes_on_newest_turn=takes,
detail=(errors[0].get("detail", "")[:120] if errors else ""))
elif step == "take_selection":
# Retry first so there is more than one take to choose between, then
# step back to the earlier one — D07's "select prior retry take".
server.stream(f"/adventures/{adv}/retry", {})
page = server.call("GET", f"/adventures/{adv}/actions?limit=5")
multi = [a for a in page["actions"] if (a.get("take_count") or 1) > 1]
if multi:
target = multi[-1]
takes = server.call(
"GET", f"/adventures/{adv}/actions/{target['id']}/variants")
chosen = server.call(
"POST", f"/adventures/{adv}/actions/{target['id']}/variant",
{"index": 0})
run.note("take_selected", action=target["id"], of=len(takes),
chose_index=0, now_live=chosen["id"])
else:
run.note("take_selection_skipped", reason="no multi-take turn found")
elif step == "undo_diverge":
server.call("POST", f"/adventures/{adv}/undo")
server.call("POST", f"/adventures/{adv}/undo")
run.turn("I turn back towards the town instead.")
_, _, can_redo = run.head()
run.note("diverged", redo_available_after_new_writing=can_redo)
elif step == "restore_save_point":
points = server.call("GET", f"/adventures/{adv}/checkpoints")
if points:
target = points[0]
before = run.head()
page = server.call(
"POST", f"/adventures/{adv}/checkpoints/{target['id']}/restore")
run.note("save_point_restored", id=target["id"], name=target["name"],
total_before=before[0], total_after=page["total"],
can_redo=page["can_redo"])
else:
run.note("restore_skipped", reason="no save point exists yet")
elif step == "failed_call":
# A real failure: point the model at a name the server does not serve,
# take a turn, and put it back. Nothing is mocked.
settings = server.call("GET", "/settings")
before_total, _, _ = run.head()
before_state = run.state()["document"]
server.call("PUT", "/settings", {"model": "no-such-model-m11"})
try:
events = server.stream(
f"/adventures/{adv}/actions",
{"type": "do", "text": "I look for a way across the water."},
timeout=180)
errors = [e for e in events if e.get("type") == "error"]
outcome = (f"reported: {errors[0].get('detail', '')[:120]}"
if errors else "NO ERROR REPORTED")
except urllib.error.HTTPError as exc:
outcome = f"HTTP {exc.code}"
except Exception as exc: # noqa: BLE001 - recorded, not swallowed
outcome = type(exc).__name__
server.call("PUT", "/settings", {"model": settings["model"]})
after_total, _, _ = run.head()
run.note("failed_call", outcome=outcome,
actions_before=before_total, actions_after=after_total,
state_unchanged=before_state == run.state()["document"])
# And prove play resumes.
run.turn("I ask the ferryman again, more politely.")
def _snapshot(run: Run) -> dict:
"""What must be identical across a restart (M02's list)."""
adv = run.adv
page = run.server.call("GET", f"/adventures/{adv}/actions?limit=3")
state = run.state()["document"]
points = run.server.call("GET", f"/adventures/{adv}/checkpoints")
knowledge = run.server.call("GET", f"/adventures/{adv}/knowledge")
settings = run.server.call("GET", "/settings")
return {
"total": page["total"],
"can_undo": page["can_undo"],
"can_redo": page["can_redo"],
"newest": [a["text"][:60] for a in page["actions"]],
"scene": (state.get("scene") or {}).get("summary"),
"entities": sorted(state.get("entities") or {}),
"facts": sorted(f.get("id") for f in state.get("facts") or []),
"save_points": sorted(p["name"] for p in points),
"knowledge": sorted(k["original_filename"] for k in knowledge),
"model": settings["model"],
"budget": settings["context_token_budget"],
}
def _recall(run: Run) -> dict:
"""M04: can the planted clue still be found, and not from the transcript?"""
adv, server = run.adv, run.server
# 1. Is the clue outside the recent-history window? (Precondition, not result.)
report = server.call("GET", f"/adventures/{adv}/context")
history_text = " ".join(
s["text"] for s in report["sections"] if s["label"] in HISTORY_LABELS)
in_history = CLUE_SENTINEL in history_text
# 2. Ask about the subject, and see what the application assembles.
run.turn("I think back to what I told Mara about the key, that first night.")
after = server.call("GET", f"/adventures/{adv}/context")
sections = {s["label"]: s["text"] for s in after["sections"]}
whole_prompt = "\n".join(sections.values())
# 3. Where did it come from? State, summary, memory, retrieval — or nowhere.
document = run.state()["document"]
fact_present = any(
CLUE_SENTINEL in json.dumps(f) for f in document.get("facts") or [])
return {
"clue_in_recent_history_window": in_history,
"clue_in_prompt": CLUE_SENTINEL in whole_prompt,
"in_state_section": CLUE_SENTINEL in sections.get(STATE_LABEL, ""),
"in_summary_section": CLUE_SENTINEL in sections.get(SUMMARY_LABEL, ""),
"in_memories_section": CLUE_SENTINEL in sections.get(MEMORIES_LABEL, ""),
"in_knowledge_sections": any(
CLUE_SENTINEL in sections.get(label, "")
for label in IMPORTED_KNOWLEDGE_LABELS),
"fact_still_in_state": fact_present,
"history_included": after["history"]["included"],
"history_total": after["history"]["total"],
"memories_used": [
m.get("text", "")[:120] for m in (after.get("memories") or {}).get("used", [])
],
# M01's summary/memory clause, stated rather than implied. A recall that
# succeeds only through narrative state, with an empty bank, has proved
# one of the two paths the design has — and the reader of this report
# should be able to see which.
"memories_in_bank": run.bank_size(),
"summary_present": bool(
(server.call("GET", f"/adventures/{adv}") or {}).get("story_summary")),
"prompt_tokens": after["tokens"]["total"],
}
if __name__ == "__main__":
raise SystemExit(main())