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

Долгая задача агента может прерваться в любой момент: процесс упал, контейнер пересоздали, ноутбук закрыли. Если состояние хранится только в памяти, после такой остановки пропадают контекст, промежуточные результаты и история действий. Тогда агент либо начинает всё сначала, либо, что хуже, повторяет уже выполненные действия. В этой статье мы соберём небольшой локальный стек и напишем воспроизводимый тест. Тест жёстко останавливает агента посреди работы, запускает его снова и проверяет три вещи: что состояние сохранилось, что задача восстановилась и что журнал действий остался целым.
Что именно проверяем
Stateful-агент хранит состояние задачи между шагами: план, статус каждого шага, результаты вызовов инструментов и журнал. Чтобы агент пережил перезапуск, нужны четыре свойства:
- Сохранность состояния. Шаг, который отмечен как выполненный, после перезапуска остаётся выполненным, и его результат не теряется.
- Восстановление задачи. Новый процесс продолжает работу с первого незавершённого шага и не пересоздаёт план.
- Целостность журнала. По журналу видно, когда задачу создали, когда возобновили и чем она закончилась. Ни один шаг не записан как завершённый дважды.
- Эквивалентность результата. Итоговый отчёт после сбоя и перезапуска совпадает с отчётом прогона без сбоев.
Тестировать будем самый неприятный сценарий — SIGKILL. Процесс не получает шанса что-то дописать или закрыть, поэтому обработчики завершения и блоки finally не помогут. Если агент переживает kill -9, мягкая остановка ему тоже не страшна.
Стек и допущения
- Python 3.10 или новее. Нужна только стандартная библиотека:
sqlite3,subprocess,unittest. - SQLite в режиме WAL — источник истины для состояния задачи.
- Журнал действий в формате JSONL, запись только в конец файла с
fsyncпосле каждого события. Это журнал для аудита и диагностики, а не для управления агентом. - Linux или macOS. Ручная проверка использует
timeoutиз GNU coreutils (на macOS команда называетсяgtimeout).
Вместо обращения к LLM инструмент агента здесь детерминированный: он считает слова и хэш текстовых файлов. Так сделано намеренно. Тест проверяет механику персистентности, и ответы модели, которые меняются от запуска к запуску, в нём только мешали бы. Как подключить локальную модель, описано в конце статьи.
Шаг 1. Рабочий каталог
mkdir -p ~/restart-lab/inbox
cd ~/restart-lab
for i in 1 2 3 4 5; do echo "заметка номер $i для проверки перезапуска" > inbox/note_$i.txt; done
ls inbox
Команды создают только каталог ~/restart-lab и файлы внутри него. Чтобы всё убрать, достаточно удалить этот каталог.
Шаг 2. Агент с чекпоинтами и журналом
Сохраните файл как agent.py. Код ниже — учебный пример, а не готовая библиотека.
"""Учебный stateful-агент: план и результаты в SQLite, журнал действий в JSONL."""
import argparse
import hashlib
import json
import os
import sqlite3
import time
from pathlib import Path
SCHEMA = """
CREATE TABLE IF NOT EXISTS task (
id TEXT PRIMARY KEY,
goal TEXT NOT NULL,
status TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS step (
task_id TEXT NOT NULL,
idx INTEGER NOT NULL,
name TEXT NOT NULL,
status TEXT NOT NULL,
result TEXT,
PRIMARY KEY (task_id, idx)
);
"""
def connect(path):
con = sqlite3.connect(path, isolation_level=None)
con.execute("PRAGMA journal_mode=WAL")
con.execute("PRAGMA synchronous=FULL")
con.executescript(SCHEMA)
return con
def log_event(path, **event):
event["ts"] = round(time.time(), 3)
with open(path, "a", encoding="utf-8") as f:
f.write(json.dumps(event, ensure_ascii=False) + "\n")
f.flush()
os.fsync(f.fileno())
def seal_journal(path):
"""Если прошлый процесс оборвал строку, закрываем её переводом строки."""
p = Path(path)
if p.exists() and p.stat().st_size:
with open(p, "rb+") as f:
f.seek(-1, os.SEEK_END)
if f.read(1) != b"\n":
f.write(b"\n")
def read_journal(path):
events, broken = [], 0
p = Path(path)
if not p.exists():
return events, broken
for line in p.read_text(encoding="utf-8", errors="replace").splitlines():
if not line.strip():
continue
try:
events.append(json.loads(line))
except json.JSONDecodeError:
broken += 1 # строка, оборванная аварийной остановкой
return events, broken
def summarize(path):
text = Path(path).read_text(encoding="utf-8")
digest = hashlib.sha256(text.encode("utf-8")).hexdigest()[:12]
return {"file": Path(path).name, "words": len(text.split()), "sha256": digest}
def reconcile(con, a):
"""Шаг завершён в БД, но в журнал не попал: дописываем событие с пометкой."""
events, _ = read_journal(a.log)
logged = {e.get("idx") for e in events
if e.get("task") == a.task and e.get("event") == "step_done"}
done = con.execute(
"SELECT idx FROM step WHERE task_id = ? AND status = 'done' ORDER BY idx",
(a.task,),
).fetchall()
for (idx,) in done:
if idx not in logged:
log_event(a.log, task=a.task, event="step_done", idx=idx, reconciled=True)
def run(a):
seal_journal(a.log)
con = connect(a.db)
row = con.execute("SELECT status FROM task WHERE id = ?", (a.task,)).fetchone()
if row is None:
names = [p.name for p in sorted(Path(a.inbox).glob("*.txt"))]
con.execute("BEGIN IMMEDIATE")
con.execute("INSERT INTO task VALUES (?, ?, 'running')",
(a.task, "summarize " + str(a.inbox)))
con.executemany(
"INSERT INTO step VALUES (?, ?, ?, 'pending', NULL)",
[(a.task, i, n) for i, n in enumerate(names)],
)
con.execute("COMMIT")
log_event(a.log, task=a.task, event="task_created", steps=len(names))
elif row[0] == "done":
log_event(a.log, task=a.task, event="already_done")
return
else:
log_event(a.log, task=a.task, event="task_resumed")
reconcile(con, a)
pending = con.execute(
"SELECT idx, name FROM step WHERE task_id = ? AND status != 'done' ORDER BY idx",
(a.task,),
).fetchall()
for idx, name in pending:
log_event(a.log, task=a.task, event="step_started", idx=idx, name=name)
result = summarize(Path(a.inbox) / name)
time.sleep(a.delay) # окно, в которое тест «роняет» процесс
con.execute(
"UPDATE step SET status = 'done', result = ? WHERE task_id = ? AND idx = ?",
(json.dumps(result, ensure_ascii=False), a.task, idx),
)
log_event(a.log, task=a.task, event="step_done", idx=idx)
rows = con.execute(
"SELECT result FROM step WHERE task_id = ? ORDER BY idx", (a.task,)
).fetchall()
tmp = Path(a.out).with_suffix(".tmp")
tmp.write_text(json.dumps([json.loads(r[0]) for r in rows],
ensure_ascii=False, indent=2), encoding="utf-8")
os.replace(tmp, a.out)
con.execute("UPDATE task SET status = 'done' WHERE id = ?", (a.task,))
log_event(a.log, task=a.task, event="task_done")
def main():
p = argparse.ArgumentParser()
p.add_argument("--task", required=True)
p.add_argument("--inbox", required=True)
p.add_argument("--db", default="state.db")
p.add_argument("--log", default="actions.jsonl")
p.add_argument("--out", default="report.json")
p.add_argument("--delay", type=float, default=0.0)
run(p.parse_args())
if __name__ == "__main__":
main()
На чём держится восстановление
- План фиксируется один раз. Список шагов записывается в одной транзакции вместе с задачей. При перезапуске агент не сканирует
inboxзаново, поэтому новые файлы не сдвинут нумерацию шагов. - Сначала результат, потом статус. Результат и статус
doneзаписываются однимUPDATE, то есть атомарно. Если процесс упадёт до этой записи, шаг выполнится повторно. Это семантика «как минимум один раз» (at-least-once), и она допустима только для идемпотентных инструментов. - Журнал вторичен. Сбой между коммитом в SQLite и записью в журнал возможен. Функция
reconcileнаходит такие шаги и дописывает для них событие с флагомreconciled, поэтому журнал остаётся полным. - Отчёт пишется атомарно. Сначала создаётся временный файл, затем
os.replaceподменяет им итоговый. Наполовину записанногоreport.jsonне бывает.
Шаг 3. Ручная проверка с жёсткой остановкой
Сначала проверим механику вручную. Каждый шаг длится около секунды, а через три секунды timeout отправляет процессу SIGKILL:
timeout --signal=KILL 3 python3 agent.py --task demo --inbox inbox --delay 1
echo "код выхода: $?"
python3 -c "import sqlite3; print(sqlite3.connect('state.db').execute('SELECT idx, status FROM step ORDER BY idx').fetchall())"
Код выхода должен быть ненулевым: для GNU timeout с сигналом KILL это 137. В выводе SQLite часть шагов будет со статусом done, остальные — pending. Сколько именно шагов успеет завершиться, зависит от скорости машины. Теперь перезапускаем агента с той же задачей:
python3 agent.py --task demo --inbox inbox
cat report.json
tail -n 15 actions.jsonl
В журнале после первого блока событий должно появиться task_resumed, затем step_started и step_done для оставшихся шагов и в конце task_done. Шаг, прерванный посреди работы, получит два события step_started и одно step_done. Так в журнале и выглядит семантика at-least-once.
Шаг 4. Воспроизводимый тест
У ручной проверки результат зависит от тайминга. Автоматический тест убирает эту зависимость: он ждёт в журнале конкретное событие и только после него останавливает процесс. Сохраните файл как test_restart.py рядом с agent.py:
import json
import sqlite3
import subprocess
import sys
import tempfile
import time
import unittest
from collections import Counter
from pathlib import Path
from agent import read_journal
AGENT = Path(__file__).with_name("agent.py")
STEPS = 6
def agent_cmd(work, task, delay, prefix):
return [sys.executable, str(AGENT), "--task", task,
"--inbox", str(work / "inbox"),
"--db", str(work / (prefix + ".db")),
"--log", str(work / (prefix + ".jsonl")),
"--out", str(work / (prefix + ".json")),
"--delay", str(delay)]
class RestartTest(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.work = Path(self._tmp.name)
inbox = self.work / "inbox"
inbox.mkdir()
for i in range(STEPS):
text = " ".join(["слово"] * (i + 3))
(inbox / f"note_{i}.txt").write_text(text, encoding="utf-8")
def tearDown(self):
self._tmp.cleanup()
def wait_steps(self, log, n, timeout=30):
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
events, _ = read_journal(log)
if sum(e.get("event") == "step_done" for e in events) >= n:
return
time.sleep(0.05)
self.fail("агент не дошёл до нужного шага")
def test_survives_kill(self):
w = self.work
# 1. Эталон: тот же агент без сбоев
subprocess.run(agent_cmd(w, "ref", 0, "ref"), check=True, timeout=60)
reference = json.loads((w / "ref.json").read_text(encoding="utf-8"))
# 2. Медленный прогон и SIGKILL после двух завершённых шагов
proc = subprocess.Popen(agent_cmd(w, "t1", 0.5, "run"))
self.wait_steps(w / "run.jsonl", 2)
proc.kill()
proc.wait(timeout=10)
self.assertNotEqual(proc.returncode, 0)
self.assertFalse((w / "run.json").exists())
con = sqlite3.connect(w / "run.db")
done = con.execute("SELECT COUNT(*) FROM step WHERE status = 'done'").fetchone()[0]
con.close()
self.assertGreaterEqual(done, 2)
self.assertLess(done, STEPS)
# 3. Перезапуск той же задачи
subprocess.run(agent_cmd(w, "t1", 0, "run"), check=True, timeout=60)
report = json.loads((w / "run.json").read_text(encoding="utf-8"))
self.assertEqual(report, reference)
# 4. Журнал действий
events, broken = read_journal(w / "run.jsonl")
kinds = Counter(e.get("event") for e in events)
self.assertEqual(kinds["task_created"], 1)
self.assertEqual(kinds["task_resumed"], 1)
self.assertEqual(kinds["task_done"], 1)
self.assertGreaterEqual(kinds["step_started"], STEPS)
done_idx = Counter(e["idx"] for e in events if e.get("event") == "step_done")
self.assertEqual(sorted(done_idx), list(range(STEPS)))
self.assertTrue(all(c == 1 for c in done_idx.values()), done_idx)
self.assertLessEqual(broken, 1)
# 5. Повторный запуск завершённой задачи ничего не выполняет
subprocess.run(agent_cmd(w, "t1", 0, "run"), check=True, timeout=60)
after, _ = read_journal(w / "run.jsonl")
self.assertEqual([e.get("event") for e in after[len(events):]], ["already_done"])
if __name__ == "__main__":
unittest.main()
Запуск:
cd ~/restart-lab
python3 -m unittest -v test_restart.py
Тест работает во временном каталоге и после себя его удаляет. Ваши файлы в ~/restart-lab он не трогает.
Как читать результат
Если тест прошёл, unittest выведет test_survives_kill ... ok и итоговое OK. Вот что за этим стоит:
- после
SIGKILLв базе сохранились как минимум два завершённых шага и при этом не все: процесс действительно прервали посреди задачи; - частичный отчёт так и не появился на диске;
- после перезапуска отчёт полностью совпал с эталонным прогоном без сбоев;
- в журнале одно создание задачи, одно возобновление и одно завершение, а каждый шаг завершён ровно один раз;
- повторный запуск завершённой задачи не выполнил ни одного шага.
Чтобы убедиться, что тест действительно ловит регрессии, сломайте агента нарочно. Например, в ветке task_resumed выбирайте все шаги, а не только незавершённые. Тест должен упасть на проверке done_idx, потому что часть шагов окажется завершённой дважды. Тест, который ни разу не падал, ничего не доказывает.
Типовые ошибки
- Статус пишется раньше результата. Если отметить шаг
done, а результат сохранить отдельной операцией, после сбоя шаг будет выглядеть выполненным, но без данных. Пишите статус и результат в одной транзакции. - План пересчитывается при старте. Если пока агент лежал, входные данные изменились, индексы шагов разъедутся с сохранёнными результатами. Фиксируйте план при создании задачи.
- Неидемпотентные инструменты. Отправка письма, платёж или запись во внешний API при повторе выполнятся второй раз. Таким инструментам нужен ключ идемпотентности, производный от
task_idиidx, или паттерн outbox: намерение сначала записывается в БД, а отправка подтверждается отдельно. - Журнал без
fsync. Данные зависают в буфере и пропадают при сбое ОС. Противkill -9хватит иflush, но для аудита журнал должен попадать на диск. - Парсер журнала падает на оборванной строке. Аварийная остановка может оставить неполную последнюю строку. Без
seal_journalследующее событие приклеится к ней, и испорченными окажутся уже две записи. - При копировании переносят только
state.db. В режиме WAL часть данных лежит вstate.db-wal. Для бэкапа используйтеsqlite3APIbackupили копируйте остановленную базу вместе со служебными файлами. - Вместо
SIGKILLтестируютSIGTERM. Мягкий сигнал даёт процессу время всё аккуратно закрыть и скрывает проблемы, которые проявятся при настоящем падении.
Подключение локальной LLM
Если заменить summarize на вызов локальной модели, например через HTTP API локального сервера, механика не изменится. Добавятся два правила:
- сохраняйте ответ модели в
step.resultдо того, как он повлияет на следующие шаги. Тогда после перезапуска агент возьмёт уже полученный ответ и не будет генерировать новый; - историю диалога (системный промпт, сообщения, вызовы инструментов) храните в отдельной таблице с порядковым номером и восстанавливайте контекст из неё, а не из памяти процесса.
Ответы модели могут отличаться от запуска к запуску, поэтому сравнение с эталоном придётся ослабить. Проверяйте структуру отчёта, набор обработанных шагов и инварианты журнала, а не дословный текст.
Ограничения подхода
- Тест моделирует падение процесса, а не отключение питания или сбой ОС. Насколько надёжен
fsync, зависит от файловой системы и окружения: в контейнерах и на сетевых дисках гарантии бывают слабее. - Агент однопроцессный. Если задачу могут взять несколько воркеров, нужна аренда задачи с таймаутом (lease) и проверка владельца при каждом коммите.
- Повтор прерванного шага заложен в саму схему. Если повторное выполнение какого-то шага недопустимо, одних чекпоинтов мало: нужна идемпотентность на стороне инструмента.
- Один сценарий остановки не покрывает все окна сбоя. Чтобы проверить точки «после коммита, но до записи в журнал» и «после отчёта, но до
task_done», добавьте в агента явные точки внедрения сбоя и параметризуйте тест. - На Windows
proc.kill()вызываетTerminateProcess, а ручная проверка черезtimeoutне подойдёт. Основной тест наunittestпереносим, но рассчитан в первую очередь на POSIX.
Что дальше
Добавьте этот тест в регрессионный набор агента и запускайте его при каждом изменении схемы состояния или логики шагов. Другие практические материалы о сборке и проверке агентов собраны в разделе гайдов, а определения терминов из статьи — в глоссарии.