"""add content attribution, feedback processing and strategy workflow Revision ID: 20260731_12 Revises: 20260731_11 Create Date: 2026-07-31 """ from __future__ import annotations from collections.abc import Sequence import sqlalchemy as sa from alembic import op revision: str = "20260731_12" down_revision: str | None = "20260731_11" branch_labels: str | Sequence[str] | None = None depends_on: str | Sequence[str] | None = None def _pipeline_table_options() -> dict[str, str]: bind = op.get_bind() if bind.dialect.name != "mysql": return {} row = bind.execute( sa.text( "SELECT CHARACTER_SET_NAME, COLLATION_NAME " "FROM INFORMATION_SCHEMA.COLUMNS " "WHERE TABLE_SCHEMA = DATABASE() " "AND TABLE_NAME = 'pipeline_run' " "AND COLUMN_NAME = 'run_id'" ) ).one() return { "mysql_charset": str(row.CHARACTER_SET_NAME), "mysql_collate": str(row.COLLATION_NAME), } def upgrade() -> None: table_options = _pipeline_table_options() op.add_column( "video_discovery_run", sa.Column("demand_package_id", sa.String(36), nullable=True), ) op.add_column( "video_discovery_run", sa.Column("daily_demand_task_id", sa.String(36), nullable=True), ) op.add_column( "video_discovery_run", sa.Column("platform_demand_version_id", sa.String(36), nullable=True), ) op.create_index( "idx_video_discovery_run_platform_version", "video_discovery_run", ["platform_demand_version_id"], ) op.add_column( "demand_feedback", sa.Column( "processing_status", sa.String(24), nullable=False, server_default="pending", ), ) op.add_column( "demand_feedback", sa.Column("platform_demand_id", sa.String(36), nullable=True), ) op.add_column( "demand_feedback", sa.Column("platform_demand_version_id", sa.String(36), nullable=True), ) op.add_column( "demand_feedback", sa.Column("processed_by", sa.String(128), nullable=True), ) op.add_column( "demand_feedback", sa.Column("processed_at", sa.DateTime(), nullable=True), ) op.add_column( "demand_feedback", sa.Column("resolution_reason", sa.Text(), nullable=True), ) op.add_column( "demand_feedback", sa.Column("consumed_run_id", sa.String(36), nullable=True), ) op.add_column( "demand_feedback", sa.Column("impact_json", sa.Text(), nullable=True), ) op.create_index( "idx_demand_feedback_processing", "demand_feedback", ["processing_status", "created_at"], ) op.create_table( "content_demand_coverage", sa.Column("coverage_id", sa.String(36), primary_key=True), sa.Column( "platform_demand_version_id", sa.String(36), sa.ForeignKey( "platform_demand_version.platform_demand_version_id", ondelete="RESTRICT", ), nullable=False, ), sa.Column( "task_id", sa.String(36), sa.ForeignKey("daily_demand_task.task_id", ondelete="RESTRICT"), nullable=False, ), sa.Column("content_id", sa.String(128), nullable=False), sa.Column("content_url", sa.String(1024), nullable=True), sa.Column("source_type", sa.String(32), nullable=False), sa.Column("coverage_role", sa.String(24), nullable=False), sa.Column("coverage_score", sa.Numeric(6, 5), nullable=False), sa.Column("evidence_json", sa.JSON(), nullable=False), sa.Column("status", sa.String(24), nullable=False), sa.Column("discovered_at", sa.DateTime(), nullable=True), sa.Column("content_published_at", sa.DateTime(), nullable=True), sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False), sa.UniqueConstraint( "platform_demand_version_id", "content_id", "coverage_role", name="uk_content_demand_coverage", ), **table_options, ) op.create_index( "idx_content_demand_coverage_content", "content_demand_coverage", ["content_id"], ) op.create_table( "content_performance_fact", sa.Column("performance_fact_id", sa.String(36), primary_key=True), sa.Column("content_id", sa.String(128), nullable=False), sa.Column("biz_dt", sa.String(8), nullable=False), sa.Column("source", sa.String(64), nullable=False), sa.Column("payload_hash", sa.String(64), nullable=False), sa.Column("exposure_count", sa.BigInteger(), nullable=True), sa.Column("sample_size", sa.BigInteger(), nullable=True), sa.Column("content_quality", sa.Numeric(6, 5), nullable=True), sa.Column("rov", sa.Numeric(16, 8), nullable=True), sa.Column("vov", sa.Numeric(16, 8), nullable=True), sa.Column("raw_payload_json", sa.JSON(), nullable=False), sa.Column("observed_at", sa.DateTime(), nullable=False), sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False), sa.UniqueConstraint( "content_id", "biz_dt", "source", "payload_hash", name="uk_content_performance_fact", ), **table_options, ) op.create_index( "idx_content_performance_fact_content", "content_performance_fact", ["content_id", "biz_dt"], ) op.create_table( "demand_attribution_snapshot", sa.Column("attribution_id", sa.String(36), primary_key=True), sa.Column( "platform_demand_version_id", sa.String(36), sa.ForeignKey( "platform_demand_version.platform_demand_version_id", ondelete="RESTRICT", ), nullable=False, ), sa.Column( "coverage_id", sa.String(36), sa.ForeignKey("content_demand_coverage.coverage_id", ondelete="RESTRICT"), nullable=False, ), sa.Column( "performance_fact_id", sa.String(36), sa.ForeignKey( "content_performance_fact.performance_fact_id", ondelete="RESTRICT", ), nullable=False, ), sa.Column( "strategy_version_id", sa.String(36), sa.ForeignKey("strategy_version.strategy_version_id", ondelete="RESTRICT"), nullable=False, ), sa.Column("support_state", sa.String(24), nullable=False), sa.Column("attributable_rov", sa.Numeric(16, 8), nullable=True), sa.Column("attributable_vov", sa.Numeric(16, 8), nullable=True), sa.Column("attribution_confidence", sa.Numeric(6, 5), nullable=False), sa.Column("validity_influence", sa.Numeric(8, 7), nullable=False), sa.Column("priority_influence", sa.Numeric(8, 7), nullable=False), sa.Column("reason", sa.Text(), nullable=False), sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False), sa.UniqueConstraint( "coverage_id", "performance_fact_id", "strategy_version_id", name="uk_demand_attribution_snapshot", ), **table_options, ) op.create_index( "idx_demand_attribution_version", "demand_attribution_snapshot", ["platform_demand_version_id", "created_at"], ) op.create_table( "strategy_change_proposal", sa.Column("proposal_id", sa.String(36), primary_key=True), sa.Column( "base_strategy_version_id", sa.String(36), sa.ForeignKey("strategy_version.strategy_version_id", ondelete="RESTRICT"), nullable=False, ), sa.Column("status", sa.String(24), nullable=False), sa.Column("requested_change_json", sa.JSON(), nullable=False), sa.Column("proposed_definition_json", sa.JSON(), nullable=False), sa.Column("diff_json", sa.JSON(), nullable=False), sa.Column("preview_json", sa.JSON(), nullable=True), sa.Column("reason", sa.Text(), nullable=False), sa.Column("created_by", sa.String(128), nullable=False), sa.Column("approved_by", sa.String(128), nullable=True), sa.Column("approved_at", sa.DateTime(), nullable=True), sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False), **table_options, ) op.create_index( "idx_strategy_proposal_status", "strategy_change_proposal", ["status", "created_at"], ) def downgrade() -> None: op.drop_index( "idx_strategy_proposal_status", table_name="strategy_change_proposal", ) op.drop_table("strategy_change_proposal") op.drop_index( "idx_demand_attribution_version", table_name="demand_attribution_snapshot", ) op.drop_table("demand_attribution_snapshot") op.drop_index( "idx_content_performance_fact_content", table_name="content_performance_fact", ) op.drop_table("content_performance_fact") op.drop_index( "idx_content_demand_coverage_content", table_name="content_demand_coverage", ) op.drop_table("content_demand_coverage") op.drop_index("idx_demand_feedback_processing", table_name="demand_feedback") for column in ( "impact_json", "consumed_run_id", "resolution_reason", "processed_at", "processed_by", "platform_demand_version_id", "platform_demand_id", "processing_status", ): op.drop_column("demand_feedback", column) op.drop_index( "idx_video_discovery_run_platform_version", table_name="video_discovery_run", ) op.drop_column("video_discovery_run", "platform_demand_version_id") op.drop_column("video_discovery_run", "daily_demand_task_id") op.drop_column("video_discovery_run", "demand_package_id")