20260727_01_pipeline_control_plane.py 8.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. """create durable pipeline control plane
  2. Revision ID: 20260727_01
  3. Revises:
  4. Create Date: 2026-07-27
  5. """
  6. from __future__ import annotations
  7. from collections.abc import Sequence
  8. from datetime import datetime
  9. import sqlalchemy as sa
  10. from alembic import op
  11. revision: str = "20260727_01"
  12. down_revision: str | None = None
  13. branch_labels: str | Sequence[str] | None = None
  14. depends_on: str | Sequence[str] | None = None
  15. def upgrade() -> None:
  16. op.create_table(
  17. "pipeline_run",
  18. sa.Column("run_id", sa.String(36), primary_key=True),
  19. sa.Column("dedupe_key", sa.String(191), nullable=False),
  20. sa.Column("pipeline_key", sa.String(64), nullable=False),
  21. sa.Column("biz_dt", sa.String(8), nullable=False),
  22. sa.Column("trigger_type", sa.String(24), nullable=False),
  23. sa.Column("trigger_source", sa.String(64), nullable=True),
  24. sa.Column("triggered_by", sa.String(128), nullable=True),
  25. sa.Column("trigger_reason", sa.Text(), nullable=True),
  26. sa.Column("parent_run_id", sa.String(36), nullable=True),
  27. sa.Column("run_mode", sa.String(24), nullable=False),
  28. sa.Column("dry_run", sa.Boolean(), nullable=False),
  29. sa.Column("status", sa.String(32), nullable=False),
  30. sa.Column("current_step", sa.String(64), nullable=True),
  31. sa.Column("scheduled_for", sa.DateTime(), nullable=True),
  32. sa.Column("deadline_at", sa.DateTime(), nullable=False),
  33. sa.Column("started_at", sa.DateTime(), nullable=True),
  34. sa.Column("finished_at", sa.DateTime(), nullable=True),
  35. sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
  36. sa.Column("lease_owner", sa.String(191), nullable=True),
  37. sa.Column("lease_until", sa.DateTime(), nullable=True),
  38. sa.Column("code_version", sa.String(128), nullable=True),
  39. sa.Column("config_snapshot_json", sa.JSON(), nullable=False),
  40. sa.Column("date_snapshot_json", sa.JSON(), nullable=False),
  41. sa.Column("summary_json", sa.JSON(), nullable=True),
  42. sa.Column("error_code", sa.String(64), nullable=True),
  43. sa.Column("error_message", sa.Text(), nullable=True),
  44. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  45. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  46. sa.UniqueConstraint("dedupe_key", name="uk_pipeline_run_dedupe_key"),
  47. mysql_charset="utf8mb4",
  48. )
  49. op.create_index(
  50. "idx_pipeline_run_biz",
  51. "pipeline_run",
  52. ["pipeline_key", "biz_dt", "created_at"],
  53. )
  54. op.create_index(
  55. "idx_pipeline_run_status",
  56. "pipeline_run",
  57. ["status", "lease_until"],
  58. )
  59. op.create_index(
  60. "idx_pipeline_run_deadline",
  61. "pipeline_run",
  62. ["status", "deadline_at"],
  63. )
  64. op.create_table(
  65. "pipeline_step_run",
  66. sa.Column("step_run_id", sa.String(36), primary_key=True),
  67. sa.Column(
  68. "run_id",
  69. sa.String(36),
  70. sa.ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
  71. nullable=False,
  72. ),
  73. sa.Column("step_key", sa.String(64), nullable=False),
  74. sa.Column("step_order", sa.Integer(), nullable=False),
  75. sa.Column("attempt", sa.Integer(), nullable=False),
  76. sa.Column("status", sa.String(32), nullable=False),
  77. sa.Column("critical", sa.Boolean(), nullable=False),
  78. sa.Column("dependency_snapshot_json", sa.JSON(), nullable=False),
  79. sa.Column("input_snapshot_json", sa.JSON(), nullable=False),
  80. sa.Column("timeout_seconds", sa.Integer(), nullable=False),
  81. sa.Column("max_attempts", sa.Integer(), nullable=False),
  82. sa.Column("retryable", sa.Boolean(), nullable=False),
  83. sa.Column("next_retry_at", sa.DateTime(), nullable=True),
  84. sa.Column("started_at", sa.DateTime(), nullable=True),
  85. sa.Column("finished_at", sa.DateTime(), nullable=True),
  86. sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
  87. sa.Column("lease_owner", sa.String(191), nullable=True),
  88. sa.Column("lease_until", sa.DateTime(), nullable=True),
  89. sa.Column("exit_code", sa.Integer(), nullable=True),
  90. sa.Column("result_summary_json", sa.JSON(), nullable=True),
  91. sa.Column("metrics_json", sa.JSON(), nullable=True),
  92. sa.Column("error_code", sa.String(64), nullable=True),
  93. sa.Column("error_message", sa.Text(), nullable=True),
  94. sa.Column("log_uri", sa.String(1024), nullable=True),
  95. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  96. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  97. sa.UniqueConstraint(
  98. "run_id",
  99. "step_key",
  100. "attempt",
  101. name="uk_pipeline_step_run_attempt",
  102. ),
  103. mysql_charset="utf8mb4",
  104. )
  105. op.create_index(
  106. "idx_pipeline_step_claim",
  107. "pipeline_step_run",
  108. ["status", "next_retry_at", "step_order"],
  109. )
  110. op.create_index(
  111. "idx_pipeline_step_run",
  112. "pipeline_step_run",
  113. ["run_id", "step_order", "attempt"],
  114. )
  115. op.create_index(
  116. "idx_pipeline_step_lease",
  117. "pipeline_step_run",
  118. ["status", "lease_until"],
  119. )
  120. op.create_table(
  121. "pipeline_lock",
  122. sa.Column("lock_key", sa.String(191), primary_key=True),
  123. sa.Column("owner_run_id", sa.String(36), nullable=False),
  124. sa.Column("owner_instance", sa.String(191), nullable=False),
  125. sa.Column("lease_until", sa.DateTime(), nullable=False),
  126. sa.Column("heartbeat_at", sa.DateTime(), nullable=False),
  127. sa.Column("version", sa.BigInteger(), nullable=False),
  128. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  129. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  130. mysql_charset="utf8mb4",
  131. )
  132. pipeline_lock = sa.table(
  133. "pipeline_lock",
  134. sa.column("lock_key", sa.String()),
  135. sa.column("owner_run_id", sa.String()),
  136. sa.column("owner_instance", sa.String()),
  137. sa.column("lease_until", sa.DateTime()),
  138. sa.column("heartbeat_at", sa.DateTime()),
  139. sa.column("version", sa.BigInteger()),
  140. )
  141. epoch = datetime(1970, 1, 1)
  142. op.bulk_insert(
  143. pipeline_lock,
  144. [
  145. {
  146. "lock_key": "pipeline:claim_guard",
  147. "owner_run_id": "control-plane",
  148. "owner_instance": "claim-coordinator",
  149. "lease_until": epoch,
  150. "heartbeat_at": epoch,
  151. "version": 1,
  152. }
  153. ],
  154. )
  155. op.create_table(
  156. "pipeline_outbox",
  157. sa.Column("outbox_id", sa.String(36), primary_key=True),
  158. sa.Column(
  159. "run_id",
  160. sa.String(36),
  161. sa.ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
  162. nullable=False,
  163. ),
  164. sa.Column(
  165. "step_run_id",
  166. sa.String(36),
  167. sa.ForeignKey("pipeline_step_run.step_run_id", ondelete="CASCADE"),
  168. nullable=False,
  169. ),
  170. sa.Column("effect_type", sa.String(64), nullable=False),
  171. sa.Column("idempotency_key", sa.String(191), nullable=False),
  172. sa.Column("payload_hash", sa.String(64), nullable=False),
  173. sa.Column("payload_json", sa.JSON(), nullable=True),
  174. sa.Column("payload_uri", sa.String(1024), nullable=True),
  175. sa.Column("status", sa.String(32), nullable=False),
  176. sa.Column("dry_run", sa.Boolean(), nullable=False),
  177. sa.Column("record_error", sa.Text(), nullable=True),
  178. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  179. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  180. sa.UniqueConstraint(
  181. "idempotency_key",
  182. name="uk_pipeline_outbox_idempotency",
  183. ),
  184. mysql_charset="utf8mb4",
  185. )
  186. op.create_index(
  187. "idx_pipeline_outbox_run",
  188. "pipeline_outbox",
  189. ["run_id", "step_run_id"],
  190. )
  191. op.create_index(
  192. "idx_pipeline_outbox_status",
  193. "pipeline_outbox",
  194. ["status", "created_at"],
  195. )
  196. def downgrade() -> None:
  197. op.drop_index("idx_pipeline_outbox_status", table_name="pipeline_outbox")
  198. op.drop_index("idx_pipeline_outbox_run", table_name="pipeline_outbox")
  199. op.drop_table("pipeline_outbox")
  200. op.drop_table("pipeline_lock")
  201. op.drop_index("idx_pipeline_step_lease", table_name="pipeline_step_run")
  202. op.drop_index("idx_pipeline_step_run", table_name="pipeline_step_run")
  203. op.drop_index("idx_pipeline_step_claim", table_name="pipeline_step_run")
  204. op.drop_table("pipeline_step_run")
  205. op.drop_index("idx_pipeline_run_deadline", table_name="pipeline_run")
  206. op.drop_index("idx_pipeline_run_status", table_name="pipeline_run")
  207. op.drop_index("idx_pipeline_run_biz", table_name="pipeline_run")
  208. op.drop_table("pipeline_run")