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}/events | Async-смета (/v1/estimates/jobs) | search:read |
GET /v1/jobs/offers-refresh/{id}/events | Async-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), поэтому оно не
укладывается в модель «один запрос — один ответ».
Протокол
- Апгрейд. Клиент коннектится, сервер апгрейдит соединение на WebSocket.
- Курсор (опционально). Первым текстовым сообщением клиент может прислать
{"after":N}— реплей начнётся со строки послеN.- дедлайн на это сообщение — 5 секунд;
- нет сообщения или таймаут → трактуется как
after=0(реплей с начала); - невалидный JSON → сервер закрывает соединение кодом 1003.
- Подписка → реплей. Сервер сначала подписывается на брокер (чтобы не
потерять события, опубликованные между подпиской и реплеем), затем реплеит
пропущенные строки
LinesSince(after)как{"type":"line","seq":N,"line":<payload>}. - Уже терминальная задача. Если задача завершилась (и реплей отдан) — сервер
синтезирует
{"type":"final","job":{"status":…}}и закрывает кодом 1000. - Live-tail. Иначе сервер тейлит канал брокера:
line— новая готовая строка (дедуп поseq);progress— обновление прогресса (total,done,percent);final— терминальный статус; после отправки соединение закрывается 1000;lagged— клиент отстал и был вытеснен из буфера брокера: сервер шлёт{"type":"lagged"}и повторно реплеитLinesSince(lastSeq)— без переподключения со стороны клиента.
- Keepalive. Ping раз в 30 с; если pong не пришёл за 60 с — соединение закрывается.
Типы сообщений (сервер → клиент)
type | Payload | Смысл |
|---|---|---|
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 фронтенда). Серверный клиент (Nodews, Gocoder/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