20260731_11_daily_demand_package.py 8.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243
  1. """add explainable evaluations and the unified daily demand package
  2. Revision ID: 20260731_11
  3. Revises: 20260731_10
  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_11"
  11. down_revision: str | None = "20260731_10"
  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.create_table(
  34. "strategy_version",
  35. sa.Column("strategy_version_id", sa.String(36), primary_key=True),
  36. sa.Column("version_key", sa.String(64), nullable=False),
  37. sa.Column("status", sa.String(24), nullable=False),
  38. sa.Column("definition_json", sa.JSON(), nullable=False),
  39. sa.Column("change_reason", sa.Text(), nullable=False),
  40. sa.Column("created_by", sa.String(128), nullable=False),
  41. sa.Column("activated_at", sa.DateTime(), nullable=True),
  42. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  43. sa.UniqueConstraint("version_key", name="uk_strategy_version_key"),
  44. **table_options,
  45. )
  46. op.create_index(
  47. "idx_strategy_version_status",
  48. "strategy_version",
  49. ["status", "activated_at"],
  50. )
  51. op.create_table(
  52. "demand_evaluation_snapshot",
  53. sa.Column("evaluation_id", sa.String(36), primary_key=True),
  54. sa.Column(
  55. "platform_demand_version_id",
  56. sa.String(36),
  57. sa.ForeignKey(
  58. "platform_demand_version.platform_demand_version_id",
  59. ondelete="RESTRICT",
  60. ),
  61. nullable=False,
  62. ),
  63. sa.Column(
  64. "run_id",
  65. sa.String(36),
  66. sa.ForeignKey("pipeline_run.run_id", ondelete="RESTRICT"),
  67. nullable=False,
  68. ),
  69. sa.Column(
  70. "strategy_version_id",
  71. sa.String(36),
  72. sa.ForeignKey("strategy_version.strategy_version_id", ondelete="RESTRICT"),
  73. nullable=False,
  74. ),
  75. sa.Column("biz_dt", sa.String(8), nullable=False),
  76. sa.Column("validity_score", sa.Numeric(6, 5), nullable=False),
  77. sa.Column("local_supply_priority", sa.Numeric(6, 5), nullable=False),
  78. sa.Column("data_confidence", sa.Numeric(6, 5), nullable=False),
  79. sa.Column("posterior_state", sa.String(24), nullable=False),
  80. sa.Column("action_tier", sa.String(32), nullable=False),
  81. sa.Column("metrics_json", sa.JSON(), nullable=False),
  82. sa.Column("missing_dimensions_json", sa.JSON(), nullable=False),
  83. sa.Column("decision_reason", sa.Text(), nullable=False),
  84. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  85. sa.UniqueConstraint(
  86. "platform_demand_version_id",
  87. name="uk_demand_evaluation_version",
  88. ),
  89. **table_options,
  90. )
  91. op.create_index(
  92. "idx_demand_evaluation_biz",
  93. "demand_evaluation_snapshot",
  94. ["biz_dt", "action_tier"],
  95. )
  96. op.create_table(
  97. "demand_score_contribution",
  98. sa.Column("contribution_id", sa.String(36), primary_key=True),
  99. sa.Column(
  100. "evaluation_id",
  101. sa.String(36),
  102. sa.ForeignKey("demand_evaluation_snapshot.evaluation_id", ondelete="RESTRICT"),
  103. nullable=False,
  104. ),
  105. sa.Column("score_type", sa.String(24), nullable=False),
  106. sa.Column("dimension", sa.String(64), nullable=False),
  107. sa.Column("raw_value_json", sa.JSON(), nullable=False),
  108. sa.Column("normalized_value", sa.Numeric(8, 7), nullable=False),
  109. sa.Column("configured_weight", sa.Numeric(8, 7), nullable=False),
  110. sa.Column("effective_weight", sa.Numeric(8, 7), nullable=False),
  111. sa.Column("contribution", sa.Numeric(8, 7), nullable=False),
  112. sa.Column("data_status", sa.String(24), nullable=False),
  113. sa.Column("reason", sa.Text(), nullable=False),
  114. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  115. sa.UniqueConstraint(
  116. "evaluation_id",
  117. "score_type",
  118. "dimension",
  119. name="uk_demand_score_contribution",
  120. ),
  121. **table_options,
  122. )
  123. op.create_index(
  124. "idx_demand_score_contribution_eval",
  125. "demand_score_contribution",
  126. ["evaluation_id"],
  127. )
  128. op.create_table(
  129. "daily_demand_package",
  130. sa.Column("demand_package_id", sa.String(36), primary_key=True),
  131. sa.Column(
  132. "run_id",
  133. sa.String(36),
  134. sa.ForeignKey("pipeline_run.run_id", ondelete="RESTRICT"),
  135. nullable=False,
  136. ),
  137. sa.Column(
  138. "evidence_package_id",
  139. sa.String(36),
  140. sa.ForeignKey("evidence_package.package_id", ondelete="RESTRICT"),
  141. nullable=False,
  142. ),
  143. sa.Column(
  144. "strategy_version_id",
  145. sa.String(36),
  146. sa.ForeignKey("strategy_version.strategy_version_id", ondelete="RESTRICT"),
  147. nullable=False,
  148. ),
  149. sa.Column("biz_dt", sa.String(8), nullable=False),
  150. sa.Column("package_version", sa.Integer(), nullable=False),
  151. sa.Column("status", sa.String(24), nullable=False),
  152. sa.Column("content_hash", sa.String(64), nullable=False),
  153. sa.Column("item_count", sa.Integer(), nullable=False),
  154. sa.Column("published_at", sa.DateTime(), nullable=True),
  155. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  156. sa.UniqueConstraint("run_id", name="uk_daily_demand_package_run"),
  157. **table_options,
  158. )
  159. op.create_index(
  160. "idx_daily_demand_package_biz",
  161. "daily_demand_package",
  162. ["biz_dt", "status"],
  163. )
  164. op.create_table(
  165. "daily_demand_task",
  166. sa.Column("task_id", sa.String(36), primary_key=True),
  167. sa.Column(
  168. "demand_package_id",
  169. sa.String(36),
  170. sa.ForeignKey("daily_demand_package.demand_package_id", ondelete="RESTRICT"),
  171. nullable=False,
  172. ),
  173. sa.Column(
  174. "platform_demand_version_id",
  175. sa.String(36),
  176. sa.ForeignKey(
  177. "platform_demand_version.platform_demand_version_id",
  178. ondelete="RESTRICT",
  179. ),
  180. nullable=False,
  181. ),
  182. sa.Column(
  183. "evaluation_id",
  184. sa.String(36),
  185. sa.ForeignKey("demand_evaluation_snapshot.evaluation_id", ondelete="RESTRICT"),
  186. nullable=False,
  187. ),
  188. sa.Column("action_tier", sa.String(32), nullable=False),
  189. sa.Column("search_terms_json", sa.JSON(), nullable=False),
  190. sa.Column("exclude_terms_json", sa.JSON(), nullable=False),
  191. sa.Column("hit_rules_json", sa.JSON(), nullable=False),
  192. sa.Column("hypotheses_json", sa.JSON(), nullable=False),
  193. sa.Column("evidence_gaps_json", sa.JSON(), nullable=False),
  194. sa.Column("existing_content_json", sa.JSON(), nullable=False),
  195. sa.Column("allocation_json", sa.JSON(), nullable=False),
  196. sa.Column("task_payload_json", sa.JSON(), nullable=False),
  197. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  198. sa.UniqueConstraint(
  199. "demand_package_id",
  200. "platform_demand_version_id",
  201. name="uk_daily_demand_task_version",
  202. ),
  203. **table_options,
  204. )
  205. op.create_index(
  206. "idx_daily_demand_task_tier",
  207. "daily_demand_task",
  208. ["demand_package_id", "action_tier"],
  209. )
  210. def downgrade() -> None:
  211. op.drop_index("idx_daily_demand_task_tier", table_name="daily_demand_task")
  212. op.drop_table("daily_demand_task")
  213. op.drop_index(
  214. "idx_daily_demand_package_biz",
  215. table_name="daily_demand_package",
  216. )
  217. op.drop_table("daily_demand_package")
  218. op.drop_index(
  219. "idx_demand_score_contribution_eval",
  220. table_name="demand_score_contribution",
  221. )
  222. op.drop_table("demand_score_contribution")
  223. op.drop_index(
  224. "idx_demand_evaluation_biz",
  225. table_name="demand_evaluation_snapshot",
  226. )
  227. op.drop_table("demand_evaluation_snapshot")
  228. op.drop_index("idx_strategy_version_status", table_name="strategy_version")
  229. op.drop_table("strategy_version")