The example works through a queue of questions across separate calls to run_session, each one
standing in for a separate process. A session never receives the previous session’s messages:
only QueueState.notes, a short list of one-line summaries the earlier sessions wrote, the
mechanism Anthropic’s writing calls structured note-taking[1]. MAX_NOTES caps that
list at six, so notes are a compaction, not a growing log: the oldest one is dropped, the same way
Anthropic describes compaction as reinitiating a window from a summary once the old one nears its
limit: this example does it by count instead of by token limit, to keep the code small.
Everything the queue needs to resume lives in one JSON file, QueueState, written by replacing a
temp file rather than overwriting in place, so a session killed mid-write can never hand the next
one a half-written checkpoint. Two real systems do the same job differently. Temporal keeps “a
complete, ordered record of everything that happened in a Workflow Execution”, and a new process
“rebuilds the state of the execution and resumes at the point where it stopped, with local
variables and progress intact”[4]. Inngest persists each step’s own result instead:
“The steps that successfully executed are memoized,” and a retry “is re-executed from the point
of failure with the state of all previous step executions”[5]. QueueState is closer
to Inngest’s shape (one record of what is done) without a workflow engine underneath it.
Letta’s archival memory goes past a flat notes list: “a semantically searchable database where
agents can store facts, knowledge, and information for long-term retrieval”, whose fragments
“must be queried on-demand via tools”[6]. Six notes never need a search index;
thousands would.
The one model-decided step in a session is a single choice between two tools, answer and
flag_for_review. This is the same kind of decision function
calling makes once, not a loop like single
agent’s. What makes this level 7 and not level 4 is everything around that one call: the
trigger that started the session, the checkpoint the session leaves behind, and the fact that
nobody has to be there for either.
examples/long_horizon/run.py · lines 126–182
def run_session(
state_path: Path,
model: Model,
tracer: Tracer,
*,
questions: list[str],
corpus_dir: Path = DEFAULT_CORPUS_DIR,
) -> Answer | None:
"""One scheduled session. Returns the answer it produced, or None if it flagged the question
for a person, or if the queue was already empty."""
tracer.record(kind="code", decided_by="code", title="Scheduler starts a session", detail="no person asked for this run")
state = QueueState.load(state_path, questions=questions)
if not state.queue:
return None
started_from = state.sessions_run # what the checkpoint must still say when this session writes
state.sessions_run += 1
question = state.queue[0]
sections = load_sections(corpus_dir)
sources = [s for s, score in bm25_search(sections, question, k=RETRIEVE_K) if score > 0]
tracer.record(kind="code", decided_by="code", title="Retrieve sources for the next queued question", detail=", ".join(s.cite for s in sources) or "none")
notes_block = "\n".join(f"- {n}" for n in state.notes) or "(no notes yet)"
blocks = "\n\n".join(f"[{s.cite}] {s.title}\n{s.text}" for s in sources)
messages = [
Message(role="system", content=SYSTEM),
Message(role="user", content=f"Notes from earlier sessions:\n{notes_block}\n\nSources:\n\n{blocks}\n\nQuestion: {question}"),
]
tracer.record(kind="code", decided_by="code", title="Rebuild context from notes, not the transcript", detail=f"{len(state.notes)} notes carried forward, no prior session's messages included")
completion = model.complete(messages, tools=TOOLS, max_tokens=300)
call = completion.tool_calls[0] if completion.tool_calls else None
call_desc = f"{call.name}({json.dumps(call.arguments, sort_keys=True)})" if call else "(no tool call)"
tracer.record(
kind="model", decided_by="model", title="Model decides whether to answer or flag this question",
detail=call_desc, tokens_in=completion.tokens_in, tokens_out=completion.tokens_out, ms=completion.ms,
)
if call and call.name == "flag_for_review":
reason = str(call.arguments.get("reason", "unspecified"))
state.pending[question] = reason
state.queue.pop(0)
tracer.record(kind="code", decided_by="code", title="Hand the question to a person", detail=reason)
result = None
else:
text = str(call.arguments.get("text", "")) if call else completion.text
citations = list(call.arguments.get("citations", [])) if call else []
state.answers[question] = {"text": text, "citations": citations}
state.notes.append(f"{question} -> {text[:80]}")
state.notes = state.notes[-MAX_NOTES:]
state.queue.pop(0)
tracer.record(kind="code", decided_by="code", title="Record the answer and compact a note", detail=text[:120])
result = Answer(text=text, citations=citations)
tracer.record(kind="code", decided_by="code", title="Checkpoint the queue to disk", detail=f"{len(state.queue)} left in queue, session {state.sessions_run}")
state.save(state_path, expect_sessions_run=started_from)
return result
Be precise about what that buys. An interrupted session writes nothing, so every finished answer
survives and is recorded once; the question it was working on returns to the queue and is asked
again, so the model call can happen twice. The work is at-least-once, the record is once. That is
safe only because a session’s one effect outside its own memory is the checkpoint: a session
that also sent an email would need the send to be idempotent.
Write-then-replace does not cover two other failures. A checkpoint damaged by anything else
(truncated, hand-edited, the wrong shape) raises a CheckpointError naming the file rather than
starting fresh, because starting over silently would drop the queue, re-answer everything, and
still report success on every tick after. And two overlapping sessions would both load the same
checkpoint, the second erasing the first, so save refuses unless the counter on disk is still
the one the session read. That is a check before a write, not a lock: it catches the overlap, it
does not make concurrent sessions safe.
tests/test_example_long_horizon.py proves each of these: a crash mid-session with the checkpoint
compared byte for byte afterward, a retry that answers the interrupted question once, seven kinds
of damaged checkpoint, and two overlapping sessions where the loser is refused. Run it yourself:
examples/long_horizon/README.md · lines 19–19
python -m examples.long_horizon --model stub:scripted