pipeline_lock.py 1.1 KB

1234567891011121314151617181920212223242526272829303132
  1. from __future__ import annotations
  2. from datetime import datetime
  3. from sqlalchemy import BigInteger, DateTime, String, func
  4. from sqlalchemy.orm import Mapped, mapped_column
  5. from supply_infra.db.base import Base
  6. class PipelineLock(Base):
  7. """Auditable lease used for cross-process pipeline coordination."""
  8. __tablename__ = "pipeline_lock"
  9. lock_key: Mapped[str] = mapped_column(String(191), primary_key=True)
  10. owner_run_id: Mapped[str] = mapped_column(String(36), nullable=False)
  11. owner_instance: Mapped[str] = mapped_column(String(191), nullable=False)
  12. lease_until: Mapped[datetime] = mapped_column(DateTime, nullable=False)
  13. heartbeat_at: Mapped[datetime] = mapped_column(DateTime, nullable=False)
  14. version: Mapped[int] = mapped_column(BigInteger, nullable=False, default=1)
  15. created_at: Mapped[datetime] = mapped_column(
  16. DateTime,
  17. nullable=False,
  18. server_default=func.now(),
  19. )
  20. updated_at: Mapped[datetime] = mapped_column(
  21. DateTime,
  22. nullable=False,
  23. server_default=func.now(),
  24. onupdate=func.now(),
  25. )