models.py 2.6 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586
  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_table: str | None = None
  40. target_id: UUID | None = None
  41. attempt_count: int = 0
  42. metadata: dict[str, Any] = Field(default_factory=dict)
  43. error_message: str | None = None
  44. started_at: datetime | None = None
  45. finished_at: datetime | None = None
  46. class PipelineRunEvent(PipelineModel):
  47. id: UUID | None = None
  48. pipeline_run_id: UUID | None = None
  49. pipeline_job_id: UUID | None = None
  50. stage: PipelineStage | str
  51. event_type: str
  52. status: str | None = None
  53. severity: str = "info"
  54. target_table: str | None = None
  55. target_id: UUID | None = None
  56. acquisition_run_id: UUID | None = None
  57. acquisition_job_id: UUID | None = None
  58. query_id: UUID | None = None
  59. item_id: UUID | None = None
  60. decode_job_id: UUID | None = None
  61. decode_result_id: UUID | None = None
  62. payload_draft_id: UUID | None = None
  63. ingest_record_id: UUID | None = None
  64. platform: str | None = None
  65. attempt_index: int | None = None
  66. message: str | None = None
  67. payload: dict[str, Any] = Field(default_factory=dict)
  68. error_message: str | None = None
  69. duration_ms: int | None = None
  70. created_at: datetime | None = None