«Память агента» часто обозначает четыре несовместимые вещи. Из-за этого историю чата кладут в vector DB, summary принимают за источник истины, а checkpoint — за сериализованный prompt. Затем процесс падает после внешнего действия и при resume повторяет его.
Четыре слоя с разными контрактами
- RAG отвечает: какие внешние документы релевантны текущему запросу. Истина живёт в исходной системе; индекс — производная проекция.
- Context compaction отвечает: как уложить рабочую историю run в окно модели. Summary — lossy cache, а не durable state.
- Checkpoint state отвечает: с какого детерминированного шага продолжить workflow после падения. Это транзакционные данные оркестратора.
- Long-term memory отвечает: какие подтверждённые сведения между runs разрешено сохранить о пользователе или процессе. У каждой записи есть provenance, scope, TTL и invalidation.
source systems ──indexer──> retrieval index ──retrieve──> context builder
┌──────────┘
conversation events ──compactor──> summary ───┤──> model
checkpoint store ────────────────> exact state┘
validated observations ──────────> long-term memoryНе допускайте обратной стрелки summary → source of truth. Модельный пересказ не должен становиться фактом только потому, что пережил несколько turns.
Durable run: журнал событий плюс materialized state
Для управляемого workflow храните неизменяемые события и текущую проекцию. PostgreSQL подходит лучше «памяти» внутри процесса.
create table agent_runs (
run_id uuid primary key,
tenant_id uuid not null,
workflow_name text not null,
workflow_version text not null,
state_schema_version integer not null,
status text not null check (status in
('running','suspended','completed','failed','cancelled')),
next_step text,
state jsonb not null,
state_revision bigint not null,
lease_owner text,
lease_expires_at timestamptz,
created_at timestamptz not null,
updated_at timestamptz not null
);
create table agent_events (
run_id uuid not null references agent_runs(run_id),
seq bigint not null,
event_type text not null,
step_id text,
attempt integer,
payload jsonb not null,
created_at timestamptz not null,
primary key (run_id, seq)
);
create table tool_operations (
operation_id text primary key,
run_id uuid not null,
step_id text not null,
tool_name text not null,
request_digest text not null,
status text not null check (status in
('prepared','started','succeeded','failed','unknown')),
external_ref text,
result jsonb,
updated_at timestamptz not null
);state содержит только необходимое для продолжения: validated inputs, результаты завершённых шагов, next_step, бюджеты, approvals и ссылки на артефакты. Не сериализуйте SDK client, coroutine, exception и секреты. Большие payload кладите в object store, в state оставляйте digest и URI.
Обновление проекции и append event выполняются одной DB-транзакцией с optimistic concurrency:
async def commit_step(run_id, expected_revision, event, new_state):
async with db.transaction():
changed = await db.execute("""
update agent_runs
set state = $1, next_step = $2,
state_revision = state_revision + 1, updated_at = now()
where run_id = $3 and state_revision = $4
""", new_state.data, new_state.next_step, run_id, expected_revision)
if changed != 1:
raise ConcurrentResume(run_id)
await append_event(run_id, event) # seq выделяется внутри той же транзакцииLease ограничивает одновременных исполнителей, revision защищает при истёкшем lease. Одного Redis lock недостаточно: владелец может остановиться, а старый процесс — ожить после выдачи нового lock.
Side effect не объединяется с вашей транзакцией
Между вызовом платежного API и записью checkpoint существует окно падения. Exactly-once через две независимые системы обычно недостижим; нужен идемпотентный протокол.
async def refund_step(run, order):
op_id = f"refund:{order.id}" # стабилен при retry/resume
op = await operations.prepare_once(
op_id, run.id, "refund", digest(order.refund_request)
)
if op.status == "succeeded":
return op.result
if op.status in {"started", "unknown"}:
remote = await payments.lookup_by_idempotency_key(op_id)
if remote.applied:
await operations.mark_succeeded(op_id, remote.ref, remote.result)
return remote.result
if not remote.authoritative:
raise SuspendForReconciliation(op_id)
await operations.mark_started(op_id)
try:
result = await payments.refund(order, idempotency_key=op_id)
except AmbiguousTransportError:
await operations.mark_unknown(op_id)
raise
await operations.mark_succeeded(op_id, result.id, result.data)
return result.dataTimeout после отправки — unknown, не failed. Resume сначала делает reconciliation по idempotency key или external reference. Если API не поддерживает идемпотентность и authoritative lookup, опасный шаг нельзя автоматически повторять.
Compaction: сокращать контекст, не состояние
Храните raw conversation/tool events отдельно. Compaction создаёт версионированный артефакт:
{
"summary_id": "sum_01J…",
"run_id": "…",
"covers_event_seq": [1, 84],
"source_digest": "sha256:…",
"compactor": {
"prompt_version": "compact-v4",
"model": "provider/model",
"schema_version": 2
},
"facts": [
{"text": "клиент просит возврат", "source_seq": [12, 19]}
],
"open_questions": ["подтверждена ли доставка"],
"decisions": [
{"code": "WAIT_FOR_TRACKING", "source_seq": [51, 58]}
],
"created_at": "…"
}Context builder берёт system/tool contracts, последние raw events, summary старой части и нужные retrieved documents. Неподтверждённые tool calls, approvals, operation IDs, денежные суммы и условия безопасности лучше переносить структурно, а не только прозой summary.
Compaction запускается по измеренному token budget конкретной модели. Оставьте резерв под tool schemas, retrieved context и output; границы выводятся из use case и tokenizer, а не из универсального процента. Проверяйте, что summary покрывает непрерывный диапазон событий и его source_digest совпадает. При изменении compactor-а старый summary можно перестроить из raw events.
RAG: индекс — кэш с provenance
Результат retrieval должен позволять проверить происхождение и свежесть:
{
"retrieval_snapshot_id": "rs_01J…",
"query_digest": "sha256:…",
"index_version": "kb-2026-08-17-3",
"acl_principal": "tenant:…/user:…",
"hits": [{
"document_id": "policy-42",
"source_uri": "s3://kb/policy-42.pdf",
"source_version": "etag:9f…",
"chunk_id": "policy-42:p7:c2",
"content_digest": "sha256:…",
"retrieved_at": "…",
"score": 0.81
}]
}TTL ограничивает допустимый возраст, но не заменяет invalidation. При изменении документа, ACL или tenant membership индексные записи помечаются stale/удаляются событием. Retrieval обязан фильтровать ACL до передачи текста модели. Для воспроизведения run сохраняйте snapshot metadata и immutable source version; «снова выполнить поиск» может вернуть другой набор.
Long-term memory: утверждения, а не transcript dump
Запись памяти должна быть узкой и проверяемой:
create table memory_items (
memory_id uuid primary key,
tenant_id uuid not null,
subject_id uuid not null,
namespace text not null,
key text not null,
value jsonb not null,
schema_version integer not null,
provenance jsonb not null,
confidence real,
valid_from timestamptz not null,
expires_at timestamptz,
invalidated_at timestamptz,
invalidation_reason text,
supersedes uuid references memory_items(memory_id),
consent_basis text,
created_at timestamptz not null
);provenance содержит тип источника, source ID/version, run_id, event sequence и способ валидации. Предпочтение «отвечать по-русски», явно заданное пользователем, может жить долго. «Пользователь рассержен», выведенное моделью из одного сообщения, не должно становиться долговременным фактом. Чувствительные выводы либо запрещаются политикой, либо требуют явного подтверждения.
Чтение памяти фильтрует tenant/subject/namespace, expires_at, invalidation и совместимую schema version. Новая запись не перезаписывает старую: она supersedes её, чтобы сохранить provenance и аудит. Удаление по запросу пользователя должно охватывать primary store, vector projection и кэши; tombstone предотвращает повторную индексацию из старого event stream.
Версии и миграции
Checkpoint привязан сразу к трём версиям: workflow code, state schema и contracts внешних инструментов. Деплой не должен слепо загружать старый JSON в новый код.
MIGRATIONS = {
(1, 2): migrate_v1_to_v2,
(2, 3): migrate_v2_to_v3,
}
def load_state(row, target=3):
data, version = row.state, row.state_schema_version
while version < target:
fn = MIGRATIONS.get((version, version + 1))
if fn is None:
raise UnsupportedCheckpoint(version)
data = validate(version + 1, fn(data))
version += 1
return validate(target, data)Миграции чистые, детерминированные и тестируются на сохранённых fixtures. Если изменена семантика незавершённого шага, безопаснее оставить старый worker для drain либо перевести run в ручную очередь. In-place update без сохранения исходной версии лишает возможности отката; храните событие миграции и digest до/после.
Resume verification
Перед продолжением worker не просто десериализует state:
- атомарно получает lease и проверяет revision;
- валидирует tenant, workflow/state schema и доступность миграции;
- сверяет digest внешних артефактов и версии tool contracts;
- проверяет непрерывность event sequence и соответствие materialized state последнему committed event;
- для
started/unknownoperations делает reconciliation; - повторно проверяет approvals: не истекли ли они и относятся ли к тому же request digest;
- пересчитывает оставшиеся step/token/cost/time budgets;
- только затем переводит run из
suspendedвrunning.
Resume не обязан продолжаться автоматически. Несовпавший digest, исчезнувший source document, неизвестный внешний эффект или несовместимый workflow — причины остановить run с диагностируемым статусом.
Failure modes
- История чата — checkpoint. В ней нет точного next step, operation status и budgets; модель вынуждена угадывать.
- Summary заменил raw events. Потерянную сумму или отрицание уже нельзя восстановить. Raw events сохраняются согласно retention, summary пересоздаётся.
- Vector DB считается памятью. Similarity не задаёт истинность, свежесть и права. Metadata filters и canonical store обязательны.
- TTL без invalidation. Отозванное разрешение остаётся активно до истечения срока. Изменения источника должны выпускать invalidation event.
- Два resume одновременно. Lease без revision допускает старого владельца. Нужны fencing/optimistic revision.
- Retry неизвестного side effect. Двойное списание. Сначала reconciliation, затем manual suspension при неопределённости.
- Новый код читает старый state «как получится». Silent semantic corruption опаснее явной ошибки. Только schema validation и миграция.
Практический сценарий: процесс умер после refund
Платёжный API применил возврат и отправил ответ, но worker умер до mark_succeeded и checkpoint. Новый worker видит tool_operations.status=started, вызывает lookup_by_idempotency_key, получает внешний refund ID, фиксирует succeeded, затем атомарно commit-ит step event и state. Модель не вызывается повторно для решения уже принятого шага. Если lookup недоступен или неавторитетен, run становится suspended; оператор сверяет платёж и возобновляет его с тем же operation_id.
Правило главы
RAG даёт внешние свидетельства, compaction экономит окно, checkpoint обеспечивает точное продолжение, long-term memory хранит подтверждённые межсессионные утверждения. У них разные схемы, сроки жизни, правила доступа и восстановления; объединение в один «memory store» стирает необходимые гарантии.
Чеклист главы
- Durable state хранит next step, revisions, budgets, approvals и operation IDs, а не Python-объекты.
- Event append и обновление state атомарны; конкурентный resume отсекается revision/fencing.
- Неоднозначный внешний результат имеет статус
unknownи проходит reconciliation. - Compaction ссылается на непрерывный диапазон raw events и проверяемый source digest.
- RAG hits содержат source/version/chunk provenance и проходят ACL-фильтрацию.
- Memory items имеют scope, provenance, TTL, invalidation, schema version и consent basis.
- Миграции детерминированы и протестированы на реальных старых checkpoints.
- Resume verification умеет отказать безопасно, а не только «продолжить».
Практикум
Задача 33.1 — crash-safe workflow. ⭑⭑⭑
Реализуйте workflow lookup order → decide → approve → refund → verify поверх PostgreSQL. Вставьте управляемые падения до вызова refund, после отправки запроса, после ответа и после checkpoint. Добавьте compaction длинного диалога и одну подтверждённую long-term preference.
Критерии приёмки:
- После каждого crash run возобновляется с корректного шага; refund применяется не более одного раза.
- Два worker-а, одновременно взявшие run, не могут оба commit-ить один revision.
- При неавторитетном lookup статус остаётся
unknown, автоматического retry нет. - Из summary удаление детали не меняет durable state и operation protocol.
- Изменение source document/ACL инвалидирует retrieval projection.
- Fixture checkpoint старой версии проходит миграцию; несовместимая версия явно блокируется.
- Истёкшая approval не используется после resume.
Первоисточники
- PostgreSQL: Explicit Locking
- PostgreSQL: Transaction Isolation
- PostgreSQL: JSON Types
- Python:
asynciosynchronization primitives - RFC 9110: Idempotent Methods