2.1 KiB
2.1 KiB
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:
run_pipeline(task_name: str, user_id: str) -> str— start a pipeline run, return run_idget_run_status(run_id: str) -> dict— return run status (pending/running/completed/failed)cancel_run(run_id: str)— cancel a running pipeline- Internal: execute tasks respecting dependency order
- Store run state in memory (dict) — no DB needed yet
Create src/backend/routers/pipeline.py
Add endpoints:
-
POST /api/v1/pipeline/run— trigger a pipeline run- Body:
{"task": "news_ingestion"}or"full"for all tasks - Returns:
{"run_id": "...", "status": "pending"}
- Body:
-
GET /api/v1/pipeline/run/{run_id}— get run status- Returns:
{"run_id": "...", "status": "...", "tasks": [...]}
- Returns:
-
POST /api/v1/pipeline/run/{run_id}/cancel— cancel a run- Returns:
{"message": "cancelled"}
- Returns:
-
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 runPipelineRunResponse— run_id and statusPipelineRunStatus— detailed run status with task resultsPipelineTaskInfo— task metadata
Acceptance Criteria
- All 4 endpoints work correctly
- Dependency order is respected when running full pipeline
- Run status updates in real-time
- Files stay under 200 lines each
Files to Create/Modify
src/backend/services/pipeline_orchestrator.pysrc/backend/services/pipeline_tasks.py(from previous task)src/backend/routers/pipeline.pysrc/backend/schemas/pipeline.py
Files to Read First
src/backend/services/pipeline_tasks.py— task definitionssrc/backend/routers/alerts.py— follow routing patternsrc/backend/schemas/alert.py— follow schema pattern