Skip to content

Latest commit

 

History

17 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Automata

Распределённая система автоматизации бизнес-процессов, предназначенная для описания, тестирования и безопасного внедрения сценариев обработки данных и интеграций с внешними сервисами.

Возможности

  • Описание и выполнение бизнес-процессов в виде сценариев (flows)
  • Поддержка DAG-зависимостей между шагами
  • PR-workflow для контроля изменений в сценариях
  • Sandbox для тестового выполнения перед продом
  • Запуск по расписанию (cron/interval)
  • Retry с exponential backoff
  • Dead Letter Queue (DLQ) для обработки ошибок (automata dlq list/show/replay)
  • Горизонтальное масштабирование воркеров

Компоненты

Компонент Порт Описание
API Server :8080 REST API для управления flows, runs, schedules
Scheduler :8081 Планировщик с leader election, создаёт runs по расписанию
Orchestrator :8083 Парсит DAG, создаёт tasks, управляет выполнением
Worker :8082 Выполняет tasks (HTTP, delay, transform)
CLI Утилита командной строки для пользователей

Потоки данных

  1. API → создаёт flows, runs, schedules в PostgreSQL
  2. Scheduler → опрашивает due schedules → создаёт runs → публикует в RabbitMQ
  3. Orchestrator → потребляет runs → парсит DAG → создаёт tasks → публикует в RabbitMQ
  4. Worker → потребляет tasks → выполняет → публикует результат

Принцип устойчивости

PostgreSQL — source of truth, RabbitMQ — оптимизация.

┌─────────────┐     ┌──────────────┐     ┌──────────────┐
│  Scheduler  │────▶│  PostgreSQL  │────▶│ Orchestrator │
│             │     │   (runs)     │     │   (polling)  │
└──────┬──────┘     └──────────────┘     └──────▲───────┘
       │                                        │
       │            ┌──────────────┐            │
       └───────────▶│   RabbitMQ   │────────────┘
                    │ (run.pending)│  event-driven
                    └──────────────┘
  • Нормальный режим: RabbitMQ доставляет события мгновенно
  • При сбое MQ: Orchestrator подхватывает runs через polling из БД
  • Гарантия: ни один run не потеряется, даже если RabbitMQ недоступен

Структура проекта

Automata/
├── cmd/
│   ├── automata-api/           # HTTP API сервер
│   ├── automata-scheduler/     # Планировщик задач
│   ├── automata-orchestrator/  # Оркестратор выполнения
│   ├── automata-worker/        # Воркер
│   └── automata-cli/           # CLI утилита
│
├── internal/
│   ├── domain/       # Доменные модели (Flow, Run, Task, Schedule)
│   ├── repo/         # PostgreSQL репозитории
│   ├── mq/           # RabbitMQ (connection, publisher, consumer)
│   ├── engine/       # Парсер FlowSpec, DAG, templates
│   ├── steps/        # Реализации шагов (http, delay, transform)
│   ├── scheduler/    # Логика планировщика
│   ├── orchestrator/ # Управление состоянием run
│   ├── worker/       # Выполнение tasks
│   ├── api/          # HTTP handlers, middleware, DTOs
│   ├── sandbox/      # Изолированное выполнение
│   ├── config/       # Конфигурация
│   └── telemetry/    # Логирование, метрики
│
├── tests/
│   └── e2e/          # E2E тесты (testcontainers)
│
├── migrations/       # SQL миграции
└── deploy/           # Docker Compose

Технологии

Компонент Технология
Язык Go 1.25
HTTP Router net/http
Логирование log/slog
PostgreSQL pgx/v5
RabbitMQ amqp091-go
Cron robfig/cron/v3
UUID google/uuid
Метрики prometheus/client_golang
Мониторинг Prometheus + Grafana
CLI github.com/spf13/cobra
E2E тесты testcontainers-go

Observability

Система экспортирует Prometheus метрики со всех 4 сервисов на эндпоинте /metrics.

Метрики

Сервис Метрика Тип Описание
API automata_api_http_requests_total counter Количество HTTP-запросов по method/path/status
API automata_api_http_request_duration_seconds histogram Latency запросов
Orchestrator automata_orchestrator_runs_total counter Завершённые runs по статусу
Orchestrator automata_orchestrator_runs_active gauge Активные runs в памяти
Orchestrator automata_orchestrator_dispatch_duration_seconds histogram Время dispatch шагов
Worker automata_worker_tasks_total counter Задачи по step_type/status
Worker automata_worker_task_duration_seconds histogram Время выполнения задач
Worker automata_worker_task_retries_total counter Количество retry
Scheduler automata_scheduler_ticks_total counter Тики планировщика
Scheduler automata_scheduler_runs_created_total counter Runs, созданные по расписанию

Grafana

После make up-all дашборд доступен на http://localhost:3000 (admin/automata).

Дашборд "Automata Overview" содержит 9 панелей: request rate, latency p95, runs по статусам, active runs, task duration по типам, tasks processed, retries, dispatch duration, scheduler runs.


Flow Spec

Flow — это "шаблон" или "рецепт" автоматизации. Один flow может иметь множество версий (FlowVersion). Каждый запуск (Run) выполняет конкретную версию flow. Flow описывается в формате JSON:

{
  "name": "sync-orders",
  "inputs": {
    "source_id": { "type": "string", "required": true }
  },
  "defaults": {
    "retry": { "max_attempts": 3, "backoff": "exponential" },
    "timeout_sec": 300
  },
  "steps": [
    {
      "id": "fetch",
      "type": "http",
      "config": {
        "method": "GET",
        "url": "https://crm.example.com/orders"
      },
      "outputs": {
        "orders": "{{ .response.body.data }}"
      }
    },
    {
      "id": "save",
      "type": "http",
      "depends_on": ["fetch"],
      "config": {
        "method": "POST",
        "url": "https://erp.example.com/import",
        "body": "{{ .steps.fetch.outputs.orders }}"
      }
    }
  ]
}

Типы шагов

Тип Описание
http HTTP запросы к внешним API
delay Пауза между шагами
transform Трансформация данных
parallel Параллельное выполнение веток

Фазы реализации

Фаза 0: Подготовка

  • Добавить зависимости: amqp091-go, google/uuid
  • Создать структуру internal/ пакетов
  • Настроить structured logging (slog)

Фаза 1: Domain + Repository

  • internal/domain/* — доменные структуры (Flow, Run, Task, Schedule, Proposal)
  • internal/repo/* — CRUD для flows, runs, tasks, schedules, proposals

Фаза 2: RabbitMQ

  • Подключение с reconnect
  • Publisher и Consumer
  • Topology (exchanges, queues)

Фаза 3: REST API

  • CRUD /flows, /runs, /schedules
  • Версионирование flows

Фаза 4: Scheduler

  • Leader election (advisory locks)
  • Обработка due schedules
  • Cron parsing (robfig/cron/v3)
  • Публикация run.pending в RabbitMQ

Фаза 5: Engine + Steps

  • Парсер FlowSpec
  • DAG структура
  • Go templates для данных
  • Реализация шагов

Фаза 6: Orchestrator

  • Гибридный подход: Event-driven + Polling
    • Event-driven: слушает runs.pending из RabbitMQ (низкая latency)
    • Polling: периодически проверяет ListPending() (fallback при сбое MQ)
  • State machine для run (PENDING → RUNNING → SUCCEEDED/FAILED)
  • Управление зависимостями между tasks (DAG)
  • RunState для управления состоянием в памяти
  • Восстановление состояния после рестарта

Фаза 7: Worker

  • Executor interface + Registry (http, delay, transform)
  • HTTPExecutor (GET/POST/PUT/DELETE, headers, body, timeout, JSON parsing)
  • DelayExecutor (context-aware timer)
  • TransformExecutor (pass-through rendered payload)
  • Retry с exponential backoff (in-process, OnStatus для HTTP)
  • Гибридный подход: Consumer (tasks.ready) + Polling fallback
  • Graceful shutdown, publisher nil-safety
  • Полная точка входа cmd/automata-worker

Фаза 8: CLI

  • Cobra CLI framework (github.com/spf13/cobra)
  • HTTP-клиент для API (Client с полным покрытием эндпоинтов)
  • Форматирование вывода (таблицы + --json)
  • Команды flow: list, create, show, update, delete, versions, publish
  • Команды run: list, start, show, cancel, tasks
  • Команды schedule: list, create, show, update, delete, enable, disable
  • Точка входа cmd/automata-cli с PersistentFlags (--api-url, --json)

Фаза 9: PR-Workflow + Sandbox

  • Proposal API (CRUD + submit/approve/reject/apply/sandbox)
  • Proposal CLI (proposal list, create, show, update, delete, submit, approve, reject, apply, sandbox)
  • Sandbox пакет (Collector, Differ, WaitForResult, CompareWithBaseline)
  • Полный PR-workflow: DRAFT → PENDING_REVIEW → APPROVED → APPLIED
  • Sandbox run с is_sandbox=true, сбор результатов, deep diff
  • spec_override в runs — sandbox запускает ProposedSpec напрямую, без создания временной версии flow

Фаза 10: E2E Testing

  • Инфраструктура: testcontainers (Postgres + RabbitMQ), in-process сервисы
  • APIClient, mock HTTP-сервер, фабрики спецификаций
  • 17 тестов: CRUD, runs (basic, chain, parallel, retry, cancel), schedules, proposals, sandbox, multi-flow

Фаза 11: Observability

  • Prometheus метрики для всех сервисов (API, Orchestrator, Worker, Scheduler)
  • HTTP metrics middleware с нормализацией путей (UUID → :id)
  • Prometheus + Grafana в docker-compose
  • Grafana дашборд "Automata Overview" (9 панелей)

Быстрый старт

make up-all       # Собрать образы и поднять всё: Postgres, RabbitMQ, миграции, 4 сервиса
make down-all     # Остановить всё и удалить контейнеры
make logs-all     # Логи всех сервисов (follow, последние 100 строк)
make test         # Unit тесты (без Docker)
make test-e2e     # E2E тесты (поднимает контейнеры через testcontainers)

CLI-команды

CLI собирается из cmd/automata-cli/ и взаимодействует с API через HTTP.

go build -o automata ./cmd/automata-cli/

Глобальные флаги:

  • --api-url — адрес API (по умолчанию http://localhost:8080)
  • --json — вывод в JSON-формате

Flows

automata flow list                          # Список всех flows
automata flow create --name "my-flow"       # Создать flow
automata flow show <ID>                     # Детали flow
automata flow update <ID> --name "new-name" # Обновить имя
automata flow update <ID> --active true     # Активировать
automata flow delete <ID>                   # Удалить
automata flow versions <ID>                 # Список версий
automata flow publish <ID> --spec-file f.json  # Опубликовать версию

Runs

automata run list --flow-id <FLOW_ID>       # Список runs
automata run start <FLOW_ID>                # Запустить run
automata run start <FLOW_ID> --input "key=value" --input "k2=v2"  # С параметрами
automata run start <FLOW_ID> --sandbox      # Запуск в sandbox
automata run start <FLOW_ID> --version 2    # Конкретная версия
automata run show <RUN_ID>                  # Детали run
automata run tasks <RUN_ID>                 # Список задач в run
automata run cancel <RUN_ID>                # Отменить run

Proposals (PR-workflow)

automata proposal list                              # Список proposals
automata proposal list --flow-id <FLOW_ID>          # Фильтр по flow
automata proposal list --status DRAFT               # Фильтр по статусу
automata proposal create <FLOW_ID> --title "Update fetch step" --spec-file spec.json
automata proposal create <FLOW_ID> --title "Fix" --spec-file spec.json --created-by "alice"
automata proposal show <ID>                         # Детали proposal
automata proposal update <ID> --title "New title"   # Обновить (только DRAFT)
automata proposal update <ID> --spec-file new.json  # Обновить spec
automata proposal delete <ID>                       # Удалить (только DRAFT)
automata proposal submit <ID>                       # Отправить на review
automata proposal approve <ID> --reviewer "bob"     # Одобрить
automata proposal approve <ID> --reviewer "bob" --comment "LGTM"
automata proposal reject <ID> --reviewer "bob" --comment "needs fix"
automata proposal apply <ID>                        # Применить (создать новую версию flow)
automata proposal sandbox <ID>                      # Запустить тестовый sandbox run

Schedules

automata schedule list                      # Все расписания
automata schedule list --flow-id <ID>       # Расписания для flow
automata schedule create <FLOW_ID> --name "daily" --cron "0 9 * * *"    # По cron
automata schedule create <FLOW_ID> --name "5min" --interval 300          # По интервалу
automata schedule create <FLOW_ID> --name "tz" --cron "0 9 * * *" --timezone "Europe/Moscow"
automata schedule create <FLOW_ID> --name "with-inputs" --cron "* * * * *" --input "city=Moscow"
automata schedule show <ID>                 # Детали расписания
automata schedule update <ID> --name "new"  # Обновить
automata schedule delete <ID>               # Удалить
automata schedule enable <ID>               # Включить
automata schedule disable <ID>              # Выключить

Модель данных

flows ──────────→ flow_versions (spec JSONB)
  │
  ├──────────────→ schedules (cron/interval)
  │
  └──────────────→ runs ──────────→ tasks
                     │
proposals ───────────┘ (PR-workflow)

Статусы

Run: PENDINGRUNNINGSUCCEEDED | FAILED | CANCELLED

Task: QUEUEDRUNNINGSUCCEEDED | FAILED

Proposal: DRAFTPENDING_REVIEWAPPROVEDAPPLIEDREJECTED

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages