models.py 1.7 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758
  1. """Pipeline run and resume models for formal orchestration."""
  2. from __future__ import annotations
  3. from datetime import datetime
  4. from typing import Any, Literal
  5. from uuid import UUID
  6. from pydantic import BaseModel, ConfigDict, Field
  7. PipelineStage = Literal["query", "search", "classify", "decode", "payload", "ingest"]
  8. PipelineStatus = Literal["pending", "running", "done", "failed", "skipped", "partial"]
  9. PipelineStep = PipelineStage
  10. RunStatus = PipelineStatus
  11. JobStatus = PipelineStatus
  12. class PipelineModel(BaseModel):
  13. model_config = ConfigDict(extra="forbid")
  14. class ResumePolicy(PipelineModel):
  15. resume: bool = True
  16. skip_done: bool = True
  17. retry_failed: bool = False
  18. max_attempts: int = 3
  19. class ResumeCursor(PipelineModel):
  20. stage: PipelineStage
  21. target_id: UUID | None = None
  22. last_successful_job_id: UUID | None = None
  23. metadata: dict[str, Any] = Field(default_factory=dict)
  24. class PipelineRun(PipelineModel):
  25. id: UUID | None = None
  26. run_key: str | None = None
  27. batch_id: UUID | None = None
  28. status: PipelineStatus = "pending"
  29. current_stage: PipelineStage | None = None
  30. metadata: dict[str, Any] = Field(default_factory=dict)
  31. error_message: str | None = None
  32. started_at: datetime | None = None
  33. finished_at: datetime | None = None
  34. class PipelineJob(PipelineModel):
  35. id: UUID | None = None
  36. run_id: UUID | None = None
  37. stage: PipelineStage
  38. status: PipelineStatus = "pending"
  39. target_id: UUID | None = None
  40. attempt_count: int = 0
  41. metadata: dict[str, Any] = Field(default_factory=dict)
  42. error_message: str | None = None
  43. started_at: datetime | None = None
  44. finished_at: datetime | None = None