20260731_10_platform_demand_mvp.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288
  1. """add evidence packages and versioned platform demand objects
  2. Revision ID: 20260731_10
  3. Revises: 20260731_09
  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_10"
  11. down_revision: str | None = "20260731_09"
  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. "evidence_package",
  35. sa.Column("package_id", sa.String(36), primary_key=True),
  36. sa.Column(
  37. "run_id",
  38. sa.String(36),
  39. sa.ForeignKey("pipeline_run.run_id", ondelete="RESTRICT"),
  40. nullable=False,
  41. ),
  42. sa.Column("biz_dt", sa.String(8), nullable=False),
  43. sa.Column("package_version", sa.Integer(), nullable=False),
  44. sa.Column("status", sa.String(24), nullable=False),
  45. sa.Column("source_snapshot_hash", sa.String(64), nullable=False),
  46. sa.Column("source_versions_json", sa.JSON(), nullable=False),
  47. sa.Column("evidence_count", sa.Integer(), nullable=False),
  48. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  49. sa.UniqueConstraint("run_id", name="uk_evidence_package_run"),
  50. **table_options,
  51. )
  52. op.create_index(
  53. "idx_evidence_package_biz",
  54. "evidence_package",
  55. ["biz_dt", "created_at"],
  56. )
  57. op.create_table(
  58. "raw_demand_expression",
  59. sa.Column("expression_id", sa.String(36), primary_key=True),
  60. sa.Column(
  61. "package_id",
  62. sa.String(36),
  63. sa.ForeignKey("evidence_package.package_id", ondelete="RESTRICT"),
  64. nullable=False,
  65. ),
  66. sa.Column("source_type", sa.String(64), nullable=False),
  67. sa.Column("source_record_id", sa.String(128), nullable=False),
  68. sa.Column("raw_text", sa.String(512), nullable=False),
  69. sa.Column("original_payload_json", sa.JSON(), nullable=False),
  70. sa.Column("content_hash", sa.String(64), nullable=False),
  71. sa.Column("observed_biz_dt", sa.String(8), nullable=False),
  72. sa.Column("data_quality", sa.String(24), nullable=False),
  73. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  74. sa.UniqueConstraint(
  75. "package_id",
  76. "source_type",
  77. "source_record_id",
  78. name="uk_raw_expression_source",
  79. ),
  80. **table_options,
  81. )
  82. op.create_index(
  83. "idx_raw_expression_package",
  84. "raw_demand_expression",
  85. ["package_id"],
  86. )
  87. op.create_index(
  88. "idx_raw_expression_hash",
  89. "raw_demand_expression",
  90. ["content_hash"],
  91. )
  92. op.create_table(
  93. "standard_demand_term",
  94. sa.Column("term_id", sa.String(36), primary_key=True),
  95. sa.Column("canonical_key", sa.String(64), nullable=False),
  96. sa.Column("canonical_text", sa.String(256), nullable=False),
  97. sa.Column("normalized_text", sa.String(256), nullable=False),
  98. sa.Column("status", sa.String(24), nullable=False),
  99. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  100. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  101. sa.UniqueConstraint("canonical_key", name="uk_standard_demand_term_key"),
  102. **table_options,
  103. )
  104. op.create_index(
  105. "idx_standard_demand_term_text",
  106. "standard_demand_term",
  107. ["normalized_text"],
  108. )
  109. op.create_table(
  110. "standard_demand_alias",
  111. sa.Column("alias_id", sa.String(36), primary_key=True),
  112. sa.Column(
  113. "term_id",
  114. sa.String(36),
  115. sa.ForeignKey("standard_demand_term.term_id", ondelete="RESTRICT"),
  116. nullable=False,
  117. ),
  118. sa.Column(
  119. "expression_id",
  120. sa.String(36),
  121. sa.ForeignKey("raw_demand_expression.expression_id", ondelete="RESTRICT"),
  122. nullable=False,
  123. ),
  124. sa.Column("alias_text", sa.String(512), nullable=False),
  125. sa.Column("relation_type", sa.String(32), nullable=False),
  126. sa.Column("confidence", sa.Numeric(6, 5), nullable=False),
  127. sa.Column("reason", sa.Text(), nullable=False),
  128. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  129. sa.UniqueConstraint(
  130. "term_id",
  131. "expression_id",
  132. name="uk_standard_demand_alias_expression",
  133. ),
  134. **table_options,
  135. )
  136. op.create_index(
  137. "idx_standard_demand_alias_term",
  138. "standard_demand_alias",
  139. ["term_id"],
  140. )
  141. op.create_table(
  142. "platform_demand",
  143. sa.Column("platform_demand_id", sa.String(36), primary_key=True),
  144. sa.Column("canonical_key", sa.String(64), nullable=False),
  145. sa.Column("name", sa.String(256), nullable=False),
  146. sa.Column("description", sa.Text(), nullable=False),
  147. sa.Column("status", sa.String(24), nullable=False),
  148. sa.Column("lifecycle_state", sa.String(24), nullable=False),
  149. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  150. sa.Column("updated_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  151. sa.UniqueConstraint("canonical_key", name="uk_platform_demand_key"),
  152. **table_options,
  153. )
  154. op.create_index("idx_platform_demand_name", "platform_demand", ["name"])
  155. op.create_index("idx_platform_demand_status", "platform_demand", ["status"])
  156. op.create_table(
  157. "platform_demand_version",
  158. sa.Column("platform_demand_version_id", sa.String(36), primary_key=True),
  159. sa.Column(
  160. "platform_demand_id",
  161. sa.String(36),
  162. sa.ForeignKey("platform_demand.platform_demand_id", ondelete="RESTRICT"),
  163. nullable=False,
  164. ),
  165. sa.Column(
  166. "term_id",
  167. sa.String(36),
  168. sa.ForeignKey("standard_demand_term.term_id", ondelete="RESTRICT"),
  169. nullable=False,
  170. ),
  171. sa.Column(
  172. "package_id",
  173. sa.String(36),
  174. sa.ForeignKey("evidence_package.package_id", ondelete="RESTRICT"),
  175. nullable=False,
  176. ),
  177. sa.Column(
  178. "run_id",
  179. sa.String(36),
  180. sa.ForeignKey("pipeline_run.run_id", ondelete="RESTRICT"),
  181. nullable=False,
  182. ),
  183. sa.Column("source_demand_grade_id", sa.BigInteger(), nullable=False),
  184. sa.Column("version_no", sa.Integer(), nullable=False),
  185. sa.Column("biz_dt", sa.String(8), nullable=False),
  186. sa.Column("name", sa.String(256), nullable=False),
  187. sa.Column("description", sa.Text(), nullable=False),
  188. sa.Column("cognition_confidence", sa.Numeric(6, 5), nullable=False),
  189. sa.Column("reason", sa.Text(), nullable=False),
  190. sa.Column("change_type", sa.String(24), nullable=False),
  191. sa.Column("evidence_hash", sa.String(64), nullable=False),
  192. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  193. sa.UniqueConstraint(
  194. "platform_demand_id",
  195. "version_no",
  196. name="uk_platform_demand_version_no",
  197. ),
  198. sa.UniqueConstraint(
  199. "platform_demand_id",
  200. "run_id",
  201. name="uk_platform_demand_version_run",
  202. ),
  203. **table_options,
  204. )
  205. op.create_index(
  206. "idx_platform_demand_version_biz",
  207. "platform_demand_version",
  208. ["biz_dt", "run_id"],
  209. )
  210. op.create_table(
  211. "platform_demand_category_rel",
  212. sa.Column("rel_id", sa.String(36), primary_key=True),
  213. sa.Column(
  214. "platform_demand_version_id",
  215. sa.String(36),
  216. sa.ForeignKey(
  217. "platform_demand_version.platform_demand_version_id",
  218. ondelete="RESTRICT",
  219. ),
  220. nullable=False,
  221. ),
  222. sa.Column("category_id", sa.BigInteger(), nullable=False),
  223. sa.Column("relation_type", sa.String(32), nullable=False),
  224. sa.Column("relation_source", sa.String(64), nullable=False),
  225. sa.Column("reason", sa.Text(), nullable=False),
  226. sa.Column("confidence", sa.Numeric(6, 5), nullable=False),
  227. sa.Column("is_inferred", sa.Boolean(), nullable=False),
  228. sa.Column("status", sa.String(24), nullable=False),
  229. sa.Column("valid_from_biz_dt", sa.String(8), nullable=False),
  230. sa.Column("created_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
  231. sa.UniqueConstraint(
  232. "platform_demand_version_id",
  233. "category_id",
  234. "relation_type",
  235. name="uk_platform_demand_category_rel",
  236. ),
  237. **table_options,
  238. )
  239. op.create_index(
  240. "idx_platform_demand_category_node",
  241. "platform_demand_category_rel",
  242. ["category_id"],
  243. )
  244. def downgrade() -> None:
  245. op.drop_index(
  246. "idx_platform_demand_category_node",
  247. table_name="platform_demand_category_rel",
  248. )
  249. op.drop_table("platform_demand_category_rel")
  250. op.drop_index(
  251. "idx_platform_demand_version_biz",
  252. table_name="platform_demand_version",
  253. )
  254. op.drop_table("platform_demand_version")
  255. op.drop_index("idx_platform_demand_status", table_name="platform_demand")
  256. op.drop_index("idx_platform_demand_name", table_name="platform_demand")
  257. op.drop_table("platform_demand")
  258. op.drop_index(
  259. "idx_standard_demand_alias_term",
  260. table_name="standard_demand_alias",
  261. )
  262. op.drop_table("standard_demand_alias")
  263. op.drop_index(
  264. "idx_standard_demand_term_text",
  265. table_name="standard_demand_term",
  266. )
  267. op.drop_table("standard_demand_term")
  268. op.drop_index("idx_raw_expression_hash", table_name="raw_demand_expression")
  269. op.drop_index("idx_raw_expression_package", table_name="raw_demand_expression")
  270. op.drop_table("raw_demand_expression")
  271. op.drop_index("idx_evidence_package_biz", table_name="evidence_package")
  272. op.drop_table("evidence_package")