Skip to content

WebSocket Events

Все real-time события в платформе доставляются клиентам через WebSocket. Источник событий — NATS. Gateway подписывается на subjects и транслирует их WebSocket-клиентам.

Архитектура

Flow Manager / Command Manager / Scheduler

        ▼ (NATS Publish)
      NATS

        ▼ (NATS Subscribe)
     Gateway

        ▼ (WebSocket Broadcast)
     Клиент

Flow Manager — Тестовые запуски

Префикс: pipeline.test.<pipe_uuid>

Все события тестового запуска привязаны к pipe_id (UUID запуска) и user_id (пользователь, запустивший тест).


pipeline.test.<pipe_uuid>.node.status — Статус ноды

Публикуется при каждом изменении статуса ноды: queuedstartedcompleted/failed.

json
{
  "pipe_id": "550e8400-e29b-41d4-a716-446655440000",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "node_id": "770e8400-e29b-41d4-a716-446655440002",
  "status": "completed",
  "error": null,
  "timestamp": "2024-01-15T10:30:00Z"
}
ПолеТипОписание
pipe_iduuidID тестового запуска
user_iduuidID пользователя
node_iduuidID ноды
statusstringqueued / started / completed / failed
errorstring | nullТекст ошибки (только при failed)
timestampstringISO 8601 UTC

Go-структура: entities.TestNodeStatusEvent
Файл: internal/domain/entities/test_pipeline.go:58
Публикация: internal/usecase/event_publisher.go:36


pipeline.test.<pipe_uuid>.node.transition — Переход между нодами

Публикуется при переходе выполнения от одной ноды к другой.

json
{
  "pipe_id": "550e8400-e29b-41d4-a716-446655440000",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "from_node_id": "770e8400-e29b-41d4-a716-446655440002",
  "to_node_id": "880e8400-e29b-41d4-a716-446655440003",
  "timestamp": "2024-01-15T10:30:01Z"
}
ПолеТипОписание
pipe_iduuidID тестового запуска
user_iduuidID пользователя
from_node_iduuidID ноды, с которой выполнен переход
to_node_iduuidID целевой ноды
timestampstringISO 8601 UTC

Go-структура: entities.TestNodeTransitionEvent
Файл: internal/domain/entities/test_pipeline.go:69
Публикация: internal/usecase/event_publisher.go:52


pipeline.test.<pipe_uuid>.node.log — Лог ноды

Публикуется для каждого сообщения лога, сгенерированного нодой (включая print() из Starlark).

json
{
  "pipe_id": "550e8400-e29b-41d4-a716-446655440000",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "node_id": "770e8400-e29b-41d4-a716-446655440002",
  "message": "Executing node 770e8400-e29b-41d4-a716-446655440002 (type: transformCodeNode)",
  "timestamp": "2024-01-15T10:30:00Z"
}
ПолеТипОписание
pipe_iduuidID тестового запуска
user_iduuidID пользователя
node_iduuidID ноды
messagestringТекст лог-сообщения
timestampstringISO 8601 UTC

Go-структура: entities.TestNodeLogEvent
Файл: internal/domain/entities/test_pipeline.go:79
Публикация: internal/usecase/event_publisher.go:67


pipeline.test.<pipe_uuid>.node.error — Ошибка ноды

Публикуется при ошибке выполнения ноды. Структура идентична node.log.

json
{
  "pipe_id": "550e8400-e29b-41d4-a716-446655440000",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "node_id": "770e8400-e29b-41d4-a716-446655440002",
  "message": "fetch failed: connection refused",
  "timestamp": "2024-01-15T10:30:05Z"
}
ПолеТипОписание
pipe_iduuidID тестового запуска
user_iduuidID пользователя
node_iduuidID ноды
messagestringТекст ошибки
timestampstringISO 8601 UTC

Go-структура: entities.TestNodeLogEvent (переиспользуется)
Файл: internal/domain/entities/test_pipeline.go:79
Публикация: internal/usecase/event_publisher.go:82


pipeline.test.<pipe_uuid>.node.job — Создание задачи

Публикуется когда jobNode создаёт задачу через command-manager. Клиент может подписаться на логи задачи по job_id.

json
{
  "pipe_id": "550e8400-e29b-41d4-a716-446655440000",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "node_id": "770e8400-e29b-41d4-a716-446655440002",
  "job_id": "990e8400-e29b-41d4-a716-446655440004",
  "command_id": "aa0e8400-e29b-41d4-a716-446655440005",
  "timestamp": "2024-01-15T10:30:02Z"
}
ПолеТипОписание
pipe_iduuidID тестового запуска
user_iduuidID пользователя
node_iduuidID job-ноды
job_iduuidID созданной задачи
command_iduuidID команды
timestampstringISO 8601 UTC

Go-структура: entities.TestNodeJobCreatedEvent
Файл: internal/domain/entities/test_pipeline.go:89
Публикация: internal/usecase/event_publisher.go:97


Flow Manager — Production Run

Префикс: pipeline.run.<instance_id>

Production-запуски используют ту же структуру событий, но с префиксом pipeline.run.<instance_id>:

SubjectТип события
pipeline.run.<instance_id>.node.statusСтатус ноды
pipeline.run.<instance_id>.node.transitionПереход между нодами
pipeline.run.<instance_id>.node.logЛог ноды
pipeline.run.<instance_id>.node.errorОшибка ноды
pipeline.run.<instance_id>.node.jobСоздание задачи

Payload идентичен тестовым запускам, за исключением:

  • pipe_idinstance_id (UUID инстанса из БД)
  • user_id присутствует для корреляции

Платформенные события

Gateway подписывается на эти subjects и транслирует всем WebSocket-клиентам.


job.created — Задача создана

Публикуется при создании новой задачи в command-manager.

json
{
  "id": "990e8400-e29b-41d4-a716-446655440004",
  "command_id": "aa0e8400-e29b-41d4-a716-446655440005",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "runner_id": "bb0e8400-e29b-41d4-a716-446655440006",
  "execution_number": 1,
  "status": "pending",
  "created_at": "2024-01-15T10:30:00Z",
  "updated_at": "2024-01-15T10:30:00Z"
}
ПолеТипОписание
iduuidID задачи
command_iduuidID команды
user_iduuidID пользователя
runner_iduuid | nullID ренератора (если назначен)
execution_numberintНомер попытки выполнения
statusstringСтатус задачи
created_atstringISO 8601 UTC
updated_atstringISO 8601 UTC

Go-структура: events.JobCreatedEvent
Файл: service-common/events/job_events.go:10


job.changed — Статус задачи обновлён

Публикуется при изменении статуса задачи (выполнение, завершение, ошибка).

json
{
  "id": "990e8400-e29b-41d4-a716-446655440004",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "status": "completed",
  "updated_at": "2024-01-15T10:35:00Z",
  "output": ["Build successful", "Deployed to staging"],
  "error": []
}
ПолеТипОписание
iduuidID задачи
user_iduuidID пользователя
statusstringpending / running / completed / failed / stopped
updated_atstringISO 8601 UTC
outputstring[]Строки вывода задачи
errorstring[]Строки ошибок

Go-структура: events.JobChangedEvent
Файл: service-common/events/job_events.go:29


executor.status — Статус исполнителя

Публикуется при изменении статуса исполнителя задачи (runner).

json
{
  "id": "cc0e8400-e29b-41d4-a716-446655440007",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "status": "online",
  "runner_id": "bb0e8400-e29b-41d4-a716-446655440006",
  "job_id": "990e8400-e29b-41d4-a716-446655440004",
  "output": ["Step 1/5: Compiling..."],
  "error": "",
  "updated_at": "2024-01-15T10:30:05Z"
}
ПолеТипОписание
iduuidID исполнителя
user_iduuidID пользователя
statusstringСтатус исполнителя
runner_iduuid | nullID ренератора
job_iduuid | nullID текущей задачи
outputstring[]Строки вывода
errorstringТекст ошибки
updated_atstringISO 8601 UTC

Go-структура: events.ExecutorStatusEvent
Файл: service-common/events/executor_events.go:11


runner.change — Изменение ренератора

Публикуется при подключении/отключении/обновлении ренератора.

json
{
  "id": "bb0e8400-e29b-41d4-a716-446655440006",
  "user_id": "660e8400-e29b-41d4-a716-446655440001",
  "status": "online",
  "tags": "linux,docker",
  "descriptions": "Production runner",
  "session_timeout": 300,
  "concurrency": 5,
  "type": "shell",
  "version": "1.2.0",
  "name": "runner-prod-1",
  "insecure": false,
  "untagged": false,
  "temporary": false,
  "ip": "192.168.1.100",
  "jobs": 2,
  "last_connection": "2024-01-15T10:30:00Z",
  "created_at": "2024-01-01T00:00:00Z",
  "updated_at": "2024-01-15T10:30:00Z"
}
ПолеТипОписание
iduuidID ренератора
user_iduuidID владельца
statusstringonline / offline / busy
tagsstringТеги через запятую
descriptionsstringОписание ренератора
session_timeoutintТаймаут сессии (сек)
concurrencyintМакс. параллельных задач
typestringТип ренератора (shell, docker)
versionstringВерсия софта ренератора
namestringИмя ренератора
insecureboolРежим небезопасного выполнения
untaggedboolПринимает задачи без тегов
temporaryboolВременный ренератор
ipstringIP-адрес ренератора
jobsintКоличество текущих задач
last_connectionstringПоследнее подключение (ISO 8601)
created_atstringДата создания (ISO 8601)
updated_atstringДата обновления (ISO 8601)

Go-структура: events.RunnerChandeEvent
Файл: service-common/events/runner_events.go:10


Жизненный цикл событий тестового запуска

Типичная последовательность событий для одного пайплайна:

1. node.status    { status: "queued"    }  — нода поставлена в очередь
2. node.transition { from: prev, to: node } — переход от предыдущей ноды
3. node.status    { status: "started"   }  — нода начала выполнение
4. node.log       { message: "..."      }  — лог-сообщения (ноль или более)
5. node.job       { job_id: "..."       }  — (только для jobNode) задача создана
6. node.status    { status: "completed" }  — нода завершена успешно
                  ИЛИ
6. node.error     { message: "..."      }  — текст ошибки
7. node.status    { status: "failed"    }  — нода завершена с ошибкой

Для параллельных веток (runProcess) события разных веток публикуются параллельно и перемешиваются в одном WebSocket-потоке. Идентификация ветки — по node_id в каждом событии.