20260731_12_feedback_attribution_strategy.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296
  1. """add content attribution, feedback processing and strategy workflow
  2. Revision ID: 20260731_12
  3. Revises: 20260731_11
  4. Create Date: 2026-07-31
  5. """
  6. from __future__ import annotations
  7. from collections.abc import Sequence
  8. import sqlalchemy as sa
  9. from alembic import op
  10. revision: str = "20260731_12"
  11. down_revision: str | None = "20260731_11"
  12. branch_labels: str | Sequence[str] | None = None
  13. depends_on: str | Sequence[str] | None = None
  14. def _pipeline_table_options() -> dict[str, str]:
  15. bind = op.get_bind()
  16. if bind.dialect.name != "mysql":
  17. return {}
  18. row = bind.execute(
  19. sa.text(
  20. "SELECT CHARACTER_SET_NAME, COLLATION_NAME "
  21. "FROM INFORMATION_SCHEMA.COLUMNS "
  22. "WHERE TABLE_SCHEMA = DATABASE() "
  23. "AND TABLE_NAME = 'pipeline_run' "
  24. "AND COLUMN_NAME = 'run_id'"
  25. )
  26. ).one()
  27. return {
  28. "mysql_charset": str(row.CHARACTER_SET_NAME),
  29. "mysql_collate": str(row.COLLATION_NAME),
  30. }
  31. def upgrade() -> None:
  32. table_options = _pipeline_table_options()
  33. op.add_column(
  34. "video_discovery_run",
  35. sa.Column("demand_package_id", sa.String(36), nullable=True),
  36. )
  37. op.add_column(
  38. "video_discovery_run",
  39. sa.Column("daily_demand_task_id", sa.String(36), nullable=True),
  40. )
  41. op.add_column(
  42. "video_discovery_run",
  43. sa.Column("platform_demand_version_id", sa.String(36), nullable=True),
  44. )
  45. op.create_index(
  46. "idx_video_discovery_run_platform_version",
  47. "video_discovery_run",
  48. ["platform_demand_version_id"],
  49. )
  50. op.add_column(
  51. "demand_feedback",
  52. sa.Column(
  53. "processing_status",
  54. sa.String(24),
  55. nullable=False,
  56. server_default="pending",
  57. ),
  58. )
  59. op.add_column(
  60. "demand_feedback",
  61. sa.Column("platform_demand_id", sa.String(36), nullable=True),
  62. )
  63. op.add_column(
  64. "demand_feedback",
  65. sa.Column("platform_demand_version_id", sa.String(36), nullable=True),
  66. )
  67. op.add_column(
  68. "demand_feedback",
  69. sa.Column("processed_by", sa.String(128), nullable=True),
  70. )
  71. op.add_column(
  72. "demand_feedback",
  73. sa.Column("processed_at", sa.DateTime(), nullable=True),
  74. )
  75. op.add_column(
  76. "demand_feedback",
  77. sa.Column("resolution_reason", sa.Text(), nullable=True),
  78. )
  79. op.add_column(
  80. "demand_feedback",
  81. sa.Column("consumed_run_id", sa.String(36), nullable=True),
  82. )
  83. op.add_column(
  84. "demand_feedback",
  85. sa.Column("impact_json", sa.Text(), nullable=True),
  86. )
  87. op.create_index(
  88. "idx_demand_feedback_processing",
  89. "demand_feedback",
  90. ["processing_status", "created_at"],
  91. )
  92. op.create_table(
  93. "content_demand_coverage",
  94. sa.Column("coverage_id", sa.String(36), primary_key=True),
  95. sa.Column(
  96. "platform_demand_version_id",
  97. sa.String(36),
  98. sa.ForeignKey(
  99. "platform_demand_version.platform_demand_version_id",
  100. ondelete="RESTRICT",
  101. ),
  102. nullable=False,
  103. ),
  104. sa.Column(
  105. "task_id",
  106. sa.String(36),
  107. sa.ForeignKey("daily_demand_task.task_id", ondelete="RESTRICT"),
  108. nullable=False,
  109. ),
  110. sa.Column("content_id", sa.String(128), nullable=False),
  111. sa.Column("content_url", sa.String(1024), nullable=True),
  112. sa.Column("source_type", sa.String(32), nullable=False),
  113. sa.Column("coverage_role", sa.String(24), nullable=False),
  114. sa.Column("coverage_score", sa.Numeric(6, 5), nullable=False),
  115. sa.Column("evidence_json", sa.JSON(), nullable=False),
  116. sa.Column("status", sa.String(24), nullable=False),
  117. sa.Column("discovered_at", sa.DateTime(), nullable=True),
  118. sa.Column("content_published_at", sa.DateTime(), nullable=True),
  119. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  120. sa.UniqueConstraint(
  121. "platform_demand_version_id",
  122. "content_id",
  123. "coverage_role",
  124. name="uk_content_demand_coverage",
  125. ),
  126. **table_options,
  127. )
  128. op.create_index(
  129. "idx_content_demand_coverage_content",
  130. "content_demand_coverage",
  131. ["content_id"],
  132. )
  133. op.create_table(
  134. "content_performance_fact",
  135. sa.Column("performance_fact_id", sa.String(36), primary_key=True),
  136. sa.Column("content_id", sa.String(128), nullable=False),
  137. sa.Column("biz_dt", sa.String(8), nullable=False),
  138. sa.Column("source", sa.String(64), nullable=False),
  139. sa.Column("payload_hash", sa.String(64), nullable=False),
  140. sa.Column("exposure_count", sa.BigInteger(), nullable=True),
  141. sa.Column("sample_size", sa.BigInteger(), nullable=True),
  142. sa.Column("content_quality", sa.Numeric(6, 5), nullable=True),
  143. sa.Column("rov", sa.Numeric(16, 8), nullable=True),
  144. sa.Column("vov", sa.Numeric(16, 8), nullable=True),
  145. sa.Column("raw_payload_json", sa.JSON(), nullable=False),
  146. sa.Column("observed_at", sa.DateTime(), nullable=False),
  147. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  148. sa.UniqueConstraint(
  149. "content_id",
  150. "biz_dt",
  151. "source",
  152. "payload_hash",
  153. name="uk_content_performance_fact",
  154. ),
  155. **table_options,
  156. )
  157. op.create_index(
  158. "idx_content_performance_fact_content",
  159. "content_performance_fact",
  160. ["content_id", "biz_dt"],
  161. )
  162. op.create_table(
  163. "demand_attribution_snapshot",
  164. sa.Column("attribution_id", sa.String(36), primary_key=True),
  165. sa.Column(
  166. "platform_demand_version_id",
  167. sa.String(36),
  168. sa.ForeignKey(
  169. "platform_demand_version.platform_demand_version_id",
  170. ondelete="RESTRICT",
  171. ),
  172. nullable=False,
  173. ),
  174. sa.Column(
  175. "coverage_id",
  176. sa.String(36),
  177. sa.ForeignKey("content_demand_coverage.coverage_id", ondelete="RESTRICT"),
  178. nullable=False,
  179. ),
  180. sa.Column(
  181. "performance_fact_id",
  182. sa.String(36),
  183. sa.ForeignKey(
  184. "content_performance_fact.performance_fact_id",
  185. ondelete="RESTRICT",
  186. ),
  187. nullable=False,
  188. ),
  189. sa.Column(
  190. "strategy_version_id",
  191. sa.String(36),
  192. sa.ForeignKey("strategy_version.strategy_version_id", ondelete="RESTRICT"),
  193. nullable=False,
  194. ),
  195. sa.Column("support_state", sa.String(24), nullable=False),
  196. sa.Column("attributable_rov", sa.Numeric(16, 8), nullable=True),
  197. sa.Column("attributable_vov", sa.Numeric(16, 8), nullable=True),
  198. sa.Column("attribution_confidence", sa.Numeric(6, 5), nullable=False),
  199. sa.Column("validity_influence", sa.Numeric(8, 7), nullable=False),
  200. sa.Column("priority_influence", sa.Numeric(8, 7), nullable=False),
  201. sa.Column("reason", sa.Text(), nullable=False),
  202. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  203. sa.UniqueConstraint(
  204. "coverage_id",
  205. "performance_fact_id",
  206. "strategy_version_id",
  207. name="uk_demand_attribution_snapshot",
  208. ),
  209. **table_options,
  210. )
  211. op.create_index(
  212. "idx_demand_attribution_version",
  213. "demand_attribution_snapshot",
  214. ["platform_demand_version_id", "created_at"],
  215. )
  216. op.create_table(
  217. "strategy_change_proposal",
  218. sa.Column("proposal_id", sa.String(36), primary_key=True),
  219. sa.Column(
  220. "base_strategy_version_id",
  221. sa.String(36),
  222. sa.ForeignKey("strategy_version.strategy_version_id", ondelete="RESTRICT"),
  223. nullable=False,
  224. ),
  225. sa.Column("status", sa.String(24), nullable=False),
  226. sa.Column("requested_change_json", sa.JSON(), nullable=False),
  227. sa.Column("proposed_definition_json", sa.JSON(), nullable=False),
  228. sa.Column("diff_json", sa.JSON(), nullable=False),
  229. sa.Column("preview_json", sa.JSON(), nullable=True),
  230. sa.Column("reason", sa.Text(), nullable=False),
  231. sa.Column("created_by", sa.String(128), nullable=False),
  232. sa.Column("approved_by", sa.String(128), nullable=True),
  233. sa.Column("approved_at", sa.DateTime(), nullable=True),
  234. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  235. **table_options,
  236. )
  237. op.create_index(
  238. "idx_strategy_proposal_status",
  239. "strategy_change_proposal",
  240. ["status", "created_at"],
  241. )
  242. def downgrade() -> None:
  243. op.drop_index(
  244. "idx_strategy_proposal_status",
  245. table_name="strategy_change_proposal",
  246. )
  247. op.drop_table("strategy_change_proposal")
  248. op.drop_index(
  249. "idx_demand_attribution_version",
  250. table_name="demand_attribution_snapshot",
  251. )
  252. op.drop_table("demand_attribution_snapshot")
  253. op.drop_index(
  254. "idx_content_performance_fact_content",
  255. table_name="content_performance_fact",
  256. )
  257. op.drop_table("content_performance_fact")
  258. op.drop_index(
  259. "idx_content_demand_coverage_content",
  260. table_name="content_demand_coverage",
  261. )
  262. op.drop_table("content_demand_coverage")
  263. op.drop_index("idx_demand_feedback_processing", table_name="demand_feedback")
  264. for column in (
  265. "impact_json",
  266. "consumed_run_id",
  267. "resolution_reason",
  268. "processed_at",
  269. "processed_by",
  270. "platform_demand_version_id",
  271. "platform_demand_id",
  272. "processing_status",
  273. ):
  274. op.drop_column("demand_feedback", column)
  275. op.drop_index(
  276. "idx_video_discovery_run_platform_version",
  277. table_name="video_discovery_run",
  278. )
  279. op.drop_column("video_discovery_run", "platform_demand_version_id")
  280. op.drop_column("video_discovery_run", "daily_demand_task_id")
  281. op.drop_column("video_discovery_run", "demand_package_id")