Appearance
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 / лог). |
displayMediaNode | Image/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 / changeTaskStatusNode | Task Tracker CRUD / status transition. |
createDocumentNode | Создаёт документ. |
sendNotificationNode | POST notifications /send по шаблону. |
runSubflowNode | Nested Prepare+Start другого пайплайна; depth ≤ 3. |
Network
| Нода | Описание |
|---|---|
httpRequestNode | Исходящий GET/POST с allowlist/SSRF (WEB_DOMAIN_ALLOWLIST). |
fetchUrlNode | GET URL → текст. |
webSearchNode | Поиск (SearXNG и т.п.) без Agent. |
Output
| Нода | Описание |
|---|---|
outputHttpNode | HTTP-выход (отправка ответа). В тестовом режиме — pass-through. |
outputEmailNode | Email-выход (отправка письма). В тестовом режиме — pass-through. |
Job
| Нода | Описание |
|---|---|
jobNode | Создаёт и мониторит задачу через API command-manager. Ожидает завершения (таймаут 5 минут, интервал опроса 2 сек). Собирает артефакты по завершении. |
AI
Нативные ноды через AiRouterClient (формы в UI). Для кастомной логики по-прежнему доступен Code-нода и up.ai_router в Starlark.
| Нода | Описание |
|---|---|
aiChatCompletionNode | Chat completion: prompt → ai.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. |
aiProvidersNode | list → ai.providers; set_active → ai.provider. |
При ошибке ноды пишут сообщение в ai.error и падают как failed.
Retry
| Нода | Описание |
|---|---|
retryNode | Останавливает ветку после N проходов через ноду (лимит в data.retries). |
Контекст выполнения
Каждый пайплайн имеет общий контекст (execution context) — хранилище переменных, передаваемое между нодами. При переходе к следующей ноде контекст клонируется, что позволяет ветвление с изолированными состояниями.
Автоматические переменные
| Переменная | Описание |
|---|---|
artefacts.previous | Список ID артефактов от последнего jobNode |
artefacts.all | Все артефакты, накопленные с начала пайплайна |
ai.response | Текст ответа AI Chat / Agent |
ai.usage | Usage 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.
Конфигурация
| Переменная | По умолчанию | Описание |
|---|---|---|
PORT | 8098 | Порт сервера |
HOST | 0.0.0.0 | Адрес сервера |
GIN_MODE | debug | Режим Gin (debug/release) |
DB_HOST | localhost | Адрес PostgreSQL |
DB_PORT | 5432 | Порт PostgreSQL |
DB_USER | flow_manager_user | Пользователь PostgreSQL |
DB_PASSWORD | — | Пароль PostgreSQL |
DB_NAME | flow_manager_db | Имя базы данных |
DB_SSL_MODE | disable | Режим SSL |
JWT_SECRET | — | Секрет для JWT |
JWT_ALGORITHM | HS256 | Алгоритм JWT |
JWT_TTL | 24h | Время жизни JWT |
SERVICE_TOKEN | — | Токен межсервисного взаимодействия |
MESSAGING_URL | nats://localhost:4222 | URL NATS |
ARTIFACTS_TMP_DIR | ./tmp/flow-manager | Каталог для временных файлов артефактов |
SANDBOX_DIR | ./sandbox | Постоянный каталог для files.export() и artefacts.save() (структура: {SANDBOX_DIR}/{pipeline_id}/) |
TEST_PIPELINE_TTL | 10m | Время жизни тестового пайплайна в памяти |
LOG_LEVEL | info | Уровень логирования |
LOG_FORMAT | json | Формат логов |