pipeline_run.py 2.9 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768
  1. from __future__ import annotations
  2. from datetime import datetime
  3. from typing import Any
  4. from sqlalchemy import Boolean, DateTime, Index, JSON, String, Text, UniqueConstraint, func
  5. from sqlalchemy.orm import Mapped, mapped_column
  6. from supply_infra.db.base import Base
  7. class PipelineRun(Base):
  8. """One durable execution request for the complete supply pipeline."""
  9. __tablename__ = "pipeline_run"
  10. __table_args__ = (
  11. UniqueConstraint("dedupe_key", name="uk_pipeline_run_dedupe_key"),
  12. Index("idx_pipeline_run_biz", "pipeline_key", "biz_dt", "created_at"),
  13. Index("idx_pipeline_run_status", "status", "lease_until"),
  14. Index("idx_pipeline_run_deadline", "status", "deadline_at"),
  15. )
  16. run_id: Mapped[str] = mapped_column(String(36), primary_key=True)
  17. dedupe_key: Mapped[str] = mapped_column(String(191), nullable=False)
  18. pipeline_key: Mapped[str] = mapped_column(
  19. String(64),
  20. nullable=False,
  21. default="supply_pipeline",
  22. )
  23. biz_dt: Mapped[str] = mapped_column(String(8), nullable=False)
  24. trigger_type: Mapped[str] = mapped_column(String(24), nullable=False)
  25. trigger_source: Mapped[str | None] = mapped_column(String(64), nullable=True)
  26. trigger_reason: Mapped[str | None] = mapped_column(Text, nullable=True)
  27. run_mode: Mapped[str] = mapped_column(String(24), nullable=False, default="full")
  28. dry_run: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
  29. status: Mapped[str] = mapped_column(String(32), nullable=False, default="queued")
  30. current_step: Mapped[str | None] = mapped_column(String(64), nullable=True)
  31. deadline_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
  32. started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
  33. finished_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
  34. heartbeat_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
  35. lease_owner: Mapped[str | None] = mapped_column(String(191), nullable=True)
  36. lease_until: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
  37. code_version: Mapped[str | None] = mapped_column(String(128), nullable=True)
  38. config_snapshot_json: Mapped[dict[str, Any]] = mapped_column(
  39. JSON,
  40. nullable=False,
  41. default=dict,
  42. )
  43. date_snapshot_json: Mapped[dict[str, Any]] = mapped_column(
  44. JSON,
  45. nullable=False,
  46. default=dict,
  47. )
  48. summary_json: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True)
  49. error_code: Mapped[str | None] = mapped_column(String(64), nullable=True)
  50. error_message: Mapped[str | None] = mapped_column(Text, nullable=True)
  51. created_at: Mapped[datetime] = mapped_column(
  52. DateTime,
  53. nullable=False,
  54. server_default=func.now(),
  55. )
  56. updated_at: Mapped[datetime] = mapped_column(
  57. DateTime,
  58. nullable=False,
  59. server_default=func.now(),
  60. onupdate=func.now(),
  61. )