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(), enumerate(), fetch(), b64encode() / b64decode(), json_encode() / json_decode(), setupProcess(), runProcess(), общий ctx
  • Артефакты: глобальный модуль artefacts — загрузка, извлечение, создание ZIP-архивов, сохранение в песочницу
  • Файловая система: глобальный модуль files — чтение, запись, удаление файлов во временном каталоге, экспорт в песочницу
  • Двойная аутентификация: JWT (пользователи) и Service Token (сервисы)
  • Swagger UI

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

СтатусОписание
idleСоздан, ожидает запуска
pendingПодготовка к запуску
runningВыполняется
waitingПауза на waitNode (delay/until/event); resume через таймер, NATS или POST /pipelines/run/:id/resume
completedУспешно завершён
failedЗавершён с ошибкой
cancelledОтменён

Типы нод

Control

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

Input

НодаОписание
inputButtonNodeВход по нажатию кнопки. В тестовом режиме — pass-through (выполняется без логики).
inputEventNodeТриггер: старт instance по NATS subject (data.subject). Payload → ctx.trigger.
inputFormNodeВход по отправке формы. В тестовом режиме — pass-through.
inputHttpRequestNodeТриггер: POST /api/v1/hooks/pipelines/:id (+ X-Hook-Secret или ?secret=). Async 202 {instance_id}; payload → ctx.trigger (body, json, method, query, опционально headers).

Display

НодаОписание
displayTextNodeТекст-аннотация на canvas (runtime no-op / лог).
displayMediaNodeImage/audio/video на canvas (url или file_id).
logMessageNodeПишет message) в логи ноды.

Transform

НодаОписание
transformCodeNodeВыполняет Starlark-скрипт (Python-подобный язык) над общим контекстом. Доступные функции описаны ниже в разделе transformCodeNode.
transformSetVariableNodeУстанавливает одну или несколько переменных в контексте из JSON-источника (data.source).
transformExtractArtifactNodeИзвлекает артефакты из выполненного jobNode. В тестовом режиме — pass-through.
transformMergeProcessNodeСинхронизация: ожидает завершения всех параллельных веток runProcess перед продолжением.
transformCollectProcessNodeКак merge, но собирает из каждой ветки только указанные ключи ctx в список (output_key, по умолчанию process.results).
jsonPathNodeИзвлекает поля из всего ctx (engine: jsonpath|jq, выражение от корня контекста).

Logic

НодаОписание
logicConditionNodeВычисляет Starlark-выражения как условия и направляет выполнение на соответствующий выходной порт (input-N). Если ни одно условие не выполнено — переход на порт error.
logicSwitchNodeУзел маршрутизации/переключения. В тестовом режиме — pass-through.
waitNodeПауза: mode=delay|until|event. Production — park (waiting) + resume; test — blocking delay (cap 30s).

Platform

НодаОписание
createTaskNode / updateTaskNode / changeTaskStatusNodeTask Tracker CRUD / status transition.
createDocumentNodeСоздаёт документ.
sendNotificationNodePOST notifications /send по шаблону.
runSubflowNodeNested Prepare+Start другого пайплайна; depth ≤ 3.

Network

НодаОписание
httpRequestNodeИсходящий GET/POST с allowlist/SSRF (WEB_DOMAIN_ALLOWLIST).
fetchUrlNodeGET URL → текст.
webSearchNodeПоиск (SearXNG и т.п.) без Agent.

Output

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

Job

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

AI

Нативные ноды через AiRouterClient (формы в UI). Для кастомной логики по-прежнему доступен Code-нода и up.ai_router в Starlark.

НодаОписание
aiChatCompletionNodeChat completion: promptai.response (+ ai.usage). Поддержка из ctx.
aiRagSearchNodeПоиск по knowledge base (scopes + query) → ai.results.
aiRagIngestNodeИндексация документа в scope → ai.document_id, ai.chunk_count.
aiAgentNodeСинхронный agent loop (POST /agent/run) с tools/scopes → ai.response, ai.tool_calls.
aiProvidersNodelistai.providers; set_activeai.provider.

При ошибке ноды пишут сообщение в ai.error и падают как failed.

Retry

НодаОписание
retryNodeОстанавливает ветку после N проходов через ноду (лимит в data.retries).

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

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

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

ПеременнаяОписание
artefacts.previousСписок ID артефактов от последнего jobNode
artefacts.allВсе артефакты, накопленные с начала пайплайна
ai.responseТекст ответа AI Chat / Agent
ai.usageUsage tokens
ai.resultsРезультаты RAG Search
ai.document_id / ai.chunk_countРезультат RAG Ingest
ai.providers / ai.providerСписок / обновлённый провайдер
ai.tool_callsВызовы tools от Agent
ai.errorОшибка последней AI-ноды

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

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


transformCodeNode

Нода выполняет Starlark-скрипт (Python-подобный sandboxed-язык) в контексте текущего выполнения пайплайна. Скрипт имеет доступ к общему контексту (ctx), глобальным модулям (files, artefacts) и набору встроенных функций (print, setVar/getVar, enumerate, fetch, b64encode/b64decode, json_encode/json_decode, setupProcess/runProcess).

Полное описание всех функций, контекста и хелперов — в 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Формат логов