1.9 KiB
1.9 KiB
Task: Pipeline task definitions
Goal
Create the task definition classes and registry for the pipeline orchestrator. This is Phase 1 — define what tasks exist and their metadata.
Requirements
Create src/backend/services/pipeline_tasks.py
Implement:
-
PipelineTaskdataclass:name: str— unique task identifierdescription: str— human-readable namefunc: Callable— async function to executedepends_on: list[str]— task names that must complete firsttimeout: int— max execution time in secondsretry_count: int— number of retries on failure
-
TaskRegistryclass:register(task: PipelineTask)— add task to registryget(task_name: str) -> PipelineTask— lookup taskget_all() -> list[PipelineTask]— list all registered tasksget_dependencies(task_name: str) -> list[str]— get task dependenciesget_ready_tasks( completed: set[str]) -> list[str]— find tasks whose deps are met
-
Register built-in tasks:
news_ingestion— runs news data ingestionfinancials_ingestion— runs financials data ingestionsector_rotation— runs sector rotation analysisprice_update— runs price data updatesentiment_analysis— runs sentiment analysis
Constraints
- File under 200 lines
- Use existing async patterns
- No real DB calls in task definitions
Acceptance Criteria
TaskRegistrycorrectly tracks tasks and dependenciesget_ready_tasks()returns correct task order- All 5 built-in tasks are registered
- File is under 200 lines
Files to Create
src/backend/services/pipeline_tasks.py
Files to Read First
src/backend/services/news_ingestion_service.py— existing task functionssrc/backend/services/financials_ingestion_service.py— existing task functionssrc/backend/services/rotation_service.py— existing task functions