Async выполнение сценариев с SSE-прогрессом + каталоги tools/scenarios
POST /api/runs теперь планирует исполнение в фоновой asyncio.Task и
возвращает run_id (202 Accepted) — UI больше не блокируется на время
всего workflow.
Новый модуль src/run_registry.py держит in-memory LRU (лимит
RUN_REGISTRY_MAX_SIZE, default 200) с RunRecord на каждый запуск:
append-only буфер событий для replay + список подписчиков-очередей
для live tail. EventEmitter пишет в буфер и фан-аутит по очередям.
Новые endpoints:
- GET /api/runs/{run_id} снапшот состояния (частичный для running)
- GET /api/runs/{run_id}/events SSE: run_started, step_started,
step_finished, run_finished
- GET /api/scenarios список сценариев с метаданными
- GET /api/scenarios/{id} полное определение для UI-графа
- GET /api/tools проксирование MCP list_tools
mcp_workflow_runner дополнен хуком emitter'а в session_state и
обёрткой run_scenario_async, которая управляет лайфсайклом RunRecord:
queued → running → success/failed + terminal sentinel в очереди
подписчиков. На shutdown lifespan отменяет активные таски.
Все модели в schemas.py и dict-endpoints получили реалистичные
examples для /docs вместо дефолтного additionalProp1.
This commit is contained in:
@@ -73,46 +73,57 @@ cd /home/worker/projects/prisma_platform
|
||||
- `http://127.0.0.1:7777/docs`
|
||||
- `http://127.0.0.1:7777/redoc`
|
||||
|
||||
## Запуск сценария через HTTP
|
||||
## HTTP API
|
||||
|
||||
- `POST http://127.0.0.1:7777/api/runs`
|
||||
### Запуск сценария (async)
|
||||
|
||||
Тело запроса:
|
||||
|
||||
```json
|
||||
{
|
||||
"scenario_id": "news_source_discovery_v1",
|
||||
"input": {
|
||||
"url": "https://example.com/news"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Пример:
|
||||
`POST /api/runs` — планирует выполнение сценария и **сразу** возвращает `run_id`. Само выполнение идёт в фоне.
|
||||
|
||||
```bash
|
||||
curl -s -X POST "http://127.0.0.1:7777/api/runs" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d '{
|
||||
"scenario_id": "news_source_discovery_v1",
|
||||
"input": {
|
||||
"url": "https://example.com/news"
|
||||
}
|
||||
"input": { "url": "https://example.com/news" }
|
||||
}'
|
||||
```
|
||||
|
||||
Успешный ответ содержит:
|
||||
Ответ (`202 Accepted`):
|
||||
|
||||
- `status=success`
|
||||
- `message=""`
|
||||
- список `steps` со статусами и временем шагов
|
||||
- `output_summary`
|
||||
- `result` итогового шага
|
||||
```json
|
||||
{
|
||||
"run_id": "f3d9…",
|
||||
"scenario_id": "news_source_discovery_v1",
|
||||
"status": "queued",
|
||||
"input": { "url": "..." },
|
||||
"started_at": "2026-04-24T..."
|
||||
}
|
||||
```
|
||||
|
||||
При ошибке:
|
||||
### Снапшот состояния
|
||||
|
||||
- `status=failed`
|
||||
- `message` содержит текст ошибки
|
||||
`GET /api/runs/{run_id}` — текущее состояние: `status` (`queued|running|success|failed`), список `steps` со статусами (`success|failed|skipped|queued`), `result` и `output_summary` при завершении.
|
||||
|
||||
### Live-прогресс (SSE)
|
||||
|
||||
`GET /api/runs/{run_id}/events` — Server-Sent Events. Поздние подписчики получают replay уже накопленных событий, затем tail до завершения.
|
||||
|
||||
```bash
|
||||
curl -N http://127.0.0.1:7777/api/runs/$RUN_ID/events
|
||||
```
|
||||
|
||||
Типы событий:
|
||||
|
||||
- `run_started` — `{run_id, scenario_id, started_at}`
|
||||
- `step_started` — `{run_id, step_name, index, started_at}`
|
||||
- `step_finished` — `{run_id, step_name, index, status, started_at, finished_at, message}`
|
||||
- `run_finished` — `{run_id, status, finished_at, message}` (терминальное, поток закрывается)
|
||||
|
||||
### Каталоги
|
||||
|
||||
- `GET /api/scenarios` — список сценариев с метаданными (`scenario_id`, `name`, `description`, `input_schema`).
|
||||
- `GET /api/scenarios/{scenario_id}` — полное определение сценария (для визуализации графа в UI).
|
||||
- `GET /api/tools` — MCP tool catalog: `[{name, description, input_schema}]` (проксируется на `MCP_BASE_URL`).
|
||||
|
||||
## Переменные окружения
|
||||
|
||||
@@ -148,6 +159,11 @@ MCP:
|
||||
- `MCP_BASE_URL` (default: `http://127.0.0.1:8081/mcp`)
|
||||
- `MCP_TIMEOUT_SECONDS` (default: `10`)
|
||||
|
||||
Runtime caches:
|
||||
|
||||
- `WORKFLOW_CACHE_MAX_SIZE` (default: `64`) — лимит LRU кэша построенных workflow.
|
||||
- `RUN_REGISTRY_MAX_SIZE` (default: `200`) — лимит LRU истории run'ов в памяти.
|
||||
|
||||
Phoenix tracing:
|
||||
|
||||
- `PHOENIX_TRACING_ENABLED` (default: `false`)
|
||||
|
||||
Reference in New Issue
Block a user