|
@@ -0,0 +1,219 @@
|
|
|
|
|
+"""create durable pipeline control plane
|
|
|
|
|
+
|
|
|
|
|
+Revision ID: 20260727_01
|
|
|
|
|
+Revises:
|
|
|
|
|
+Create Date: 2026-07-27
|
|
|
|
|
+"""
|
|
|
|
|
+from __future__ import annotations
|
|
|
|
|
+
|
|
|
|
|
+from collections.abc import Sequence
|
|
|
|
|
+from datetime import datetime
|
|
|
|
|
+
|
|
|
|
|
+import sqlalchemy as sa
|
|
|
|
|
+from alembic import op
|
|
|
|
|
+
|
|
|
|
|
+revision: str = "20260727_01"
|
|
|
|
|
+down_revision: str | None = None
|
|
|
|
|
+branch_labels: str | Sequence[str] | None = None
|
|
|
|
|
+depends_on: str | Sequence[str] | None = None
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+def upgrade() -> None:
|
|
|
|
|
+ op.create_table(
|
|
|
|
|
+ "pipeline_run",
|
|
|
|
|
+ sa.Column("run_id", sa.String(36), primary_key=True),
|
|
|
|
|
+ sa.Column("dedupe_key", sa.String(191), nullable=False),
|
|
|
|
|
+ sa.Column("pipeline_key", sa.String(64), nullable=False),
|
|
|
|
|
+ sa.Column("biz_dt", sa.String(8), nullable=False),
|
|
|
|
|
+ sa.Column("trigger_type", sa.String(24), nullable=False),
|
|
|
|
|
+ sa.Column("trigger_source", sa.String(64), nullable=True),
|
|
|
|
|
+ sa.Column("triggered_by", sa.String(128), nullable=True),
|
|
|
|
|
+ sa.Column("trigger_reason", sa.Text(), nullable=True),
|
|
|
|
|
+ sa.Column("parent_run_id", sa.String(36), nullable=True),
|
|
|
|
|
+ sa.Column("run_mode", sa.String(24), nullable=False),
|
|
|
|
|
+ sa.Column("dry_run", sa.Boolean(), nullable=False),
|
|
|
|
|
+ sa.Column("status", sa.String(32), nullable=False),
|
|
|
|
|
+ sa.Column("current_step", sa.String(64), nullable=True),
|
|
|
|
|
+ sa.Column("scheduled_for", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("deadline_at", sa.DateTime(), nullable=False),
|
|
|
|
|
+ sa.Column("started_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("finished_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("lease_owner", sa.String(191), nullable=True),
|
|
|
|
|
+ sa.Column("lease_until", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("code_version", sa.String(128), nullable=True),
|
|
|
|
|
+ sa.Column("config_snapshot_json", sa.JSON(), nullable=False),
|
|
|
|
|
+ sa.Column("date_snapshot_json", sa.JSON(), nullable=False),
|
|
|
|
|
+ sa.Column("summary_json", sa.JSON(), nullable=True),
|
|
|
|
|
+ sa.Column("error_code", sa.String(64), nullable=True),
|
|
|
|
|
+ sa.Column("error_message", sa.Text(), nullable=True),
|
|
|
|
|
+ sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.UniqueConstraint("dedupe_key", name="uk_pipeline_run_dedupe_key"),
|
|
|
|
|
+ mysql_charset="utf8mb4",
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_run_biz",
|
|
|
|
|
+ "pipeline_run",
|
|
|
|
|
+ ["pipeline_key", "biz_dt", "created_at"],
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_run_status",
|
|
|
|
|
+ "pipeline_run",
|
|
|
|
|
+ ["status", "lease_until"],
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_run_deadline",
|
|
|
|
|
+ "pipeline_run",
|
|
|
|
|
+ ["status", "deadline_at"],
|
|
|
|
|
+ )
|
|
|
|
|
+
|
|
|
|
|
+ op.create_table(
|
|
|
|
|
+ "pipeline_step_run",
|
|
|
|
|
+ sa.Column("step_run_id", sa.String(36), primary_key=True),
|
|
|
|
|
+ sa.Column(
|
|
|
|
|
+ "run_id",
|
|
|
|
|
+ sa.String(36),
|
|
|
|
|
+ sa.ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
|
|
|
|
|
+ nullable=False,
|
|
|
|
|
+ ),
|
|
|
|
|
+ sa.Column("step_key", sa.String(64), nullable=False),
|
|
|
|
|
+ sa.Column("step_order", sa.Integer(), nullable=False),
|
|
|
|
|
+ sa.Column("attempt", sa.Integer(), nullable=False),
|
|
|
|
|
+ sa.Column("status", sa.String(32), nullable=False),
|
|
|
|
|
+ sa.Column("critical", sa.Boolean(), nullable=False),
|
|
|
|
|
+ sa.Column("dependency_snapshot_json", sa.JSON(), nullable=False),
|
|
|
|
|
+ sa.Column("input_snapshot_json", sa.JSON(), nullable=False),
|
|
|
|
|
+ sa.Column("timeout_seconds", sa.Integer(), nullable=False),
|
|
|
|
|
+ sa.Column("max_attempts", sa.Integer(), nullable=False),
|
|
|
|
|
+ sa.Column("retryable", sa.Boolean(), nullable=False),
|
|
|
|
|
+ sa.Column("next_retry_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("started_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("finished_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("lease_owner", sa.String(191), nullable=True),
|
|
|
|
|
+ sa.Column("lease_until", sa.DateTime(), nullable=True),
|
|
|
|
|
+ sa.Column("exit_code", sa.Integer(), nullable=True),
|
|
|
|
|
+ sa.Column("result_summary_json", sa.JSON(), nullable=True),
|
|
|
|
|
+ sa.Column("metrics_json", sa.JSON(), nullable=True),
|
|
|
|
|
+ sa.Column("error_code", sa.String(64), nullable=True),
|
|
|
|
|
+ sa.Column("error_message", sa.Text(), nullable=True),
|
|
|
|
|
+ sa.Column("log_uri", sa.String(1024), nullable=True),
|
|
|
|
|
+ sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.UniqueConstraint(
|
|
|
|
|
+ "run_id",
|
|
|
|
|
+ "step_key",
|
|
|
|
|
+ "attempt",
|
|
|
|
|
+ name="uk_pipeline_step_run_attempt",
|
|
|
|
|
+ ),
|
|
|
|
|
+ mysql_charset="utf8mb4",
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_step_claim",
|
|
|
|
|
+ "pipeline_step_run",
|
|
|
|
|
+ ["status", "next_retry_at", "step_order"],
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_step_run",
|
|
|
|
|
+ "pipeline_step_run",
|
|
|
|
|
+ ["run_id", "step_order", "attempt"],
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_step_lease",
|
|
|
|
|
+ "pipeline_step_run",
|
|
|
|
|
+ ["status", "lease_until"],
|
|
|
|
|
+ )
|
|
|
|
|
+
|
|
|
|
|
+ op.create_table(
|
|
|
|
|
+ "pipeline_lock",
|
|
|
|
|
+ sa.Column("lock_key", sa.String(191), primary_key=True),
|
|
|
|
|
+ sa.Column("owner_run_id", sa.String(36), nullable=False),
|
|
|
|
|
+ sa.Column("owner_instance", sa.String(191), nullable=False),
|
|
|
|
|
+ sa.Column("lease_until", sa.DateTime(), nullable=False),
|
|
|
|
|
+ sa.Column("heartbeat_at", sa.DateTime(), nullable=False),
|
|
|
|
|
+ sa.Column("version", sa.BigInteger(), nullable=False),
|
|
|
|
|
+ sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ mysql_charset="utf8mb4",
|
|
|
|
|
+ )
|
|
|
|
|
+ pipeline_lock = sa.table(
|
|
|
|
|
+ "pipeline_lock",
|
|
|
|
|
+ sa.column("lock_key", sa.String()),
|
|
|
|
|
+ sa.column("owner_run_id", sa.String()),
|
|
|
|
|
+ sa.column("owner_instance", sa.String()),
|
|
|
|
|
+ sa.column("lease_until", sa.DateTime()),
|
|
|
|
|
+ sa.column("heartbeat_at", sa.DateTime()),
|
|
|
|
|
+ sa.column("version", sa.BigInteger()),
|
|
|
|
|
+ )
|
|
|
|
|
+ epoch = datetime(1970, 1, 1)
|
|
|
|
|
+ op.bulk_insert(
|
|
|
|
|
+ pipeline_lock,
|
|
|
|
|
+ [
|
|
|
|
|
+ {
|
|
|
|
|
+ "lock_key": "pipeline:claim_guard",
|
|
|
|
|
+ "owner_run_id": "control-plane",
|
|
|
|
|
+ "owner_instance": "claim-coordinator",
|
|
|
|
|
+ "lease_until": epoch,
|
|
|
|
|
+ "heartbeat_at": epoch,
|
|
|
|
|
+ "version": 1,
|
|
|
|
|
+ }
|
|
|
|
|
+ ],
|
|
|
|
|
+ )
|
|
|
|
|
+
|
|
|
|
|
+ op.create_table(
|
|
|
|
|
+ "pipeline_outbox",
|
|
|
|
|
+ sa.Column("outbox_id", sa.String(36), primary_key=True),
|
|
|
|
|
+ sa.Column(
|
|
|
|
|
+ "run_id",
|
|
|
|
|
+ sa.String(36),
|
|
|
|
|
+ sa.ForeignKey("pipeline_run.run_id", ondelete="CASCADE"),
|
|
|
|
|
+ nullable=False,
|
|
|
|
|
+ ),
|
|
|
|
|
+ sa.Column(
|
|
|
|
|
+ "step_run_id",
|
|
|
|
|
+ sa.String(36),
|
|
|
|
|
+ sa.ForeignKey("pipeline_step_run.step_run_id", ondelete="CASCADE"),
|
|
|
|
|
+ nullable=False,
|
|
|
|
|
+ ),
|
|
|
|
|
+ sa.Column("effect_type", sa.String(64), nullable=False),
|
|
|
|
|
+ sa.Column("idempotency_key", sa.String(191), nullable=False),
|
|
|
|
|
+ sa.Column("payload_hash", sa.String(64), nullable=False),
|
|
|
|
|
+ sa.Column("payload_json", sa.JSON(), nullable=True),
|
|
|
|
|
+ sa.Column("payload_uri", sa.String(1024), nullable=True),
|
|
|
|
|
+ sa.Column("status", sa.String(32), nullable=False),
|
|
|
|
|
+ sa.Column("dry_run", sa.Boolean(), nullable=False),
|
|
|
|
|
+ sa.Column("record_error", sa.Text(), nullable=True),
|
|
|
|
|
+ sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
|
|
|
|
+ sa.UniqueConstraint(
|
|
|
|
|
+ "idempotency_key",
|
|
|
|
|
+ name="uk_pipeline_outbox_idempotency",
|
|
|
|
|
+ ),
|
|
|
|
|
+ mysql_charset="utf8mb4",
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_outbox_run",
|
|
|
|
|
+ "pipeline_outbox",
|
|
|
|
|
+ ["run_id", "step_run_id"],
|
|
|
|
|
+ )
|
|
|
|
|
+ op.create_index(
|
|
|
|
|
+ "idx_pipeline_outbox_status",
|
|
|
|
|
+ "pipeline_outbox",
|
|
|
|
|
+ ["status", "created_at"],
|
|
|
|
|
+ )
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+def downgrade() -> None:
|
|
|
|
|
+ op.drop_index("idx_pipeline_outbox_status", table_name="pipeline_outbox")
|
|
|
|
|
+ op.drop_index("idx_pipeline_outbox_run", table_name="pipeline_outbox")
|
|
|
|
|
+ op.drop_table("pipeline_outbox")
|
|
|
|
|
+ op.drop_table("pipeline_lock")
|
|
|
|
|
+ op.drop_index("idx_pipeline_step_lease", table_name="pipeline_step_run")
|
|
|
|
|
+ op.drop_index("idx_pipeline_step_run", table_name="pipeline_step_run")
|
|
|
|
|
+ op.drop_index("idx_pipeline_step_claim", table_name="pipeline_step_run")
|
|
|
|
|
+ op.drop_table("pipeline_step_run")
|
|
|
|
|
+ op.drop_index("idx_pipeline_run_deadline", table_name="pipeline_run")
|
|
|
|
|
+ op.drop_index("idx_pipeline_run_status", table_name="pipeline_run")
|
|
|
|
|
+ op.drop_index("idx_pipeline_run_biz", table_name="pipeline_run")
|
|
|
|
|
+ op.drop_table("pipeline_run")
|