1.8 KiB
1.8 KiB
Task: Wire Pipeline Orchestrator
Current State
src/data-pipeline/pipeline.py exists but is not connected to Celery Beat or the backend. The four ingestion scripts (prices, sec, financials, news) run standalone — no dependency ordering, no error handling between stages, no visibility into pipeline health.
Goal
Create a proper pipeline orchestrator that:
- Chains ingestion tasks in correct dependency order
- Provides health/status endpoints
- Integrates with existing Celery infrastructure
Requirements
- Define dependency graph: prices → sec → financials → news (prices must run first)
- Create
src/backend/tasks/pipeline.pywith orchestration logic - Add
/api/v1/sync/statusendpoint (already exists per README — verify it works) - Create Celery chain/group for full pipeline run
- Add manual trigger endpoint:
POST /api/v1/sync/pipeline/run - Add to docker-compose worker service
Acceptance Criteria
- Full pipeline runs end-to-end via single API call
- Dependency ordering enforced (prices before sec, etc.)
- Pipeline status visible via
/api/v1/sync/status - Failed stage does not block other independent stages
- Pipeline health check returns last run time, status, errors
Constraints
- Use Celery chains for ordered stages, Celery groups for parallel stages
- Follow existing task patterns — don't reinvent error handling
- Keep under 200 lines per framework rule
- Read existing migrations before writing schema
Files to Create/Modify
src/backend/tasks/pipeline.py(new)src/backend/routers/data_sync.py(add pipeline trigger endpoint)src/backend/celery_app.py(add pipeline schedule)docker-compose.dev.yml(verify worker includes task)
Next Steps After This Task
Wire backtest integration (phase2-backtest-integration)