| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758 |
- """Pipeline run and resume models for formal orchestration."""
- from __future__ import annotations
- from datetime import datetime
- from typing import Any, Literal
- from uuid import UUID
- from pydantic import BaseModel, ConfigDict, Field
- PipelineStage = Literal["query", "search", "classify", "decode", "payload", "ingest"]
- PipelineStatus = Literal["pending", "running", "done", "failed", "skipped", "partial"]
- PipelineStep = PipelineStage
- RunStatus = PipelineStatus
- JobStatus = PipelineStatus
- class PipelineModel(BaseModel):
- model_config = ConfigDict(extra="forbid")
- class ResumePolicy(PipelineModel):
- resume: bool = True
- skip_done: bool = True
- retry_failed: bool = False
- max_attempts: int = 3
- class ResumeCursor(PipelineModel):
- stage: PipelineStage
- target_id: UUID | None = None
- last_successful_job_id: UUID | None = None
- metadata: dict[str, Any] = Field(default_factory=dict)
- class PipelineRun(PipelineModel):
- id: UUID | None = None
- run_key: str | None = None
- batch_id: UUID | None = None
- status: PipelineStatus = "pending"
- current_stage: PipelineStage | None = None
- metadata: dict[str, Any] = Field(default_factory=dict)
- error_message: str | None = None
- started_at: datetime | None = None
- finished_at: datetime | None = None
- class PipelineJob(PipelineModel):
- id: UUID | None = None
- run_id: UUID | None = None
- stage: PipelineStage
- status: PipelineStatus = "pending"
- target_id: UUID | None = None
- attempt_count: int = 0
- metadata: dict[str, Any] = Field(default_factory=dict)
- error_message: str | None = None
- started_at: datetime | None = None
- finished_at: datetime | None = None
|