◆ ForgeVis
Skip to content

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'а:

#НодаОкноДействие
0S1T..T+12hPOST /export на S1
1S2T+12h..T+14hPOST /export на S2
2S1T+14h..endPOST /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 больше чем нод.

Известные ограничения текущей версии ​

  1. Аутентификация выключена. Платформа управления вызывает рекордеров и они вызывают друг друга без Authorization заголовка. Сценарий рассчитан на закрытую внутреннюю сеть. Включение auth — следующая итерация.
  2. Каждый рекордер должен присутствовать в таблице servers с корректным name (он же node_name в конфиге рекордера). Без этой записи sub-экспорт на эту ноду отправить нельзя.
  3. Маленькие окна, целиком лежащие на одной ноде, всё равно проходят через фазу assembling. Это лишний concat-pass для одного входного файла — миллисекунды на 4-часовом окне, но в логах это видно.
  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 → 502Assembler-нода прямо сейчас недоступна. Транзиент; через минуту попробуйте снова.
runs[].attempts > 0Был как минимум один рестарт рекордера во время этой задачи. Не ошибка, инфо.

Совместимость с одиночным режимом ​

Прямые вызовы рекордера POST /api/archive/{cam}/export (см. Management API) продолжают работать ровно как раньше — это single-node экспорт по локальной файловой системе. Multi-node поведение появляется только когда экспорт инициируется через сервис управления.

Это значит, что если вы перешли с single-node на multi-node развёртывание, никаких изменений в UI не требуется: оба обращения дают один и тот же формат job_id/state/download_url. Просто адрес теперь — сервис управления, а не конкретный рекордер.

Proprietary software.