"""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_table: str | None = None 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 class PipelineRunEvent(PipelineModel): id: UUID | None = None pipeline_run_id: UUID | None = None pipeline_job_id: UUID | None = None stage: PipelineStage | str event_type: str status: str | None = None severity: str = "info" target_table: str | None = None target_id: UUID | None = None acquisition_run_id: UUID | None = None acquisition_job_id: UUID | None = None query_id: UUID | None = None item_id: UUID | None = None decode_job_id: UUID | None = None decode_result_id: UUID | None = None payload_draft_id: UUID | None = None ingest_record_id: UUID | None = None platform: str | None = None attempt_index: int | None = None message: str | None = None payload: dict[str, Any] = Field(default_factory=dict) error_message: str | None = None duration_ms: int | None = None created_at: datetime | None = None