pipeline_outbox.py 2.0 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465
  1. from __future__ import annotations
  2. from datetime import datetime
  3. from typing import Any
  4. from sqlalchemy import (
  5. Boolean,
  6. DateTime,
  7. ForeignKey,
  8. Index,
  9. JSON,
  10. String,
  11. Text,
  12. UniqueConstraint,
  13. func,
  14. )
  15. from sqlalchemy.orm import Mapped, mapped_column
  16. from supply_infra.db.base import Base
  17. class PipelineOutbox(Base):
  18. """Ledger of dispatched AIGC write effects."""
  19. __tablename__ = "pipeline_outbox"
  20. __table_args__ = (
  21. UniqueConstraint("idempotency_key", name="uk_pipeline_outbox_idempotency"),
  22. Index("idx_pipeline_outbox_run", "run_id", "step_run_id"),
  23. Index("idx_pipeline_outbox_status", "status", "created_at"),
  24. )
  25. outbox_id: Mapped[str] = mapped_column(String(36), primary_key=True)
  26. run_id: Mapped[str] = mapped_column(
  27. String(36),
  28. ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
  29. nullable=False,
  30. )
  31. step_run_id: Mapped[str] = mapped_column(
  32. String(36),
  33. ForeignKey("pipeline_step_run.step_run_id", ondelete="CASCADE"),
  34. nullable=False,
  35. )
  36. effect_type: Mapped[str] = mapped_column(String(64), nullable=False)
  37. idempotency_key: Mapped[str] = mapped_column(String(191), nullable=False)
  38. payload_hash: Mapped[str] = mapped_column(String(64), nullable=False)
  39. payload_json: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True)
  40. payload_uri: Mapped[str | None] = mapped_column(String(1024), nullable=True)
  41. status: Mapped[str] = mapped_column(
  42. String(32),
  43. nullable=False,
  44. default="dispatched",
  45. )
  46. dry_run: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
  47. record_error: Mapped[str | None] = mapped_column(Text, nullable=True)
  48. created_at: Mapped[datetime] = mapped_column(
  49. DateTime,
  50. nullable=False,
  51. server_default=func.now(),
  52. )
  53. updated_at: Mapped[datetime] = mapped_column(
  54. DateTime,
  55. nullable=False,
  56. server_default=func.now(),
  57. onupdate=func.now(),
  58. )