Skip to content

Flow Manager

Назначение

Сервис создания, хранения, версионирования и исполнения процессных пайплайнов (flows). Оркестрация пайплайнов с поддержкой Starlark-скриптов, in-memory тестового запуска и real-time стриминга событий через NATS.

Порты

  • HTTP: 8098
  • БД: flow_manager_db

Ключевые возможности

  • CRUD пайплайнов с контролем доступа по владельцу
  • Версионированные шаблоны — каждое сохранение создаёт иммутабельную версию
  • Модель графа: DAG из нод, связанных рёбрами (edges)
  • In-memory тестовый запуск (без записи в БД) с потоковой передачей событий через NATS
  • Production-запуск с сохранением в БД и отслеживанием статуса
  • Тестирование отдельных нод (code, job) через /try/...
  • Ограничение параллельных тестовых запусков: максимум 50 одновременных
  • Starlark-исполнитель: print(), setVar(), getVar(), fetch(), setupProcess(), runProcess(), общий ctx
  • Артефакты: глобальный модуль artefacts — загрузка, извлечение, создание ZIP-архивов, сохранение в песочницу
  • Файловая система: глобальный модуль files — чтение, запись, удаление файлов во временном каталоге, экспорт в песочницу
  • Двойная аутентификация: JWT (пользователи) и Service Token (сервисы)
  • Swagger UI

Статусы инстансов

СтатусОписание
idleСоздан, ожидает запуска
pendingПодготовка к запуску
runningВыполняется
completedУспешно завершён
failedЗавершён с ошибкой
cancelledОтменён

Типы нод

Control

НодаОписание
startNodeТочка входа пайплайна. Не выполняет никакой логики, только помечает начало выполнения.

Input

НодаОписание
inputButtonNodeВход по нажатию кнопки. В тестовом режиме — pass-through (выполняется без логики).
inputEventNodeВход по событию (NATS). В тестовом режиме — pass-through.
inputFormNodeВход по отправке формы. В тестовом режиме — pass-through.
inputHttpRequestNodeВход по HTTP-запросу. В тестовом режиме — pass-through.

Transform

НодаОписание
transformCodeNodeВыполняет Starlark-скрипт (Python-подобный язык) над общим контекстом. Доступные функции описаны ниже в разделе transformCodeNode.
transformSetVariableNodeУстанавливает одну или несколько переменных в контексте из JSON-источника (data.source).
transformExtractArtifactNodeИзвлекает артефакты из выполненного jobNode. В тестовом режиме — pass-through.
transformMergeProcessNodeСинхронизация: ожидает завершения всех параллельных веток runProcess перед продолжением.

Logic

НодаОписание
logicConditionNodeВычисляет Starlark-выражения как условия и направляет выполнение на соответствующий выходной порт (input-N). Если ни одно условие не выполнено — переход на порт error.
logicSwitchNodeУзел маршрутизации/переключения. В тестовом режиме — pass-through.

Output

НодаОписание
outputHttpNodeHTTP-выход (отправка ответа). В тестовом режиме — pass-through.
outputEmailNodeEmail-выход (отправка письма). В тестовом режиме — pass-through.

Job

НодаОписание
jobNodeСоздаёт и мониторит задачу через API command-manager. Ожидает завершения (таймаут 5 минут, интервал опроса 2 сек). Собирает артефакты по завершении.

Контекст выполнения

Каждый пайплайн имеет общий контекст (execution context) — хранилище переменных, передаваемое между нодами. При переходе к следующей ноде контекст клонируется, что позволяет ветвление с изолированными состояниями.

Автоматические переменные

ПеременнаяОписание
artefacts.previousСписок ID артефактов от последнего jobNode
artefacts.allВсе артефакты, накопленные с начала пайплайна

Рабочая директория

Каждый инстанс пайплайна получает временную директорию ARTIFACTS_TMP_DIR/<instance_id>. Директория удаляется после выполнения ноды. Модули files и artefacts работают только внутри этой директории.


transformCodeNode

Нода выполняет Starlark-скрипт (Python-подобный sandboxed-язык) в контексте текущего выполнения пайплайна. Скрипт имеет доступ к общему контексту (ctx), глобальным модулям (files, artefacts) и набору встроенных функций.

Полное описание всех функций, контекста и хелперов — в Starlark API.

WebSocket Events

Flow Manager публикует события в NATS. Gateway подписывается на эти subjects и транслирует их клиентам через WebSocket.

  • Тестовые запуски: pipeline.test.<pipe_uuid>.node.*
  • Production Run: pipeline.run.<instance_id>.node.*
  • Платформенные: job.created, job.changed, executor.status, runner.change

Полное описание всех событий с пейлоадами — в WebSocket Events.

Конфигурация

ПеременнаяПо умолчаниюОписание
PORT8098Порт сервера
HOST0.0.0.0Адрес сервера
GIN_MODEdebugРежим Gin (debug/release)
DB_HOSTlocalhostАдрес PostgreSQL
DB_PORT5432Порт PostgreSQL
DB_USERflow_manager_userПользователь PostgreSQL
DB_PASSWORDПароль PostgreSQL
DB_NAMEflow_manager_dbИмя базы данных
DB_SSL_MODEdisableРежим SSL
JWT_SECRETСекрет для JWT
JWT_ALGORITHMHS256Алгоритм JWT
JWT_TTL24hВремя жизни JWT
SERVICE_TOKENТокен межсервисного взаимодействия
MESSAGING_URLnats://localhost:4222URL NATS
ARTIFACTS_TMP_DIR./tmp/flow-managerКаталог для временных файлов артефактов
SANDBOX_DIR./sandboxПостоянный каталог для files.export() и artefacts.save() (структура: {SANDBOX_DIR}/{pipeline_id}/)
TEST_PIPELINE_TTL10mВремя жизни тестового пайплайна в памяти
LOG_LEVELinfoУровень логирования
LOG_FORMATjsonФормат логов