85 lines
2.8 KiB
Python
85 lines
2.8 KiB
Python
"""Durable WorkflowRun — resume semantics, step state, idempotency (G-RUN, G-CTX).
|
|
|
|
Extends workflow_instances with resume_at, resume_reason, step_state,
|
|
idempotency_key, and lock_owner for durable/resumable workflow execution.
|
|
|
|
Revision ID: 0125
|
|
Revises: 0124
|
|
"""
|
|
|
|
from alembic import op
|
|
import sqlalchemy as sa
|
|
from sqlalchemy.dialects.postgresql import JSONB, UUID as PGUUID
|
|
|
|
revision = "0125"
|
|
down_revision = "0124"
|
|
branch_labels = None
|
|
depends_on = None
|
|
|
|
|
|
def upgrade() -> None:
|
|
# ── Add durable/resumable columns to workflow_instances ──────────────
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("resume_at", sa.DateTime(timezone=True), nullable=True),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("resume_reason", sa.String(50), nullable=True),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("step_state", JSONB, nullable=False, server_default="{}"),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("idempotency_key", sa.String(255), nullable=True),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("lock_owner", sa.String(100), nullable=True),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("lock_expires_at", sa.DateTime(timezone=True), nullable=True),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("error_message", sa.Text, nullable=True),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("retry_count", sa.Integer, nullable=False, server_default="0"),
|
|
)
|
|
op.add_column(
|
|
"workflow_instances",
|
|
sa.Column("max_retries", sa.Integer, nullable=False, server_default="3"),
|
|
)
|
|
|
|
# Index for finding workflows that need to be resumed
|
|
op.create_index(
|
|
"ix_wf_instances_resume",
|
|
"workflow_instances",
|
|
["tenant_id", "status", "resume_at"],
|
|
)
|
|
# Index for idempotency key lookup
|
|
op.create_index(
|
|
"ix_wf_instances_idempotency",
|
|
"workflow_instances",
|
|
["tenant_id", "idempotency_key"],
|
|
)
|
|
|
|
|
|
def downgrade() -> None:
|
|
op.drop_index("ix_wf_instances_idempotency", table_name="workflow_instances")
|
|
op.drop_index("ix_wf_instances_resume", table_name="workflow_instances")
|
|
op.drop_column("workflow_instances", "max_retries")
|
|
op.drop_column("workflow_instances", "retry_count")
|
|
op.drop_column("workflow_instances", "error_message")
|
|
op.drop_column("workflow_instances", "lock_expires_at")
|
|
op.drop_column("workflow_instances", "lock_owner")
|
|
op.drop_column("workflow_instances", "idempotency_key")
|
|
op.drop_column("workflow_instances", "step_state")
|
|
op.drop_column("workflow_instances", "resume_reason")
|
|
op.drop_column("workflow_instances", "resume_at")
|