Проверяем устойчивость stateful-агента к перезапуску

Проверяем устойчивость stateful-агента к перезапуску
Temporary fallback cover; replace in editorial pass.

Testing how a stateful agent survives restarts: durable checkpoints, task resumption and an action journal, all verified with a hard-kill test suite.

Level: advanced · Reading and hands-on time: about 90 minutes · Updated 11 October 2026 · Stack: Python 3.10+, SQLite, pytest, optional Ollama

Long agent tasks often fail the same way. The process stops halfway through. On restart, the agent has lost its context, its intermediate results and any record of what it already did. It either starts over, paying again for model calls and repeating side effects, or it carries on from a guessed position. This guide builds a small local agent stack in which every step is checkpointed and journaled. It then gives you a test suite that kills the process with SIGKILL at nine fixed points and one random point, restarts it, and checks that the task finishes exactly once and that the journal explains what happened.

1. The case: a report agent that gets killed halfway through

Take a small but realistic agent task (this is an illustrative scenario, not a customer report). Every night the agent:

  1. reads a notes file (fetch);
  2. summarizes it, optionally with a local LLM (summarize);
  3. writes a report file that some other system picks up (write_artifact). This step is a side effect, so it changes the world outside the agent.

Suppose the process is killed by an out-of-memory event, a deploy or a laptop lid closing. If that happens between steps 2 and 3, a naive agent that holds its plan and results in memory loses all of it. If it happens after step 3 wrote the file but before the agent recorded that, a naive retry writes a second report. Both failures are quiet. The usual result is a duplicate email, a doubled database row or a second ticket.

A stateful agent fixes this by keeping its plan, step outputs and progress in durable storage outside the process. The fix is only real if you test it, though. The test has to stop the process abruptly at every dangerous boundary.

2. The model: what "survives a restart" means

"Restart resilience" can mean almost anything, so the tests check five concrete properties:

  1. State is preserved. After a kill, the database still holds the task, its plan and the output of every committed step.
  2. The task resumes, not restarts. Steps committed before the kill are not run again.
  3. Interrupted steps are recovered explicitly. A step that was in progress at kill time is moved back to pending, and the journal records that recovery.
  4. Side effects happen exactly once in outcome. The executor is at-least-once: a step may run twice. The side effect is made idempotent with a stable key, so the external result appears once.
  5. The journal is complete and consistent. Every step_started event is matched by exactly one step_finished, step_recovered or step_error, and every process run is visible as its own run_started.

Two design rules make these properties achievable:

Because of this, there are exactly three interesting places to kill a step:

Crash pointDurable state at kill timeCorrect behaviour on restart
before_effectstep is running, tool not calledrecover the step, run it again
after_effectstep is running, tool already produced its effectrecover, run again, and the effect must deduplicate
after_commitstep is doneno recovery, continue with the next step

With three steps, that gives 3 × 3 = 9 deterministic crash scenarios. One more test sends SIGKILL from outside while a step is sleeping, to confirm that the internal crash hooks aren't hiding anything.

3. Setup

You need Linux or macOS, because the code uses fcntl file locks and POSIX signals. You also need Python 3.10 or newer, and optionally the sqlite3 command-line tool for inspection. The agent itself uses only the standard library. pytest is needed only for the tests.

mkdir restart-lab && cd restart-lab
python3 -m venv .venv
. .venv/bin/activate
pip install pytest
mkdir input
printf 'Line one\nLine two\n\nLine three\nLine four\n' > input/notes.txt

The final layout:

restart-lab/
├── agent.py          # CLI + agent loop + recovery
├── store.py          # SQLite schema, transactions, journal
├── tools.py          # fetch / summarize / write_artifact
├── plan.json         # the task plan
├── test_restart.py   # the restart test suite
└── input/notes.txt

Create plan.json:

{
  "goal": "Summarize notes into a report",
  "steps": [
    {"tool": "fetch", "input": {"path": "input/notes.txt"}},
    {"tool": "summarize", "input": {"from": 0}},
    {"tool": "write_artifact", "input": {"from": 1}}
  ]
}

The plan is fixed on purpose. If an LLM produced the plan, it would have to be persisted at task creation for the same reason as everything else. A restart must not re-plan, because a second plan can differ from the first.

4. Durable state and the action journal (store.py)

SQLite is used in WAL mode with synchronous=FULL. With those settings, a committed transaction survives the process being killed. The connection uses isolation_level=None, which lets the code issue BEGIN IMMEDIATE / COMMIT itself, so it is always clear what is atomic.

"""Durable state for the agent: tasks, steps and an append-only event journal."""
import json
import sqlite3
import time

SCHEMA = """
CREATE TABLE IF NOT EXISTS tasks(
  id TEXT PRIMARY KEY,
  goal TEXT NOT NULL,
  status TEXT NOT NULL,                    -- pending | running | done | failed
  created_at REAL NOT NULL,
  updated_at REAL NOT NULL
);
CREATE TABLE IF NOT EXISTS steps(
  task_id TEXT NOT NULL REFERENCES tasks(id),
  idx INTEGER NOT NULL,
  tool TEXT NOT NULL,
  input TEXT NOT NULL,
  status TEXT NOT NULL DEFAULT 'pending',  -- pending | running | done | failed
  output TEXT,
  attempts INTEGER NOT NULL DEFAULT 0,
  PRIMARY KEY(task_id, idx)
);
CREATE TABLE IF NOT EXISTS events(
  seq INTEGER PRIMARY KEY AUTOINCREMENT,
  ts REAL NOT NULL,
  task_id TEXT NOT NULL,
  run_id TEXT NOT NULL,
  kind TEXT NOT NULL,
  step_idx INTEGER,
  data TEXT NOT NULL
);
"""


def connect(path):
    conn = sqlite3.connect(path, isolation_level=None, timeout=10)
    conn.execute("PRAGMA journal_mode=WAL")
    conn.execute("PRAGMA synchronous=FULL")
    conn.execute("PRAGMA foreign_keys=ON")
    conn.executescript(SCHEMA)
    return conn


class tx:
    """BEGIN IMMEDIATE ... COMMIT, or ROLLBACK on exception."""

    def __init__(self, conn):
        self.conn = conn

    def __enter__(self):
        self.conn.execute("BEGIN IMMEDIATE")
        return self.conn

    def __exit__(self, exc_type, exc, tb):
        self.conn.execute("COMMIT" if exc_type is None else "ROLLBACK")
        return False


def log(conn, task_id, run_id, kind, step_idx=None, **data):
    """Append one journal event. Call only inside a tx() block."""
    conn.execute(
        "INSERT INTO events(ts, task_id, run_id, kind, step_idx, data) VALUES(?,?,?,?,?,?)",
        (time.time(), task_id, run_id, kind, step_idx, json.dumps(data, sort_keys=True)),
    )


def create_task(conn, task_id, goal, plan, run_id):
    now = time.time()
    with tx(conn):
        conn.execute(
            "INSERT INTO tasks(id, goal, status, created_at, updated_at) VALUES(?,?,?,?,?)",
            (task_id, goal, "pending", now, now),
        )
        for i, step in enumerate(plan):
            conn.execute(
                "INSERT INTO steps(task_id, idx, tool, input) VALUES(?,?,?,?)",
                (task_id, i, step["tool"], json.dumps(step["input"], sort_keys=True)),
            )
        log(conn, task_id, run_id, "task_created", steps=len(plan))

The journal lives in the same database file as the state. That is deliberate. A separate JSONL file would need its own fsync discipline, and it could disagree with the state after a crash. A JSONL export for humans and log shippers is generated from the table later.

5. Tools and idempotent side effects (tools.py)

Every tool takes the same arguments: (input, ctx, idem_key). Here ctx maps step indexes to the committed outputs of earlier steps. idem_key is stable across retries (<task_id>-<step_idx>), and it is the only thing that keeps the side effect from repeating.

"""Agent tools. Every tool must be safe to call more than once with the same idem_key."""
import hashlib
import json
import os
import urllib.request
from pathlib import Path

ART = Path(os.environ.get("AGENT_ARTIFACTS", "artifacts"))


def fetch(inp, ctx, idem_key):
    text = Path(inp["path"]).read_text(encoding="utf-8")
    return {
        "chars": len(text),
        "sha256": hashlib.sha256(text.encode("utf-8")).hexdigest(),
        "text": text,
    }


def summarize(inp, ctx, idem_key):
    text = ctx[inp["from"]]["text"]
    if os.environ.get("AGENT_LLM") == "ollama":
        body = json.dumps({
            "model": os.environ.get("OLLAMA_MODEL", "llama3.2"),
            "prompt": "Summarize the following notes in three short lines:\n\n" + text,
            "stream": False,
        }).encode("utf-8")
        req = urllib.request.Request(
            os.environ.get("OLLAMA_URL", "http://127.0.0.1:11434") + "/api/generate",
            data=body,
            headers={"Content-Type": "application/json"},
        )
        with urllib.request.urlopen(req, timeout=300) as resp:
            return {"summary": json.load(resp)["response"].strip(), "engine": "ollama"}
    # Deterministic stub: first three non-empty lines.
    lines = [line.strip() for line in text.splitlines() if line.strip()]
    return {"summary": "\n".join(lines[:3]), "engine": "stub"}


def write_artifact(inp, ctx, idem_key):
    content = ctx[inp["from"]]["summary"]
    ART.mkdir(parents=True, exist_ok=True)

    if os.environ.get("AGENT_UNSAFE_EFFECTS") == "1":
        # Deliberately NOT idempotent: used only by the negative-control test.
        log_path = ART / "report.log"
        with open(log_path, "a", encoding="utf-8") as f:
            f.write(content + "\n---\n")
        return {"path": str(log_path), "deduplicated": False}

    final = ART / f"{idem_key}.md"
    tmp = ART / f".{idem_key}.{os.getpid()}.tmp"
    with open(tmp, "w", encoding="utf-8") as f:
        f.write(content)
        f.flush()
        os.fsync(f.fileno())
    try:
        os.link(tmp, final)  # atomic; raises if final already exists
        deduplicated = False
    except FileExistsError:
        deduplicated = True
    finally:
        os.unlink(tmp)
    return {"path": str(final), "deduplicated": deduplicated}

Why write → fsync → link instead of simply open(..., O_EXCL)? With O_EXCL, a kill between creating the file and finishing the write leaves a partial file behind. The retry would then "deduplicate" against that partial file and accept it as the result. os.link publishes the complete file atomically under its final name, or fails because the name already exists. Any partial work is left only in a temp file that is never treated as the result.

For a real external API, the equivalent is to send idem_key as the provider's idempotency key or request ID. If the provider has no such mechanism, the step cannot be made exactly-once from the agent's side. Section 10 covers this.

6. The agent loop with recovery (agent.py)

On every start, the loop runs recovery first. It finds steps left in running by a previous process, resets them to pending and journals that. It then executes pending steps in order. Crash hooks controlled by AGENT_CRASH_AT=<point>:<idx> send SIGKILL to the agent's own process. With SIGKILL there is no exception, no finally block and no clean shutdown, which is the condition the test needs.

A per-task lock (fcntl.flock) stops two processes from running the same task at once. The kernel releases it when the process dies, even after SIGKILL, so a crash cannot leave a stale lock.

"""Minimal stateful agent: durable steps, crash recovery, append-only journal."""
import argparse
import fcntl
import json
import os
import signal
import sys
import time
import uuid

import store
import tools

TOOLS = {
    "fetch": tools.fetch,
    "summarize": tools.summarize,
    "write_artifact": tools.write_artifact,
}


def crash_point(name, idx):
    if os.environ.get("AGENT_CRASH_AT") == f"{name}:{idx}":
        os.kill(os.getpid(), signal.SIGKILL)


def recover(conn, task_id, run_id):
    with store.tx(conn):
        rows = conn.execute(
            "SELECT idx FROM steps WHERE task_id=? AND status='running' ORDER BY idx",
            (task_id,),
        ).fetchall()
        for (idx,) in rows:
            conn.execute(
                "UPDATE steps SET status='pending' WHERE task_id=? AND idx=?", (task_id, idx)
            )
            store.log(conn, task_id, run_id, "step_recovered", idx)
        store.log(conn, task_id, run_id, "run_started", recovered=[r[0] for r in rows])


def set_task(conn, task_id, status):
    conn.execute(
        "UPDATE tasks SET status=?, updated_at=? WHERE id=?", (status, time.time(), task_id)
    )


def run(conn, task_id, run_id):
    task = conn.execute("SELECT status FROM tasks WHERE id=?", (task_id,)).fetchone()
    if task is None:
        print(f"unknown task {task_id}", file=sys.stderr)
        return 1
    if task[0] in ("done", "failed"):
        return 0 if task[0] == "done" else 2

    max_attempts = int(os.environ.get("AGENT_MAX_ATTEMPTS", "3"))
    delay = float(os.environ.get("AGENT_STEP_DELAY", "0"))
    recover(conn, task_id, run_id)

    while True:
        row = conn.execute(
            "SELECT idx, tool, input, attempts FROM steps "
            "WHERE task_id=? AND status='pending' ORDER BY idx LIMIT 1",
            (task_id,),
        ).fetchone()
        if row is None:
            with store.tx(conn):
                set_task(conn, task_id, "done")
                store.log(conn, task_id, run_id, "task_done")
            return 0

        idx, tool, inp, attempts = row
        if attempts >= max_attempts:
            with store.tx(conn):
                conn.execute(
                    "UPDATE steps SET status='failed' WHERE task_id=? AND idx=?", (task_id, idx)
                )
                set_task(conn, task_id, "failed")
                store.log(conn, task_id, run_id, "task_failed", idx, reason="max_attempts")
            return 2

        with store.tx(conn):
            conn.execute(
                "UPDATE steps SET status='running', attempts=attempts+1 "
                "WHERE task_id=? AND idx=?",
                (task_id, idx),
            )
            set_task(conn, task_id, "running")
            store.log(conn, task_id, run_id, "step_started", idx, tool=tool, attempt=attempts + 1)

        ctx = {
            i: json.loads(out)
            for i, out in conn.execute(
                "SELECT idx, output FROM steps WHERE task_id=? AND status='done'", (task_id,)
            )
        }
        idem_key = f"{task_id}-{idx}"

        crash_point("before_effect", idx)
        time.sleep(delay)
        try:
            out = TOOLS[tool](json.loads(inp), ctx, idem_key)
        except Exception as exc:
            with store.tx(conn):
                conn.execute(
                    "UPDATE steps SET status='pending' WHERE task_id=? AND idx=?", (task_id, idx)
                )
                store.log(conn, task_id, run_id, "step_error", idx, error=repr(exc))
            continue
        crash_point("after_effect", idx)

        with store.tx(conn):
            conn.execute(
                "UPDATE steps SET status='done', output=? WHERE task_id=? AND idx=?",
                (json.dumps(out, sort_keys=True), task_id, idx),
            )
            store.log(conn, task_id, run_id, "step_finished", idx, tool=tool)
        crash_point("after_commit", idx)


def status(conn, task_id):
    task = conn.execute("SELECT id, goal, status FROM tasks WHERE id=?", (task_id,)).fetchone()
    if task is None:
        print(f"unknown task {task_id}", file=sys.stderr)
        return 1
    steps = [
        {"idx": i, "tool": t, "status": s, "attempts": a}
        for i, t, s, a in conn.execute(
            "SELECT idx, tool, status, attempts FROM steps WHERE task_id=? ORDER BY idx",
            (task_id,),
        )
    ]
    print(json.dumps({"task": task[0], "goal": task[1], "status": task[2], "steps": steps}, indent=2))
    return 0


def export(conn, task_id):
    for seq, ts, run_id, kind, idx, data in conn.execute(
        "SELECT seq, ts, run_id, kind, step_idx, data FROM events WHERE task_id=? ORDER BY seq",
        (task_id,),
    ):
        record = {"seq": seq, "ts": ts, "run": run_id, "kind": kind, "step": idx}
        record.update(json.loads(data))
        print(json.dumps(record, ensure_ascii=False))
    return 0


def acquire_lock(db, task_id):
    handle = open(f"{db}.{task_id}.lock", "w")
    try:
        fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB)
    except BlockingIOError:
        handle.close()
        return None
    return handle


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--db", default="state.db")
    sub = parser.add_subparsers(dest="cmd", required=True)
    c = sub.add_parser("create")
    c.add_argument("task_id")
    c.add_argument("plan")
    for name in ("run", "status", "export"):
        sub.add_parser(name).add_argument("task_id")
    args = parser.parse_args()

    conn = store.connect(args.db)
    run_id = uuid.uuid4().hex[:8]

    if args.cmd == "create":
        with open(args.plan, encoding="utf-8") as f:
            plan = json.load(f)
        store.create_task(conn, args.task_id, plan["goal"], plan["steps"], run_id)
        return 0
    if args.cmd == "run":
        lock = acquire_lock(args.db, args.task_id)
        if lock is None:
            print(f"task {args.task_id} is already running", file=sys.stderr)
            return 3
        return run(conn, args.task_id, run_id)
    if args.cmd == "status":
        return status(conn, args.task_id)
    return export(conn, args.task_id)


if __name__ == "__main__":
    sys.exit(main())

Exit codes: 0 means done, 1 means unknown task, 2 means failed after the maximum number of attempts, and 3 means another process holds the lock. When the agent is killed by a signal, Python's subprocess reports -9, and bash reports 137 (128 + 9).

7. Manual run: kill, inspect, resume

Do this once by hand before automating it. Watching the state change makes the assertions in the next section easier to follow.

python agent.py create t1 plan.json

# Kill the agent right after the report file was written,
# but before the step was marked done.
AGENT_CRASH_AT=after_effect:2 python agent.py run t1; echo "exit=$?"

Bash should report exit=137, probably with a "Killed" line from the shell. Now look at what survived:

python agent.py status t1
sqlite3 state.db "SELECT seq, run_id, kind, step_idx FROM events ORDER BY seq;"
ls -la artifacts/

From the code, you should see the following. Steps 0 and 1 are done. Step 2 is running with attempts = 1, and the task is running. The journal ends with step_started for step 2 and has no matching finish event. artifacts/t1-2.md already exists. This is the most dangerous state: the effect happened, but the agent has no durable record that it did.

Resume:

python agent.py run t1; echo "exit=$?"
python agent.py status t1
python agent.py export t1 > journal.jsonl
cat journal.jsonl
ls artifacts/
sqlite3 state.db "SELECT output FROM steps WHERE task_id='t1' AND idx=2;"

Expected: exit=0, and all three steps done. Step 2 has attempts = 2. The journal contains a second run_started with a different run id, a step_recovered for step 2, a second step_started for step 2, then step_finished and task_done. There is still exactly one t1-2.md, and the stored output for step 2 contains "deduplicated": true. Steps 0 and 1 were not run again, which you can tell because they have no second step_started.

To repeat the experiment, delete state.db*, the lock files and artifacts/, then create the task again.

8. The automated restart test suite (test_restart.py)

The suite runs the agent as a real subprocess, so the kills are real. Each test gets a fresh temporary directory. Environment variables starting with AGENT_ are removed from the parent environment so that a stray setting in your shell can't change the results.

"""Restart-resilience tests: kill the agent at fixed points and verify recovery."""
import json
import os
import signal
import sqlite3
import subprocess
import sys
import time
from pathlib import Path

import pytest

AGENT = Path(__file__).resolve().parent / "agent.py"
PLAN = {
    "goal": "Summarize notes into a report",
    "steps": [
        {"tool": "fetch", "input": {"path": "input/notes.txt"}},
        {"tool": "summarize", "input": {"from": 0}},
        {"tool": "write_artifact", "input": {"from": 1}},
    ],
}
EXPECTED = "Line one\nLine two\nLine three"


def agent_cmd(wd, *args):
    return [sys.executable, str(AGENT), "--db", str(wd / "state.db"), *args]


def agent_env(wd, **extra):
    env = {k: v for k, v in os.environ.items() if not k.startswith("AGENT_")}
    env["AGENT_ARTIFACTS"] = str(wd / "artifacts")
    env.update(extra)
    return env


def run_agent(wd, *args, **extra):
    return subprocess.run(
        agent_cmd(wd, *args), cwd=wd, env=agent_env(wd, **extra),
        capture_output=True, text=True, timeout=60,
    )


def query(wd, sql):
    conn = sqlite3.connect(wd / "state.db")
    try:
        return conn.execute(sql).fetchall()
    finally:
        conn.close()


def events(wd):
    rows = query(wd, "SELECT kind, step_idx, data FROM events WHERE task_id='t1' ORDER BY seq")
    return [(k, i, json.loads(d)) for k, i, d in rows]


def kinds(evts, kind):
    return [i for k, i, _ in evts if k == kind]


def step_output(wd, idx):
    (out,) = query(wd, f"SELECT output FROM steps WHERE task_id='t1' AND idx={int(idx)}")[0]
    return json.loads(out)


def artifacts(wd):
    folder = wd / "artifacts"
    return sorted(folder.glob("*.md")) if folder.exists() else []


def wait_for(predicate, timeout=10.0):
    deadline = time.monotonic() + timeout
    while time.monotonic() < deadline:
        if predicate():
            return True
        time.sleep(0.05)
    return False


@pytest.fixture
def wd(tmp_path):
    (tmp_path / "input").mkdir()
    (tmp_path / "input" / "notes.txt").write_text(
        "Line one\nLine two\n\nLine three\nLine four\n", encoding="utf-8"
    )
    (tmp_path / "plan.json").write_text(json.dumps(PLAN), encoding="utf-8")
    result = run_agent(tmp_path, "create", "t1", "plan.json")
    assert result.returncode == 0, result.stderr
    return tmp_path


def assert_completed_once(wd):
    (task,) = query(wd, "SELECT status FROM tasks WHERE id='t1'")[0]
    steps = [s for (s,) in query(wd, "SELECT status FROM steps WHERE task_id='t1' ORDER BY idx")]
    assert task == "done"
    assert steps == ["done", "done", "done"]

    evts = events(wd)
    assert sorted(kinds(evts, "step_finished")) == [0, 1, 2]
    assert len(kinds(evts, "task_done")) == 1
    # Journal consistency: every start is closed by exactly one outcome.
    closed = (len(kinds(evts, "step_finished")) + len(kinds(evts, "step_recovered"))
              + len(kinds(evts, "step_error")))
    assert len(kinds(evts, "step_started")) == closed

    files = artifacts(wd)
    assert [f.name for f in files] == ["t1-2.md"]
    assert files[0].read_text(encoding="utf-8") == EXPECTED
    return evts


def test_baseline_without_crash(wd):
    result = run_agent(wd, "run", "t1")
    assert result.returncode == 0, result.stderr
    evts = assert_completed_once(wd)
    assert kinds(evts, "step_recovered") == []
    assert len(kinds(evts, "run_started")) == 1


@pytest.mark.parametrize("point", ["before_effect", "after_effect", "after_commit"])
@pytest.mark.parametrize("idx", [0, 1, 2])
def test_crash_and_resume(wd, point, idx):
    first = run_agent(wd, "run", "t1", AGENT_CRASH_AT=f"{point}:{idx}")
    assert first.returncode == -signal.SIGKILL

    second = run_agent(wd, "run", "t1")
    assert second.returncode == 0, second.stderr

    evts = assert_completed_once(wd)
    assert len(kinds(evts, "run_started")) == 2
    interrupted = point != "after_commit"
    assert kinds(evts, "step_recovered") == ([idx] if interrupted else [])
    assert kinds(evts, "step_started").count(idx) == (2 if interrupted else 1)
    for other in {0, 1, 2} - {idx}:
        assert kinds(evts, "step_started").count(other) == 1
    if (point, idx) == ("after_effect", 2):
        assert step_output(wd, 2)["deduplicated"] is True


def test_external_sigkill_mid_step(wd):
    proc = subprocess.Popen(
        agent_cmd(wd, "run", "t1"), cwd=wd, env=agent_env(wd, AGENT_STEP_DELAY="1.0"),
        stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
    )
    try:
        assert wait_for(lambda: 1 in kinds(events(wd), "step_started"))
        proc.send_signal(signal.SIGKILL)
    finally:
        proc.wait(timeout=15)
    assert proc.returncode == -signal.SIGKILL

    result = run_agent(wd, "run", "t1")
    assert result.returncode == 0, result.stderr
    evts = assert_completed_once(wd)
    assert kinds(evts, "step_recovered") == [1]


def test_second_runner_is_rejected(wd):
    proc = subprocess.Popen(
        agent_cmd(wd, "run", "t1"), cwd=wd, env=agent_env(wd, AGENT_STEP_DELAY="1.0"),
        stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
    )
    try:
        assert wait_for(lambda: kinds(events(wd), "run_started"))
        result = run_agent(wd, "run", "t1")
        assert result.returncode == 3
    finally:
        proc.kill()
        proc.wait(timeout=15)


def test_negative_control_unsafe_effect_is_duplicated(wd):
    """Proves the suite can detect duplication: a non-idempotent effect must repeat."""
    first = run_agent(wd, "run", "t1", AGENT_CRASH_AT="after_effect:2", AGENT_UNSAFE_EFFECTS="1")
    assert first.returncode == -signal.SIGKILL
    second = run_agent(wd, "run", "t1", AGENT_UNSAFE_EFFECTS="1")
    assert second.returncode == 0, second.stderr
    log = (wd / "artifacts" / "report.log").read_text(encoding="utf-8")
    assert log.count("\n---\n") == 2

Run it:

pytest -v test_restart.py

The suite collects 13 tests: 1 baseline, 9 parametrized crash scenarios, 1 external kill, 1 lock test and 1 negative control. The two tests that use AGENT_STEP_DELAY=1.0 add a few seconds of wall time each. The rest depend only on process startup speed.

Why the negative control matters

A restart test that never fails proves little. test_negative_control_unsafe_effect_is_duplicated switches the report writer to a plain append and asserts that the duplicate appears. If that test ever fails, your crash injection isn't reaching the dangerous window (for example, because a refactor moved crash_point("after_effect", ...)). In that case, the green results from the other tests can't be trusted.

9. How to read a passing and a failing run

When the whole suite is green, the five properties from section 2 hold for this plan on your machine:

Typical red results and what they mean:

Failing assertionLikely cause
first.returncode == -9 fails with 0The crash hook didn't fire. Check the spelling of AGENT_CRASH_AT and whether the hook is still placed where the test expects.
Two .md files, or a modified artifactThe idempotency key isn't stable across runs (for example, it includes run_id or a timestamp), or the writer isn't atomic.
step_started count > 1 for a step that wasn't interruptedCommitted outputs are being ignored. The loop is selecting by something other than status='pending', or recovery resets done steps.
Steps stuck in running after the second runRecovery didn't run, or it ran in a transaction that was rolled back.
Journal consistency assertion failsA state change was committed without its event, or an event was written outside the state transaction.
database is locked in stderrTwo writers overlapped. Check that the lock test passes and that nothing else holds the DB open in write mode.

10. Failure cases and what catches them

Not every failure is covered by this suite. Here is where the boundary lies:

FailureCovered here?Notes
Process killed (SIGKILL, OOM killer, docker kill)YesThe core of the suite.
Graceful stop (SIGTERM, Ctrl+C)IndirectlyA graceful stop is a subset of a hard kill as long as there are no shutdown handlers that write state. If you add such handlers, add a SIGTERM variant of the external-kill test.
Duplicate concurrent runnerYestest_second_runner_is_rejected. Valid on one host only; see limitations.
Tool raises an exceptionPartlyThe loop journals step_error and retries up to AGENT_MAX_ATTEMPTS. Add a test with a tool that always fails and assert exit code 2 and task_failed.
Input changed between runsBy design, not testedThe committed fetch output (with its sha256) is reused after restart. That is the right behaviour for resumption, but a stale-data check, if you need one, has to be written explicitly.
Non-idempotent external APINo, only demonstratedThe negative control shows the duplication. The fix depends on the provider: idempotency keys, a lookup by client reference before retrying, or marking the step as "manual confirmation required" after a recovery.
Power loss / OS crashNoSIGKILL leaves the OS page cache intact. Power loss doesn't. synchronous=FULL and fsync help, but the atomic link also needs an fsync of the directory to be durable against power loss. Test this in a VM that you hard-reset.
Corrupt or deleted databaseNoThis calls for backups (sqlite3 state.db ".backup ..."), not recovery logic.

11. Adding a local model with Ollama

The stub summarizer keeps the tests deterministic. To run the same stack against a local model, install Ollama, pull a model and switch the engine. The model name below is an example, so use any model you have pulled:

ollama pull llama3.2
rm -rf state.db* artifacts/
python agent.py create t2 plan.json
AGENT_LLM=ollama OLLAMA_MODEL=llama3.2 AGENT_CRASH_AT=after_commit:1 python agent.py run t2
AGENT_LLM=ollama OLLAMA_MODEL=llama3.2 python agent.py run t2
python agent.py export t2

What changes when a model is involved:

Don't put live model calls in the automated suite. Nondeterministic outputs make content assertions flaky, and a missing local server would turn into false failures. Keep model runs as a manual or nightly smoke check, separate from the restart test.

12. Limitations

13. Next steps

When the suite is green, keep it as a regression test. Any change to the agent loop, the tools or the storage layer should be followed by a full pytest -v test_restart.py run. Whenever you add a tool, add its crash points to the parametrization. If the tool has an external side effect, also add a negative control that proves a duplicate would be caught.

For related work, see the agent memory guide and testing AI agents. Browse all practical material in the guides section, and look up unfamiliar terms in the glossary.

We publish what works for us—and implement the same solutions for your business. We design AI automation, Telegram bots, chats, and AI agents for real-world processes. Discuss your project →