48 lines
1.8 KiB
Markdown
48 lines
1.8 KiB
Markdown
# 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:
|
|
1. Chains ingestion tasks in correct dependency order
|
|
2. Provides health/status endpoints
|
|
3. Integrates with existing Celery infrastructure
|
|
|
|
## Requirements
|
|
|
|
1. Define dependency graph: prices → sec → financials → news (prices must run first)
|
|
2. Create `src/backend/tasks/pipeline.py` with orchestration logic
|
|
3. Add `/api/v1/sync/status` endpoint (already exists per README — verify it works)
|
|
4. Create Celery chain/group for full pipeline run
|
|
5. Add manual trigger endpoint: `POST /api/v1/sync/pipeline/run`
|
|
6. 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)
|