20260729_03_schema_baseline.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268
  1. """consolidated schema baseline
  2. Revision ID: 20260729_03
  3. Revises:
  4. Create Date: 2026-07-29
  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. from sqlalchemy.dialects import mysql
  12. revision: str = "20260729_03"
  13. down_revision: str | None = None
  14. branch_labels: str | Sequence[str] | None = None
  15. depends_on: str | Sequence[str] | None = None
  16. def upgrade() -> None:
  17. op.create_table(
  18. "pipeline_run",
  19. sa.Column("run_id", sa.String(36), primary_key=True),
  20. sa.Column("dedupe_key", sa.String(191), nullable=False),
  21. sa.Column("pipeline_key", sa.String(64), nullable=False),
  22. sa.Column("biz_dt", sa.String(8), nullable=False),
  23. sa.Column("trigger_type", sa.String(24), nullable=False),
  24. sa.Column("trigger_source", sa.String(64), nullable=True),
  25. sa.Column("trigger_reason", sa.Text(), nullable=True),
  26. sa.Column("run_mode", sa.String(24), nullable=False),
  27. sa.Column("dry_run", sa.Boolean(), nullable=False),
  28. sa.Column("status", sa.String(32), nullable=False),
  29. sa.Column("current_step", sa.String(64), nullable=True),
  30. sa.Column("deadline_at", sa.DateTime(), nullable=True),
  31. sa.Column("started_at", sa.DateTime(), nullable=True),
  32. sa.Column("finished_at", sa.DateTime(), nullable=True),
  33. sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
  34. sa.Column("lease_owner", sa.String(191), nullable=True),
  35. sa.Column("lease_until", sa.DateTime(), nullable=True),
  36. sa.Column("code_version", sa.String(128), nullable=True),
  37. sa.Column("config_snapshot_json", sa.JSON(), nullable=False),
  38. sa.Column("date_snapshot_json", sa.JSON(), nullable=False),
  39. sa.Column("summary_json", sa.JSON(), nullable=True),
  40. sa.Column("error_code", sa.String(64), nullable=True),
  41. sa.Column("error_message", sa.Text(), nullable=True),
  42. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  43. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  44. sa.UniqueConstraint("dedupe_key", name="uk_pipeline_run_dedupe_key"),
  45. mysql_charset="utf8mb4",
  46. )
  47. op.create_index(
  48. "idx_pipeline_run_biz",
  49. "pipeline_run",
  50. ["pipeline_key", "biz_dt", "created_at"],
  51. )
  52. op.create_index(
  53. "idx_pipeline_run_status",
  54. "pipeline_run",
  55. ["status", "lease_until"],
  56. )
  57. op.create_index(
  58. "idx_pipeline_run_deadline",
  59. "pipeline_run",
  60. ["status", "deadline_at"],
  61. )
  62. op.create_table(
  63. "pipeline_step_run",
  64. sa.Column("step_run_id", sa.String(36), primary_key=True),
  65. sa.Column(
  66. "run_id",
  67. sa.String(36),
  68. sa.ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
  69. nullable=False,
  70. ),
  71. sa.Column("step_key", sa.String(64), nullable=False),
  72. sa.Column("step_order", sa.Integer(), nullable=False),
  73. sa.Column("attempt", sa.Integer(), nullable=False),
  74. sa.Column("status", sa.String(32), nullable=False),
  75. sa.Column("critical", sa.Boolean(), nullable=False),
  76. sa.Column("dependency_snapshot_json", sa.JSON(), nullable=False),
  77. sa.Column("input_snapshot_json", sa.JSON(), nullable=False),
  78. sa.Column("timeout_seconds", sa.Integer(), nullable=False),
  79. sa.Column("max_attempts", sa.Integer(), nullable=False),
  80. sa.Column("retryable", sa.Boolean(), nullable=False),
  81. sa.Column("next_retry_at", sa.DateTime(), nullable=True),
  82. sa.Column("started_at", sa.DateTime(), nullable=True),
  83. sa.Column("finished_at", sa.DateTime(), nullable=True),
  84. sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
  85. sa.Column("lease_owner", sa.String(191), nullable=True),
  86. sa.Column("lease_until", sa.DateTime(), nullable=True),
  87. sa.Column("exit_code", sa.Integer(), nullable=True),
  88. sa.Column("result_summary_json", sa.JSON(), nullable=True),
  89. sa.Column("error_code", sa.String(64), nullable=True),
  90. sa.Column("error_message", sa.Text(), nullable=True),
  91. sa.Column("log_uri", sa.String(1024), nullable=True),
  92. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  93. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  94. sa.UniqueConstraint(
  95. "run_id",
  96. "step_key",
  97. "attempt",
  98. name="uk_pipeline_step_run_attempt",
  99. ),
  100. mysql_charset="utf8mb4",
  101. )
  102. op.create_index(
  103. "idx_pipeline_step_claim",
  104. "pipeline_step_run",
  105. ["status", "next_retry_at", "step_order"],
  106. )
  107. op.create_index(
  108. "idx_pipeline_step_run",
  109. "pipeline_step_run",
  110. ["run_id", "step_order", "attempt"],
  111. )
  112. op.create_index(
  113. "idx_pipeline_step_lease",
  114. "pipeline_step_run",
  115. ["status", "lease_until"],
  116. )
  117. op.create_table(
  118. "pipeline_lock",
  119. sa.Column("lock_key", sa.String(191), primary_key=True),
  120. sa.Column("owner_run_id", sa.String(36), nullable=False),
  121. sa.Column("owner_instance", sa.String(191), nullable=False),
  122. sa.Column("lease_until", sa.DateTime(), nullable=False),
  123. sa.Column("heartbeat_at", sa.DateTime(), nullable=False),
  124. sa.Column("version", sa.BigInteger(), nullable=False),
  125. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  126. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  127. mysql_charset="utf8mb4",
  128. )
  129. pipeline_lock = sa.table(
  130. "pipeline_lock",
  131. sa.column("lock_key", sa.String()),
  132. sa.column("owner_run_id", sa.String()),
  133. sa.column("owner_instance", sa.String()),
  134. sa.column("lease_until", sa.DateTime()),
  135. sa.column("heartbeat_at", sa.DateTime()),
  136. sa.column("version", sa.BigInteger()),
  137. )
  138. epoch = datetime(1970, 1, 1)
  139. op.bulk_insert(
  140. pipeline_lock,
  141. [
  142. {
  143. "lock_key": "pipeline:claim_guard",
  144. "owner_run_id": "control-plane",
  145. "owner_instance": "claim-coordinator",
  146. "lease_until": epoch,
  147. "heartbeat_at": epoch,
  148. "version": 1,
  149. }
  150. ],
  151. )
  152. op.create_table(
  153. "pipeline_outbox",
  154. sa.Column("outbox_id", sa.String(36), primary_key=True),
  155. sa.Column(
  156. "run_id",
  157. sa.String(36),
  158. sa.ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
  159. nullable=False,
  160. ),
  161. sa.Column(
  162. "step_run_id",
  163. sa.String(36),
  164. sa.ForeignKey("pipeline_step_run.step_run_id", ondelete="CASCADE"),
  165. nullable=False,
  166. ),
  167. sa.Column("effect_type", sa.String(64), nullable=False),
  168. sa.Column("idempotency_key", sa.String(191), nullable=False),
  169. sa.Column("payload_hash", sa.String(64), nullable=False),
  170. sa.Column("payload_json", sa.JSON(), nullable=True),
  171. sa.Column("payload_uri", sa.String(1024), nullable=True),
  172. sa.Column("status", sa.String(32), nullable=False),
  173. sa.Column("dry_run", sa.Boolean(), nullable=False),
  174. sa.Column("record_error", sa.Text(), nullable=True),
  175. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  176. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  177. sa.UniqueConstraint(
  178. "idempotency_key",
  179. name="uk_pipeline_outbox_idempotency",
  180. ),
  181. mysql_charset="utf8mb4",
  182. )
  183. op.create_index(
  184. "idx_pipeline_outbox_run",
  185. "pipeline_outbox",
  186. ["run_id", "step_run_id"],
  187. )
  188. op.create_index(
  189. "idx_pipeline_outbox_status",
  190. "pipeline_outbox",
  191. ["status", "created_at"],
  192. )
  193. op.create_table(
  194. "agent_document_injection",
  195. sa.Column("id", sa.BigInteger(), autoincrement=True, nullable=False),
  196. sa.Column(
  197. "agent_name",
  198. sa.String(length=128),
  199. nullable=False,
  200. comment="Agent name",
  201. ),
  202. sa.Column(
  203. "content",
  204. sa.Text().with_variant(mysql.LONGTEXT(), "mysql"),
  205. nullable=False,
  206. comment="Document content injected after the base system prompt",
  207. ),
  208. sa.Column(
  209. "enabled",
  210. sa.Integer(),
  211. server_default=sa.text("1"),
  212. nullable=False,
  213. comment="Whether the document injection is active",
  214. ),
  215. sa.Column(
  216. "create_time",
  217. sa.DateTime(),
  218. server_default=sa.func.now(),
  219. nullable=False,
  220. ),
  221. sa.Column(
  222. "update_time",
  223. sa.DateTime(),
  224. server_default=sa.func.now(),
  225. nullable=False,
  226. ),
  227. sa.PrimaryKeyConstraint("id"),
  228. sa.UniqueConstraint(
  229. "agent_name",
  230. name="uk_agent_document_injection_agent_name",
  231. ),
  232. mysql_charset="utf8mb4",
  233. )
  234. op.create_index(
  235. "ix_agent_document_injection_agent_name",
  236. "agent_document_injection",
  237. ["agent_name"],
  238. )
  239. def downgrade() -> None:
  240. op.drop_index(
  241. "ix_agent_document_injection_agent_name",
  242. table_name="agent_document_injection",
  243. )
  244. op.drop_table("agent_document_injection")
  245. op.drop_index("idx_pipeline_outbox_status", table_name="pipeline_outbox")
  246. op.drop_index("idx_pipeline_outbox_run", table_name="pipeline_outbox")
  247. op.drop_table("pipeline_outbox")
  248. op.drop_table("pipeline_lock")
  249. op.drop_index("idx_pipeline_step_lease", table_name="pipeline_step_run")
  250. op.drop_index("idx_pipeline_step_run", table_name="pipeline_step_run")
  251. op.drop_index("idx_pipeline_step_claim", table_name="pipeline_step_run")
  252. op.drop_table("pipeline_step_run")
  253. op.drop_index("idx_pipeline_run_deadline", table_name="pipeline_run")
  254. op.drop_index("idx_pipeline_run_status", table_name="pipeline_run")
  255. op.drop_index("idx_pipeline_run_biz", table_name="pipeline_run")
  256. op.drop_table("pipeline_run")