WebSocket-стримы прогресса задач

Асинхронные задачи публичного API (смета и живая актуализация цен/остатков) считаются в фоне построчно. Прогресс можно забирать двумя равнозначными способами:

  • REST-опрос по курсоруGET .../jobs/{id}?after=<seq> возвращает строки, появившиеся после курсора (см. cookbook, сценарии G и H);
  • WebSocket-стрим — push тех же событий без поллинга.

Данные идентичны — это один и тот же cursor-прогресс. WebSocket выбирают ради realtime-обновлений (прогресс-бар, построчное появление в UI), REST-опрос — ради простоты и устойчивости к обрывам.

Endpoint’ы

Метод + путьЗадачаScope
GET /v1/jobs/estimates/{id}/eventsAsync-смета (/v1/estimates/jobs)search:read
GET /v1/jobs/offers-refresh/{id}/eventsAsync-refresh (/v1/offers/refresh-jobs)offers:refresh

{id}job_id, полученный при создании задачи (202 {job_id}). Задача owner-scoped: чужой или несуществующий job_id неотличимо возвращает 404 до апгрейда соединения (foreign-owner маскируется под absent — тот же приём, что в REST-эндпоинтах задач).

Авторизация — тем же заголовком, что и REST: Authorization: Bearer trk_…. Scope проверяется точно так же, как на REST-маршрутах того же bounded context. Метрики потребления (usageRec) на WS-соединение не начисляются: соединение может жить весь lifecycle задачи (до 3600 с на уровне nginx), поэтому оно не укладывается в модель «один запрос — один ответ».

Протокол

  1. Апгрейд. Клиент коннектится, сервер апгрейдит соединение на WebSocket.
  2. Курсор (опционально). Первым текстовым сообщением клиент может прислать {"after":N} — реплей начнётся со строки после N.
    • дедлайн на это сообщение — 5 секунд;
    • нет сообщения или таймаут → трактуется как after=0 (реплей с начала);
    • невалидный JSON → сервер закрывает соединение кодом 1003.
  3. Подписка → реплей. Сервер сначала подписывается на брокер (чтобы не потерять события, опубликованные между подпиской и реплеем), затем реплеит пропущенные строки LinesSince(after) как {"type":"line","seq":N,"line":<payload>}.
  4. Уже терминальная задача. Если задача завершилась (и реплей отдан) — сервер синтезирует {"type":"final","job":{"status":…}} и закрывает кодом 1000.
  5. Live-tail. Иначе сервер тейлит канал брокера:
    • line — новая готовая строка (дедуп по seq);
    • progress — обновление прогресса (total, done, percent);
    • final — терминальный статус; после отправки соединение закрывается 1000;
    • lagged — клиент отстал и был вытеснен из буфера брокера: сервер шлёт {"type":"lagged"} и повторно реплеит LinesSince(lastSeq) — без переподключения со стороны клиента.
  6. Keepalive. Ping раз в 30 с; если pong не пришёл за 60 с — соединение закрывается.

Типы сообщений (сервер → клиент)

typePayloadСмысл
line{"seq":N,"line":<verbatim JSON строки>}Готовая строка результата; seq — курсор
progress{"total":…,"done":…,"percent":…}Прогресс задачи
final{"job":{"status":…}}Терминальный статус, дальше close 1000
lagged{}Клиент отстал; далее повторный реплей

line.line — это тот же verbatim-payload строки, что уходит и в REST-опрос, и в брокер; форма строки зависит от вида задачи (смета vs refresh).

Коды закрытия

КодПричина
1000Нормальное завершение (после final)
1003Невалидный JSON в стартовом {"after":N}

Примеры

wscat

wscat -H "Authorization: Bearer $TRACIUM_TOKEN" \
  -c "wss://api.tracium.ru/v1/jobs/estimates/$JOB_ID/events"
# после коннекта, чтобы начать не с нуля:
> {"after": 42}

Минимальный клиент на JS

const ws = new WebSocket(
  `wss://api.tracium.ru/v1/jobs/offers-refresh/${jobId}/events`,
  // заголовок Authorization ставится на уровне подключения вашего WS-клиента
);
 
ws.onopen = () => ws.send(JSON.stringify({ after: cursor })); // необязательно
 
ws.onmessage = (ev) => {
  const msg = JSON.parse(ev.data);
  switch (msg.type) {
    case "line":     applyLine(msg.line); cursor = msg.seq; break;
    case "progress": setProgress(msg.total, msg.done, msg.percent); break;
    case "lagged":   /* сервер сам перезальёт строки, ничего не делаем */ break;
    case "final":    finish(msg.job.status); break;
  }
};

Браузерный WebSocket не позволяет задать произвольные заголовки — в вебе ключ прокидывается вашим бэкендом-прокси или согласованным механизмом (не кладите trk_… в URL/JS фронтенда). Серверный клиент (Node ws, Go coder/websocket, wscat) заголовок Authorization задаёт напрямую.

Когда что выбирать

СитуацияТранспорт
UI с прогресс-баром и построчным появлениемWebSocket (/events)
Фоновое/batch дочитывание, консольная утилитаREST-опрос ?after=
Нестабильная сеть, простой ретрайREST-опрос ?after=
Забрать результат позже по job_id (в пределах ~24 ч)REST-опрос ?after=0

Реализация (для сопровождения)

  • Транспорт-handler один на все job-kind’ы: backend/internal/platform/jobstream/ws/handler.go (библиотека github.com/coder/websocket). Kind-специфика вынесена в порт JobReplayer.
  • Адаптеры-реплееры: offers/refresh/api/public-http/di.go (GET /v1/jobs/offers-refresh/{id}/events) и search/proposal/di.go (GET /v1/jobs/estimates/{id}/events). Handler создаётся только при наличии брокера (d.Broker != nil); иначе маршрут не регистрируется.

Связанные документы

  • Cookbook, сценарии async-задач: cookbook.md
  • Живая актуализация цен/остатков: live-refresh.md
  • Обзор сервисной границы: README.md