59 lines
2.1 KiB
Markdown
59 lines
2.1 KiB
Markdown
# Task: Pipeline orchestrator API endpoints
|
|||
|
|
|
||
|
|
## Goal
|
||
|
|
Create the API endpoints for triggering and monitoring pipeline runs. This is Phase 2 — the API layer on top of the task definitions from the previous task.
|
||
|
|
|
||
|
|
## Requirements
|
||
|
|
|
||
|
|
### Create `src/backend/services/pipeline_orchestrator.py`
|
||
|
|
|
||
|
|
Implement `PipelineOrchestrator` class:
|
||
|
|
|
||
|
|
1. `run_pipeline(task_name: str, user_id: str) -> str` — start a pipeline run, return run_id
|
||
|
|
2. `get_run_status(run_id: str) -> dict` — return run status (pending/running/completed/failed)
|
||
|
|
3. `cancel_run(run_id: str)` — cancel a running pipeline
|
||
|
|
4. Internal: execute tasks respecting dependency order
|
||
|
|
5. Store run state in memory (dict) — no DB needed yet
|
||
|
|
|
||
|
|
### Create `src/backend/routers/pipeline.py`
|
||
|
|
|
||
|
|
Add endpoints:
|
||
|
|
|
||
|
|
1. `POST /api/v1/pipeline/run` — trigger a pipeline run
|
||
|
|
- Body: `{"task": "news_ingestion"}` or `"full"` for all tasks
|
||
|
|
- Returns: `{"run_id": "...", "status": "pending"}`
|
||
|
|
|
||
|
|
2. `GET /api/v1/pipeline/run/{run_id}` — get run status
|
||
|
|
- Returns: `{"run_id": "...", "status": "...", "tasks": [...]}`
|
||
|
|
|
||
|
|
3. `POST /api/v1/pipeline/run/{run_id}/cancel` — cancel a run
|
||
|
|
- Returns: `{"message": "cancelled"}`
|
||
|
|
|
||
|
|
4. `GET /api/v1/pipeline/tasks` — list available tasks
|
||
|
|
- Returns: list of registered tasks with descriptions
|
||
|
|
|
||
|
|
### Response models in `schemas/pipeline.py`
|
||
|
|
|
||
|
|
Create:
|
||
|
|
- `PipelineRunRequest` — task name to run
|
||
|
|
- `PipelineRunResponse` — run_id and status
|
||
|
|
- `PipelineRunStatus` — detailed run status with task results
|
||
|
|
- `PipelineTaskInfo` — task metadata
|
||
|
|
|
||
|
|
## Acceptance Criteria
|
||
|
|
1. All 4 endpoints work correctly
|
||
|
|
2. Dependency order is respected when running full pipeline
|
||
|
|
3. Run status updates in real-time
|
||
|
|
4. Files stay under 200 lines each
|
||
|
|
|
||
|
|
## Files to Create/Modify
|
||
|
|
- `src/backend/services/pipeline_orchestrator.py`
|
||
|
|
- `src/backend/services/pipeline_tasks.py` (from previous task)
|
||
|
|
- `src/backend/routers/pipeline.py`
|
||
|
|
- `src/backend/schemas/pipeline.py`
|
||
|
|
|
||
|
|
## Files to Read First
|
||
|
|
- `src/backend/services/pipeline_tasks.py` — task definitions
|
||
|
|
- `src/backend/routers/alerts.py` — follow routing pattern
|
||
|
|
- `src/backend/schemas/alert.py` — follow schema pattern
|