| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296 |
- """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")
|