Appearance
Multi-node archive (мульти-нодный экспорт)
Скоро
Функция зависит от сервиса управления (platform/backend), который пока не опубликован — скоро будет. Одиночная запись, рестриминг и экспорт на standalone-рекордере работают уже сейчас и без него.
Когда камера в течение запрошенного окна экспорта перемещалась между несколькими нодами-рекордерами, её сегменты архива физически лежат на разных серверах. Multi-node archive — это слой управления, который собирает фрагменты с нужных нод, склеивает их в один итоговый MP4 и отдаёт пользователю.
Эта функциональность работает только при наличии общей БД и сервиса управления (platform/backend). Одиночные standalone-рекордеры экспорт обслуживают сами, тут ничего не меняется.
Когда фича включается
Multi-node archive активен при выполнении обоих условий:
- Все рекордеры подключены к одной общей PostgreSQL (см. Конфигурация — База данных).
- Рядом развёрнут сервис управления
platform/backend, и UI/SDK обращается за экспортом к нему (на адрес сервиса управления), а не напрямую к рекордерам.
В этом режиме экспорт инициирует сервис управления, а рекордеры выступают как исполнители: один отвечает за свой кусочек, ещё один делает финальную склейку. Конечный пользователь видит обычный job со своим job_id и стандартным циклом «submit → poll → download».
Архитектура потока
┌────────┐ POST /api/archive/{cam}/export?from&to
│ user │ ─────────────────────────────────────────► ┌────────────┐
└────────┘ │ platform/ │
│ backend │
└─────┬──────┘
│
┌────────────────────────────────────────────────────────┤
│ 1. SELECT recordings WHERE camera AND start_time IN window
│ 2. Группировка по хост-нодам в хронологические «runs»
│ 3. INSERT archive_jobs + N×archive_job_runs (state=pending)
│ 4. asyncio фоновая задача-оркестратор
│
▼ фаза 1: экспорт каждого фрагмента на своей ноде
┌─────────┐ ┌─────────┐ ┌─────────┐
│ node S1 │ POST /export?from&to (run #0) │ node S2 │ │ node S3 │
│ │ ──────────────────────────────► │ │ │ │
│ ffmpeg │ │ ffmpeg │ │ ffmpeg │
│ -c copy │ │ -c copy │ │ -c copy │
└────┬────┘ └────┬────┘ └────┬────┘
│ part_0.mp4 готов │ │
│ │ │
фаза 2: backend выбирает «assembler» — ноду с самой большой долей
│ │ │
┌────▼──────────────────────────────────────────────────────────────┐
│ assembler (например S2) │
│ POST /api/archive/{cam}/assemble │
│ body: [(addr_S1, job_0), (addr_S2, job_1), (addr_S3, job_2)] │
│ ──► качает 3 part-файла параллельно │
│ ──► ffmpeg concat parts.txt → final.mp4 │
│ ──► (если timelapse) ffmpeg setpts=PTS/N → final-tl.mp4 │
└──────────────────────────────────────────────────────────────────┬─┘
│
┌──────────────────────────────────────────────────────────────┘
│
▼ фаза 3: пользователь забирает результат
┌────────┐ GET /api/archive/jobs/{id}/download
│ user │ ◄────────────────────────────────── backend стримит из assembler'а
└────────┘Что такое «run»
Сервис управления собирает все сегменты камеры из таблицы recordings в выбранном окне и разбивает их на непрерывные блоки одного хоста — «runs». Каждый run — это один POST-запрос экспорта на одну ноду.
Пример. Камера 12 часов жила на S1, потом мигрировала на S2 на 2 часа, потом обратно на S1. Пользователь запросил экспорт диапазона охватывающего всю миграцию — получаются три run'а:
| # | Нода | Окно | Действие |
|---|---|---|---|
| 0 | S1 | T..T+12h | POST /export на S1 |
| 1 | S2 | T+12h..T+14h | POST /export на S2 |
| 2 | S1 | T+14h..end | POST /export на S1 ещё раз — отдельным заданием |
Важно: даже если две части лежат на одной и той же ноде, между ними есть «дырка» от другого хоста. Поэтому они идут двумя отдельными запросами с двумя разными окнами — иначе финальный MP4 содержал бы внутренний разрыв.
Кто такой «assembler»
После того как все три (или сколько их вышло) part-файла готовы, нужна одна нода, которая скачает их к себе и склеит. Сервис управления выбирает её так:
«нода, которой принадлежит максимальная по длительности доля окна»
Свойства:
- Если миграции в окне нет вовсе (все сегменты на одной ноде), assembler совпадает с этой нодой — он скачивает 1 файл со 127.0.0.1, склеивать по сути нечего. Минимальные накладные.
- При миграциях выбор минимизирует объём трафика между нодами (assembler тянет только чужие куски).
- Сервис управления никогда не выступает assembler'ом — он только оркеструет и проксирует финальный download'.
Endpoints сервиса управления
Эти эндпойнты обслуживает platform/backend. Они выглядят как обычные recorder-эндпойнты, но фактически orchestrate'ят несколько нод.
Постановка экспорта
POST /api/archive/{camera_id}/export
Content-Type: application/json
{
"from_ts": 1780054469,
"to_ts": 1780058069
}Ответ 202 Accepted:
json
{
"job_id": "9e23...",
"camera_id": "camera_001",
"kind": "export",
"from_ts": 1780054469,
"to_ts": 1780058069,
"state": "pending",
"created_at": 1780100000,
"runs": [
{"ordinal": 0, "node_name": "S1", "from_ts": 1780054469, "to_ts": 1780056269, "recorder_state": "pending"},
{"ordinal": 1, "node_name": "S2", "from_ts": 1780056269, "to_ts": 1780058069, "recorder_state": "pending"}
]
}runs показывает план сразу — пользователь видит, какие ноды будут задействованы и какое временное окно достанется каждой.
Постановка time-lapse
POST /api/archive/{camera_id}/timelapse
Content-Type: application/json
{
"from_ts": 1780054469,
"to_ts": 1780100000,
"speed": 30
}speed — целое в диапазоне [2, 120]. Sub-экспорты на нодах остаются обычными stream-copy экспортами; ускорение setpts=PTS/speed применяет assembler уже к склеенному файлу.
Статус
GET /api/archive/jobs/{job_id}Возможные state:
| Значение | Что значит |
|---|---|
pending | задача только создана, рекордерам ещё ничего не отправлено |
exporting | каждый run работает на своей ноде; см. runs[].recorder_state |
assembling | все runs завершены, assembler сейчас склеивает |
ready | финальный MP4 готов, можно скачивать |
failed | задача провалилась; см. reason |
Поле runs[].recorder_state — pending, running, done, failed — мнимое отражение recorder-side состояния. Поле attempts показывает сколько раз сервис управления переотправлял sub-экспорт на эту ноду (см. раздел «Устойчивость к перезапускам» ниже).
После переходов в failed поле reason содержит человекочитаемое объяснение (например run 2 on S3 unrecoverable after 3 attempts).
Скачивание
GET /api/archive/jobs/{job_id}/downloadСтримит готовый MP4 через сервис управления. Сервис не буферизует файл на свой диск — байты летят напрямую с assembler-ноды клиенту через прокси-стрим. Эта схема выдерживает большие файлы даже на маленьком VPS для backend'а.
Возможные ответы:
200+video/mp4— стрим файла.409 Conflict— задача ещё неready(всё ещёpending/exporting/assembling).410 Gone— задачаready, но рекордер уже почистил mp4 (по своему TTL через час). Закажите экспорт заново.502 Bad Gateway— assembler-рекордер недоступен прямо сейчас, хотя данные у него есть. Транзиентная ошибка; попробуйте через минуту.
Список заданий
GET /api/archive/jobs?camera_id=camera_001&limit=50&offset=0История по конкретной камере (или всем). Полезно для UI с историей запросов на экспорт. Возвращает компактный ArchiveJobListItemDto (без массива runs).
Отмена / уборка
DELETE /api/archive/jobs/{job_id}Удаляет запись из БД, плюс best-effort удаление промежуточных файлов со всех вовлечённых рекордеров (sub-экспорты + assembler-файл). Если какой-то рекордер недоступен в этот момент — его TTL-уборщик подметёт файл через час.
Устойчивость
Multi-node экспорт — длительная операция (минуты на больших окнах). За это время любая часть системы может перезапуститься.
Сценарий 1: рекордер падает в середине sub-экспорта
В archive_jobs сохранён план целиком, включая исходное окно каждого run'а. При следующем поллинге сервис управления получает 404 от поднявшегося рекордера → переотправляет тот же sub-экспорт с тем же окном → запоминает новый recorder_job_id. Поле attempts инкрементится. Лимит — 3 попытки на run, после чего job уходит в failed.
Файлы записи на упавшей ноде никуда не делись (они на постоянном диске, не во временной памяти), поэтому переотправка с тем же окном даёт тот же результат.
Сценарий 2: assembler падает после старта склейки
Сервис управления видит 404 от assembler'а → выбирает другую ноду (исключая упавшую) → отправляет ту же /assemble задачу с тем же списком частей. Sub-экспорты от рекордеров ещё лежат у них (по TTL), и новый assembler их подхватит. Лимит — 2 попытки assembler'а на job.
Сценарий 3: сервис управления (backend) рестартует
Это самый «тонкий» случай и его мы прорабатываем особо аккуратно.
Что сохраняется: вся таблица archive_jobs + archive_job_runs в PG. Это означает, что на момент рестарта backend знает план каждой задачи и какой recorder_job_id он сохранил у каждого рекордера.
На старте backend сканирует все задачи в нетерминальных состояниях и для каждой поднимает фоновую задачу-оркестратор. Оркестратор идемпотентен:
- Если у run'а уже есть
recorder_job_id— повторно не отправляем, просто опрашиваем. - Если рекордер вернул
404(он тоже рестартовал за время простоя backend'а) — действует обычный retry-сценарий run'а. - Если задача была в фазе
assemblingи assembler-нода ещё жива — оркестратор сразу подключается к её polling'у. - Если assembler пропал — выбирается новая нода (с учётом исключения предыдущей).
Итог: пользователь не теряет задачу при рестарте backend'а — после короткой паузы polling продолжается. В худшем случае сервис управления вообще не успел дойти до рестарта корректно, и часть задач упадёт в failed — пользователь увидит причину и переотправит.
Сценарий 4: одной из нужных нод нет в реестре servers
Если в archive_job_runs.node_name оказалось имя, которого нет в таблице servers (нода декомиссионирована, опечатка), job сразу падает в failed с reason node 'X' not in servers registry. Чините реестр или удаляйте старые записи recordings.
Производительность и параллелизм
- Сервис управления не делает тяжёлой работы: ни ffmpeg, ни буфера файлов. Только: запросы к PG, POST/GET к рекордерам, проксирование стрима. Маленькой VPS под backend обычно достаточно.
- Тяжёлая работа (декод/энкод, ffmpeg concat) распределяется по рекордерам: чем больше нод, тем больше одновременных экспортов параллельно может крутиться.
- Пер-нодные лимиты остаются как у одиночного экспорта: 4 параллельных export-задания и 2 timelapse-задания на нодду — лишние ждут в очереди с
recorder_state: pending. - Параллелизм скачивания частей у assembler'а ограничен 3 параллельными потоками — это защищает CPU/сеть assembler-ноды от перегрузки на «болтливых» миграциях, где runs больше чем нод.
Известные ограничения текущей версии
- Аутентификация выключена. Платформа управления вызывает рекордеров и они вызывают друг друга без
Authorizationзаголовка. Сценарий рассчитан на закрытую внутреннюю сеть. Включение auth — следующая итерация. - Каждый рекордер должен присутствовать в таблице
serversс корректнымname(он жеnode_nameв конфиге рекордера). Без этой записи sub-экспорт на эту ноду отправить нельзя. - Маленькие окна, целиком лежащие на одной ноде, всё равно проходят через фазу assembling. Это лишний concat-pass для одного входного файла — миллисекунды на 4-часовом окне, но в логах это видно.
- Резюмирование при рестарте backend'а возвращает задачу к polling-циклу, но не восстанавливает прогресс ffmpeg на рекордере. Если sub-экспорт был уже почти готов, мы либо подхватим его (если рекордер не упал), либо переотправим с нуля (если упал).
Диагностика
| Симптом | Что смотреть |
|---|---|
state=failed, reason содержит not in servers registry | Запись servers.name отсутствует для одной из нод-владелцев сегментов. Сравните archive_job_runs.node_name с servers.name. |
state=failed, reason содержит unrecoverable after N attempts | Рекордер либо лежит, либо его настройки фильтрации не возвращают сегмент. Проверьте GET /api/archive/{cam}/segments?from=&to= на этом рекордере. |
state=failed, reason содержит assembler ... after N tries | Все попытки assembler'а провалились. Возможно, в окне всего одна нода и она проблемная. Проверьте логи assembler'а — там будет ffmpeg stderr. |
download → 410 | Прошло больше часа — рекордер удалил mp4 по TTL. Запустите экспорт заново. |
download → 502 | Assembler-нода прямо сейчас недоступна. Транзиент; через минуту попробуйте снова. |
runs[].attempts > 0 | Был как минимум один рестарт рекордера во время этой задачи. Не ошибка, инфо. |
Совместимость с одиночным режимом
Прямые вызовы рекордера POST /api/archive/{cam}/export (см. Management API) продолжают работать ровно как раньше — это single-node экспорт по локальной файловой системе. Multi-node поведение появляется только когда экспорт инициируется через сервис управления.
Это значит, что если вы перешли с single-node на multi-node развёртывание, никаких изменений в UI не требуется: оба обращения дают один и тот же формат job_id/state/download_url. Просто адрес теперь — сервис управления, а не конкретный рекордер.