Миграция схемы checkpoint агента без потери задач
Миграция схемы checkpoint агента без потери незавершенных задач: версии записей, конвертеры, аренды, идемпотентность и процессные тесты.

Checkpoint нельзя считать техническим мусором, который можно безнаказанно переименовать вместе с полем в коде. Если агент сохраняет состояние между попытками, перезапусками воркеров и релизами, эта запись определяет, продолжит ли он уже начатую работу или выполнит ее повторно. Ошибка в миграции здесь превращается не в красивый stack trace, а в повторный платеж, два одинаковых письма, потерянный документ или зависшую задачу, которую никто не может открыть.
Я видел типичный сценарий слишком много раз: команда меняет структуру state, выкатывает сервис, а через несколько дней ретрай поднимает запись, созданную старым кодом. Новый десериализатор либо падает, либо подставляет значение по умолчанию. Второй вариант опаснее. Агент продолжает работу с неверной позицией и оставляет после себя след, который трудно связать с одной строкой миграции.
Надежная миграция не начинается с SQL-скрипта. Она начинается с признания простого факта: checkpoint является долговечным контрактом выполнения. У него должна быть версия, однозначный путь преобразования, правила конкурентного захвата и тесты, в которых одна задача переживает несколько версий программы.
Checkpoint хранит не данные, а право продолжить работу
Checkpoint должен содержать ровно те факты, которые позволяют безопасно решить, какой шаг агент имеет право выполнить следующим. Это отличается от снимка всей оперативной памяти процесса и от журнала отладки.
Представим агента, который обрабатывает заявку: получает вложения, извлекает текст, вызывает модель, создает запись во внешней системе и отправляет уведомление. После каждого шага он сохраняет состояние. Плохой checkpoint хранит только step: 3. Хороший хранит идентификатор заявки, идентификаторы уже созданных внешних объектов, входные версии, попытку, время следующего запуска и идемпотентный ключ для вызова, который еще не подтвержден.
Разница проявляется после аварии. Если процесс упал после создания внешней записи, но до сохранения ее идентификатора, один номер шага не объяснит, была ли операция выполнена. Агент либо создаст дубль, либо пропустит работу. Такая неопределенность не лечится дополнительным полем status после того, как данные уже потеряны.
Разделяйте четыре вида информации:
- факты, нужные для следующего шага;
- результат внешнего действия, который нельзя вычислить повторно;
- данные для идемпотентности и сверки;
- диагностические сведения, которые помогают расследованию, но не меняют решение агента.
Логи и трассировка не заменяют checkpoint. Лог может быть неполным, храниться меньше нужного срока или не иметь транзакционной связи с записью задачи. Checkpoint также не равен очереди: очередь говорит, что работу надо попытаться выполнить, а checkpoint говорит, в каком именно месте ее надо возобновить.
Это различие часто размывают, когда в JSON кладут все подряд. Через пару релизов такой документ становится смесью временных кешей, отладочных флагов и бизнес-состояния. Мигрировать его тяжело не потому, что JSON сложен, а потому, что никто не знает смысл каждого поля. Перед изменением схемы составьте таблицу: поле, владелец, источник истины, нужно ли оно после перезапуска, можно ли восстановить его из другого места. Поля без ответа не должны участвовать в принятии решений.
Версия должна лежать в каждой записи
Версия схемы должна быть частью самого checkpoint, а не константой, выведенной из версии контейнера, ветки Git или даты создания. В один момент времени в хранилище почти всегда живут записи от нескольких релизов: задача могла ждать ретрая, быть заблокирована лимитом, попасть в карантин или остаться после ручной остановки.
Минимальный конверт выглядит так:
{
"task_id": "job_01J9R8...",
"schema_version": 2,
"revision": 17,
"status": "running",
"lease_until": "2026-07-23T10:17:00Z",
"state": {
"phase": "classify",
"source_document_id": "doc_481",
"classification_request_id": "req_9aa"
}
}
schema_version отвечает только за формат и семантику state. revision отвечает за конкурентное обновление одной записи. status описывает жизненный цикл задания. Не смешивайте эти числа. Когда команда использует одно поле version для всех трех задач, она неизбежно пишет условие, которое невозможно прочитать: «если версия меньше 12, значит ли это старый JSON, старое состояние или чужое обновление?»
Начинайте нумерацию с 1. Отсутствие версии можно временно трактовать как 0, если такие данные уже существуют. Но этот случай должен быть отдельным конвертером, а не набором проверок вида state.foo ?? state.bar ?? "" по всему коду.
Храните метаданные конверта отдельными колонками, если база это позволяет. В PostgreSQL удобно вынести task_id, status, lease_until, revision и schema_version в типизированные поля, а предметное состояние оставить в jsonb. Тогда можно индексировать активные задачи, находить записи старой схемы и ограничивать выборку без разбора JSON в приложении.
Руководство PostgreSQL прямо описывает, что стандартный READ COMMITTED показывает только данные, зафиксированные к началу отдельного оператора. Два последовательных запроса в одной транзакции могут увидеть разные результаты. Поэтому схема «сначала прочитали готовую задачу, потом отдельно пометили ее занятой» допускает гонку между воркерами.
Конвертеры должны идти только вперед
Конвертер checkpoint не обязан уметь превращать новый формат обратно в старый. Он обязан детерминированно довести любую поддерживаемую старую запись до текущего внутреннего представления, не меняя исходник до успешного продолжения задачи.
Пусть в версии 1 агент хранил один объект delivery:
{
"schema_version": 1,
"state": {
"phase": "send",
"delivery": {
"recipient": "[email protected]",
"body": "Готово",
"sent": false
}
}
}
В версии 2 вы разделили намерение отправить сообщение и подтвержденный внешний результат. Это разумное изменение: message описывает то, что агент хочет сделать, а provider_message_id доказывает, что провайдер уже принял запрос. Конвертер должен явно зафиксировать, что он не знает о старом внешнем результате больше, чем знает исходная запись.
function upgradeToV2(v1: CheckpointV1): CheckpointV2 {
if (v1.state.phase !== "send") {
return {
schema_version: 2,
state: { ...v1.state, delivery_attempt: null }
};
}
return {
schema_version: 2,
state: {
phase: "send",
message: {
recipient: v1.state.delivery.recipient,
body: v1.state.delivery.body
},
provider_message_id: null,
delivery_attempt: v1.state.delivery.sent
? { outcome: "unknown", migrated_from: 1 }
: null
}
};
}
Обратите внимание на неприятную деталь: значение sent: true не обязано означать, что существует идентификатор сообщения у провайдера. Если старый код ставил флаг до сетевого вызова или сохранял его после вызова без отдельного подтверждения, конвертер не может честно придумать недостающий факт. Он должен перевести задачу в состояние, где исполнитель сделает сверку по идемпотентному ключу, запросит внешний сервис или отправит запись на разбор оператору.
Именно здесь команды часто делают опасный выбор: считают, что миграция должна «починить» смысл старых данных. Нет. Конвертер меняет представление. Он не получает право выдумывать историю выполнения.
Оформите цепочку преобразований как отдельный модуль:
type AnyCheckpoint = CheckpointV0 | CheckpointV1 | CheckpointV2 | CheckpointV3;
function normalize(raw: AnyCheckpoint): CheckpointV3 {
let current = raw;
while (current.schema_version < 3) {
switch (current.schema_version) {
case 0: current = upgradeV0toV1(current); break;
case 1: current = upgradeV1toV2(current); break;
case 2: current = upgradeV2toV3(current); break;
default: throw new UnsupportedCheckpointVersion(current);
}
}
if (current.schema_version !== 3) {
throw new UnsupportedCheckpointVersion(current);
}
return validateV3(current);
}
Не пишите один гигантский migrateToLatest, который распознает десятки комбинаций полей. Он быстро становится вторым, неформальным форматом данных. Маленькие переходы V1 -> V2 и V2 -> V3 проще тестировать, проще удалить и легче расследовать по журналу.
Шаблон Parallel Change, описанный Данило Сато, разбивает несовместимое изменение на расширение, переходный период и удаление старого пути. Для checkpoint это означает: сначала научить читатель понимать оба формата, затем создавать новый формат, а уже потом убирать старый.
Чтение старого формата и запись нового решают разные задачи
Совместимость чтения отвечает за судьбу уже начатых задач. Совместимость записи отвечает за то, что создают новые воркеры. Их нельзя выкатывать одним неразличимым изменением.
Рабочая последовательность для перехода с версии 1 на версию 2 выглядит так:
- Добавьте версию 2 в читатель и конвертер, но продолжайте создавать версию 1.
- Выпустите этот код и убедитесь, что он обрабатывает реальные старые записи без ошибок нормализации.
- Переключите писатель на версию 2, сохранив чтение версии 1.
- Дождитесь, пока завершатся либо будут обработаны все активные записи версии 1.
- Удалите создание и чтение версии 1 только после проверки хранилища и резервного плана для старых экспортов.
Первый этап кажется лишним, пока не понадобится откат. Если новый писатель уже создал версию 2, а предыдущий релиз не умеет ее читать, rollback превращается в остановку всего пула задач или в ручную правку строк. Иногда это допустимо, но решение надо принять до выката, а не после алерта.
Не путайте обратную совместимость с двусторонней записью. Популярная рекомендация «пишите оба формата, чтобы не рисковать» часто ухудшает ситуацию. Два представления одного состояния расходятся при частичных сбоях: новый код обновил provider_message_id, а старое поле sent осталось прежним. Затем старый воркер читает устаревший вариант и делает повторный вызов.
Двойная запись оправдана, только если у вас есть четко определенный главный источник, правило сверки и короткое окно жизни этого решения. В большинстве агентов лучше читать старое, нормализовать в памяти и записывать только текущую схему при следующем успешном сохранении. Это дает постепенную материализацию без массового переписывания таблицы.
Захват задачи и сохранение состояния должны образовывать один протокол
Миграция формата не спасет задачу, если два воркера могут одновременно считать ее своей. Нужен явный протокол аренды или блокировки, а также защита от устаревшего писателя, который завершил работу позже нового владельца.
Для очереди задач в PostgreSQL часто достаточно аренды и оптимистичной ревизии. Сначала воркер атомарно захватывает запись, затем выполняет шаг, затем сохраняет новое состояние только при совпадении revision и владельца аренды.
WITH candidate AS (
SELECT task_id
FROM agent_checkpoints
WHERE status IN ('ready', 'retry')
AND run_after <= now()
AND (lease_until IS NULL OR lease_until < now())
ORDER BY run_after, created_at
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE agent_checkpoints c
SET lease_owner = $1,
lease_until = now() + interval '60 seconds',
status = 'running',
revision = revision + 1
FROM candidate
WHERE c.task_id = candidate.task_id
RETURNING c.task_id, c.schema_version, c.revision, c.state;
SKIP LOCKED подходит для распределения независимых задач между воркерами, но он не решает бизнес-согласованность сам по себе. После захвата храните возвращенную revision. При сохранении делайте условное обновление:
UPDATE agent_checkpoints
SET schema_version = $2,
state = $3::jsonb,
status = $4,
run_after = $5,
lease_owner = NULL,
lease_until = NULL,
revision = revision + 1,
updated_at = now()
WHERE task_id = $1
AND revision = $6
AND lease_owner = $7
RETURNING revision;
Если запрос не вернул строку, воркер потерял право сохранять результат. Он не должен пробовать обновление второй раз с новой ревизией. Иначе запоздалый процесс перезапишет checkpoint, который уже продвинул другой воркер.
PostgreSQL предупреждает, что транзакции уровня SERIALIZABLE могут завершиться ошибкой сериализации, и приложение должно повторить всю транзакцию. Это относится к механике захвата и обновления: повторять можно короткую транзакцию работы с записью, но нельзя бездумно повторять уже выполненный внешний вызов.
Поэтому отделите две операции. Транзакция захватывает право выполнить шаг. Внешний вызов использует идемпотентный ключ. Вторая транзакция фиксирует подтвержденный результат. Между ними процесс может умереть, и именно это место определяет качество вашего checkpoint.
Идемпотентность нельзя восстановить из красивого JSON
Самая неприятная точка находится между «отправили запрос наружу» и «сохранили ответ». Если агент погиб в этом окне, после рестарта он видит только намерение выполнить действие. Внешняя система могла принять запрос, отвергнуть его или принять, но не успеть вернуть ответ.
Для операций, которые поддерживают идемпотентность, сформируйте ключ до вызова и положите его в checkpoint до обращения к провайдеру:
{
"phase": "create_case",
"request": {
"customer_id": "cust_104",
"summary": "Проверка документа"
},
"idempotency": {
"key": "job_01J9R8:create_case:0",
"attempt": 0
},
"external_case_id": null
}
После рестарта агент повторяет запрос с тем же ключом. Внешняя сторона либо возвращает прежний результат, либо отвергает повтор как уже обработанный. Если API не поддерживает такой ключ, используйте собственный журнал намерений и сверку по устойчивому бизнес-идентификатору. Но не подменяйте это флагом done: true.
У каждой фазы должен быть один из трех ясных статусов: действие не начато, действие подтверждено, результат неизвестен. Статус «в процессе» без дополнительного контекста почти бесполезен после падения. Он говорит, что код когда-то вошел в функцию, но не говорит, произошло ли событие на внешней стороне.
Если внешняя операция необратима и не имеет идемпотентности, добавьте явный путь разбора. Например, задача со статусом needs_reconciliation не должна бесконечно ретраиться. Оператору нужны исходный запрос, время попытки, корреляционный идентификатор, ответ при наличии и описание того, какие поля нельзя вычислить. Это дешевле, чем потом объяснять владельцу данных, почему агент выполнил действие дважды.
Тесты должны переживать не только функцию, но и процесс
Юнит-тест конвертера необходим, но он ловит только самый удобный класс ошибок. Реальные поломки возникают на стыке старой записи, нового двоичного файла, аренды, сетевого вызова, остановки процесса и следующего релиза.
Соберите фикстуры из реальных исторических checkpoint. Перед тем как обезличить данные, сохраните их форму: отсутствующие поля, пустые массивы, старые названия фаз, незавершенные попытки, записи с истекшей арендой. Ручной пример, созданный после изменения схемы, почти всегда слишком аккуратный.
Проверяйте как минимум следующие свойства:
- любая поддерживаемая фикстура нормализуется до текущей версии;
- нормализация повторяема: второй запуск не меняет результат;
- конвертер не изменяет входной объект;
- неизвестная версия останавливает задачу предсказуемо;
- обязательные инварианты текущей схемы проверяются после преобразования.
Для цепочки релизов нужен процессный тест. Он не обязан поднимать полный production-стенд, но должен использовать настоящую базу, реальное сохранение checkpoint и отдельные процессы или изолированные экземпляры приложения. Сценарий выглядит так:
1. Релиз A создает задачу версии 1 и сохраняет checkpoint после внешнего намерения.
2. Тест завершает процесс A до фиксации внешнего результата.
3. Релиз B захватывает ту же строку, преобразует версию 1 в версию 2 и повторяет вызов с тем же idempotency key.
4. Релиз B сохраняет подтвержденный результат версии 2.
5. Релиз C читает эту задачу и завершает ее, не вызывая внешнюю сторону повторно.
Проверяйте не только schema_version = 3. Проверяйте, что внешняя заглушка получила один логический запрос, что итоговый идентификатор сохранен, а повторный запуск завершенной задачи ничего не делает. Если API-заглушка умеет хранить вызовы по идемпотентному ключу, этот тест хорошо ловит дубли.
Добавьте тест на устаревшего владельца аренды. Воркер A захватывает задачу и зависает. Аренда истекает, воркер B завершает задачу. Затем A «просыпается» и пытается записать старый checkpoint. Условный UPDATE должен вернуть ноль строк, а A должен зафиксировать потерю аренды без повторного исполнения эффекта.
Есть еще один тест, который команды пропускают: переход через несколько схем. Не тестируйте только V2 -> V3, если в хранилище могла остаться V1. Проверяйте V1 -> V2 -> V3 через настоящий публичный путь normalize. Удаление промежуточного конвертера без такой проверки ломает долгие ретраи именно тогда, когда они нужны больше всего.
Массовая миграция допустима только как отдельная операция
Фоновое переписывание всех checkpoint иногда нужно, например когда вы меняете типизированные колонки, хотите построить новый индекс или обязаны убрать чувствительное поле. Но оно не заменяет совместимого чтения.
Если джоб миграции обновляет миллионы строк, он конкурирует с воркерами за те же записи. Не запускайте один огромный UPDATE без ограничений. Он создаст долгую транзакцию, раздует журнал изменений, усложнит откат и может удерживать ресурсы, нужные живым задачам.
Обрабатывайте записи пакетами, выбирайте только конкретную старую версию и используйте тот же контроль ревизии, что у исполнителя. Если воркер изменил checkpoint между чтением и миграцией, фоновый процесс должен пропустить строку и вернуться к ней позже. Его цель не победить гонку, а не сломать актуальное состояние.
Подход expand, migrate, contract полезен и здесь: сначала код понимает новую форму, затем фон меняет старые строки, после подтверждения отсутствия старых данных вы удаляете переходный путь. В материале Evolutionary Database Design этот переходный период назван отдельной частью изменения, а не побочным эффектом DDL. Это правильная дисциплина для состояния агента, даже когда сама схема хранится в JSON.
Перед массовой миграцией договоритесь о метриках. Считайте активные checkpoint по schema_version, число ошибок нормализации, количество задач needs_reconciliation, возраст самой старой активной записи и число отказов условного сохранения. Без этого вы не знаете, завершился ли переход, или просто перестали смотреть на старые данные.
Неизвестную или испорченную запись надо остановить, а не угадывать
Самая плохая реакция на неизвестную схему выглядит так: поймать ошибку, создать пустой объект состояния и позволить агенту начать с начала. Это может быть приемлемо только для операции, которую документированно можно повторить без последствий. Для большинства бизнес-процессов это скрытая потеря контекста.
Разделите ошибки на три категории. Запись старой поддерживаемой версии проходит конвертер. Запись будущей или неизвестной версии попадает в карантин, потому что текущий код не понимает ее семантику. Поврежденная запись тоже попадает в карантин, но с отдельной причиной: JSON невалиден, нарушен инвариант, отсутствует обязательный идентификатор или структура не соответствует заявленной версии.
Сохраняйте исходный payload неизменным, причину отказа, версию кода читателя и идентификатор задачи. Не перезаписывайте проблемный checkpoint «исправленным» пустым объектом. Оператору может понадобиться сравнить его с резервной копией, журналом внешнего сервиса или результатом более нового релиза.
Если состояние сериализуется бинарным форматом, риск еще выше. Документация Python прямо предупреждает, что pickle небезопасен для недоверенных данных и не должен использоваться для входа, который мог быть подменен. Даже во внутреннем хранилище бинарный снимок без явной схемы усложняет аудит и долговременную совместимость.
Выбирайте формат, который можно валидировать независимо от кода исполнителя. JSON с JSON Schema, Protobuf с продуманными правилами эволюции или типизированные колонки подходят лучше, чем сериализация объекта рантайма. Формат сам по себе не спасает от плохой семантики, но он дает вам шанс увидеть несовместимость до того, как агент выполнит действие.
Удаление старого конвертера требует доказательства
Конвертеры раздражают разработчиков, потому что код переходов выглядит временным. Но «временный» не означает «его можно удалить после двух спринтов». Он нужен ровно столько, сколько в системе могут появиться и существовать старые записи.
Сначала прекратите создание старой версии. Затем дождитесь завершения активных задач этого формата, включая отложенные ретраи и карантин. После этого проверьте основное хранилище, реплики для восстановления, экспортные очереди и инструменты ручного повторного запуска. Если инженер может взять checkpoint из архива и попытаться выполнить его новым воркером, этот путь тоже входит в контракт.
Зафиксируйте в репозитории границу поддержки: например, текущая схема 4 читает версии 2, 3 и 4, а версия 1 больше не поддерживается после того, как все ее записи прошли через отдельную процедуру разбора. Это лучше, чем бесконечное if вокруг полей. У конвертера должна быть дата удаления, но удаляйте его по наблюдаемому отсутствию данных, а не по календарю.
Первое изменение, которое стоит сделать до следующего релиза, простое: добавьте schema_version в каждую новую запись и запретите читателю молча принимать неизвестный формат. Это не решит прошлые ошибки, но остановит привычку прятать несовместимость за значениями по умолчанию. Дальше построите цепочку конвертеров, проверите аренду и заставите тест пройти через реальный рестарт процесса. Тогда незавершенная задача будет переживать релизы как рабочий объект, а не как случайный остаток старого кода.
Часто задаваемые вопросы
Нужна ли версия схемы, если checkpoint хранится в JSON?
Нельзя считать checkpoint внутренней деталью, если агент сохраняет его между перезапусками. Как только запись переживает процесс, релиз или смену воркера, ее формат становится контрактом восстановления. Версия нужна даже при хранении JSON в одной колонке.
Можно ли использовать номер релиза агента вместо версии checkpoint?
Нет. Номер релиза описывает код, а schema_version описывает форму конкретной записи. Во время поэтапного выката один релиз может читать несколько форматов, поэтому связывать формат только с версией сервиса опасно.
Как долго нужно поддерживать старые форматы состояния?
Старые записи можно читать до тех пор, пока существует хотя бы один незавершенный процесс, который мог их создать, плюс запас на ручное восстановление и отложенные ретраи. Удалять преобразователь стоит только после измеримой проверки, что таких записей больше нет. Архив без рабочего читателя не помогает восстановить задачу в инциденте.
Защищает ли транзакция базы данных от повторного выполнения шага?
Отдельная транзакция защищает только часть операции. Если агент уже вызвал внешний API, а затем упал до сохранения checkpoint, повторный запуск может сделать тот же вызов еще раз. Нужны идемпотентные ключи на внешней стороне или журнал намерений с понятным правилом повтора.
Можно ли откатить релиз после миграции checkpoint?
Да, если старый код способен корректно загрузить запись после того, как новый код ее изменил. На практике это редко верно при разделении одного поля на несколько или при изменении семантики. Поэтому сначала вводят совместимое расширение, затем переводят записи, и только потом удаляют старый формат.
Нужно ли переписывать все старые checkpoint сразу?
Лучше сохранять исходную запись без изменения и выполнять преобразование в памяти при чтении. Это делает повторный запуск безопасным и оставляет возможность сравнить старый и новый смысл. Фоновая материализация допустима позже, когда новый формат доказал свою работоспособность.
Что делать с checkpoint неизвестной версии?
Неизвестная версия не должна молча трактоваться как пустое состояние или как ближайший известный формат. Пометьте задачу как требующую вмешательства, сохраните исходный payload и поднимите сигнал. Продолжение с догадкой обычно хуже остановки одной задачи.
Какие тесты нужны для миграции состояния агента?
Минимальный набор включает чтение старых фикстур, повторный запуск преобразователя, восстановление после аварии между захватом и сохранением, а также цепочку из нескольких релизов. Юнит-тест функции преобразования нужен, но он не ловит ошибки захвата записи, транзакций и воркеров. Проверяйте итоговый бизнес-эффект, а не только форму JSON.
Какие данные агент должен сохранять в checkpoint?
Да, если процесс должен пережить рестарт воркера. Состояние только в памяти годится для коротких операций, которые можно безопасно начать заново. Для многошаговых действий с внешними вызовами храните идентификатор задачи, позицию, результат шага и данные для идемпотентности.
Когда можно удалить конвертер старой схемы?
Это зависит от максимальной жизни задачи, политики ретраев, задержек очереди и того, как долго операторы могут разбирать инцидент. Установите явное окно совместимости и измеряйте возраст активных записей. Не удаляйте конвертер только потому, что прошло несколько релизов.