"""add append-only real release telemetry and review labels Revision ID: 20260716_0018 Revises: 20260716_0017 Create Date: 2026-07-16 22:00:00 """ from collections.abc import Sequence import sqlalchemy as sa from alembic import op revision: str = "20260716_0018" down_revision: str | None = "20260716_0017" branch_labels: str | Sequence[str] | None = None depends_on: str | Sequence[str] | None = None _TELEMETRY_TABLES = ( "agent_asset_release_observations", "agent_asset_release_labels", ) def _require_postgresql() -> None: dialect_name = op.get_bind().dialect.name if dialect_name != "postgresql": raise RuntimeError( "20260716_0018 only supports PostgreSQL; " f"refusing to mutate {dialect_name} without transactional constraint DDL" ) def _require_empty_telemetry_domain_for_downgrade() -> None: bind = op.get_bind() counts = { table_name: int(bind.scalar(sa.text(f"SELECT COUNT(*) FROM {table_name}")) or 0) for table_name in _TELEMETRY_TABLES } if any(counts.values()): summary = ", ".join(f"{name}={count}" for name, count in counts.items()) raise RuntimeError( "cannot downgrade release telemetry: immutable observations or labels exist " f"({summary})" ) def upgrade() -> None: _require_postgresql() op.create_table( "agent_asset_release_observations", sa.Column("id", sa.String(length=36), nullable=False), sa.Column("tenant_id", sa.String(length=64), nullable=False), sa.Column("asset_id", sa.String(length=36), nullable=False), sa.Column("release_id", sa.String(length=64), nullable=False), sa.Column("stage", sa.String(length=16), nullable=False), sa.Column("version", sa.String(length=30), nullable=False), sa.Column("rule_code", sa.String(length=100), nullable=False), sa.Column("business_stage", sa.String(length=40), nullable=False), sa.Column( "source_kind", sa.String(length=32), server_default="expense_claim_risk", nullable=False, ), sa.Column("source_fingerprint", sa.String(length=64), nullable=False), sa.Column("candidate_hit", sa.Boolean(), nullable=False), sa.Column("baseline_hit", sa.Boolean(), nullable=True), sa.Column("runtime_status", sa.String(length=16), nullable=False), sa.Column( "failure_code", sa.String(length=40), server_default="none", nullable=False, ), sa.Column("idempotency_key", sa.String(length=80), nullable=False), sa.Column("payload_fingerprint", sa.String(length=64), nullable=False), sa.Column( "created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False, ), sa.CheckConstraint( "stage IN ('shadow', 'canary', 'active')", name="ck_agent_asset_release_observations_stage", ), sa.CheckConstraint( "runtime_status IN ('completed', 'failed')", name="ck_agent_asset_release_observations_runtime_status", ), sa.CheckConstraint( "source_kind IN ('expense_claim_risk')", name="ck_agent_asset_release_observations_source_kind", ), sa.PrimaryKeyConstraint("id"), sa.UniqueConstraint( "tenant_id", "id", name="uq_agent_asset_release_observations_tenant_id", ), sa.UniqueConstraint( "tenant_id", "idempotency_key", name="uq_agent_asset_release_observations_tenant_idempotency", ), sa.UniqueConstraint( "tenant_id", "id", "asset_id", "release_id", "stage", "version", name="uq_agent_asset_release_observations_release_identity", ), ) op.create_index( "ix_agent_asset_release_observations_release", "agent_asset_release_observations", ["tenant_id", "asset_id", "release_id", "stage", "version", "created_at"], ) op.create_index( "ix_agent_asset_release_observations_source", "agent_asset_release_observations", ["tenant_id", "source_fingerprint"], ) op.create_table( "agent_asset_release_labels", sa.Column("id", sa.String(length=36), nullable=False), sa.Column("tenant_id", sa.String(length=64), nullable=False), sa.Column("observation_id", sa.String(length=36), nullable=False), sa.Column("asset_id", sa.String(length=36), nullable=False), sa.Column("release_id", sa.String(length=64), nullable=False), sa.Column("stage", sa.String(length=16), nullable=False), sa.Column("version", sa.String(length=30), nullable=False), sa.Column("label", sa.String(length=24), nullable=False), sa.Column("verification_source", sa.String(length=32), nullable=False), sa.Column("source_event_fingerprint", sa.String(length=64), nullable=False), sa.Column("actor_fingerprint", sa.String(length=64), nullable=False), sa.Column("idempotency_key", sa.String(length=80), nullable=False), sa.Column("payload_fingerprint", sa.String(length=64), nullable=False), sa.Column( "created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False, ), sa.CheckConstraint( "label IN ('confirmed', 'false_positive')", name="ck_agent_asset_release_labels_label", ), sa.CheckConstraint( "verification_source IN ('typed_risk_disposition', 'release_review')", name="ck_agent_asset_release_labels_source", ), sa.ForeignKeyConstraint( ["tenant_id", "observation_id", "asset_id", "release_id", "stage", "version"], [ "agent_asset_release_observations.tenant_id", "agent_asset_release_observations.id", "agent_asset_release_observations.asset_id", "agent_asset_release_observations.release_id", "agent_asset_release_observations.stage", "agent_asset_release_observations.version", ], name="fk_agent_asset_release_labels_release_observation", ondelete="RESTRICT", ), sa.PrimaryKeyConstraint("id"), sa.UniqueConstraint( "tenant_id", "id", name="uq_agent_asset_release_labels_tenant_id", ), sa.UniqueConstraint( "tenant_id", "idempotency_key", name="uq_agent_asset_release_labels_tenant_idempotency", ), ) op.create_index( "ix_agent_asset_release_labels_observation_time", "agent_asset_release_labels", ["tenant_id", "observation_id", "created_at"], ) op.create_index( "ix_agent_asset_release_labels_release", "agent_asset_release_labels", ["tenant_id", "asset_id", "release_id", "stage", "version"], ) op.execute( """ CREATE FUNCTION reject_agent_asset_release_telemetry_mutation() RETURNS trigger AS $$ BEGIN RAISE EXCEPTION 'agent asset release telemetry is append-only'; END; $$ LANGUAGE plpgsql; """ ) for table_name in _TELEMETRY_TABLES: op.execute( f""" CREATE TRIGGER trg_{table_name}_append_only BEFORE UPDATE OR DELETE ON {table_name} FOR EACH ROW EXECUTE FUNCTION reject_agent_asset_release_telemetry_mutation(); """ ) def downgrade() -> None: _require_postgresql() _require_empty_telemetry_domain_for_downgrade() for table_name in reversed(_TELEMETRY_TABLES): op.execute(f"DROP TRIGGER IF EXISTS trg_{table_name}_append_only ON {table_name}") op.execute("DROP FUNCTION IF EXISTS reject_agent_asset_release_telemetry_mutation()") op.drop_index( "ix_agent_asset_release_labels_release", table_name="agent_asset_release_labels", ) op.drop_index( "ix_agent_asset_release_labels_observation_time", table_name="agent_asset_release_labels", ) op.drop_table("agent_asset_release_labels") op.drop_index( "ix_agent_asset_release_observations_source", table_name="agent_asset_release_observations", ) op.drop_index( "ix_agent_asset_release_observations_release", table_name="agent_asset_release_observations", ) op.drop_table("agent_asset_release_observations")