test_database_runtime.py 41 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120
  1. import json
  2. import re
  3. from decimal import Decimal
  4. from pathlib import Path
  5. from content_agent.integrations.database_runtime import (
  6. ContentSupplyDbConfig,
  7. DatabaseRuntimeStore,
  8. )
  9. class FakeCursor:
  10. def __init__(self, connection):
  11. self.connection = connection
  12. self._one = None
  13. self._all = []
  14. def __enter__(self):
  15. return self
  16. def __exit__(self, *_args):
  17. return None
  18. def execute(self, sql, params=None):
  19. self.connection.statements.append((sql, list(params or [])))
  20. if sql.startswith("SELECT COUNT(*)"):
  21. self._one = {"cnt": 0}
  22. self._all = []
  23. elif sql.startswith("SELECT"):
  24. self._one = self.connection.select_one_result
  25. self._all = self.connection.select_all_result
  26. def fetchone(self):
  27. return self._one
  28. def fetchall(self):
  29. return self._all
  30. class FakeConnection:
  31. def __init__(self):
  32. self.statements = []
  33. self.commit_count = 0
  34. self.select_one_result = None
  35. self.select_all_result = []
  36. def __enter__(self):
  37. return self
  38. def __exit__(self, *_args):
  39. return None
  40. def cursor(self):
  41. return FakeCursor(self)
  42. def commit(self):
  43. self.commit_count += 1
  44. def test_content_supply_db_config_reads_env_file(tmp_path):
  45. env_file = tmp_path / ".env"
  46. env_file.write_text(
  47. "\n".join(
  48. [
  49. "CONTENT_SUPPLY_DB_HOST=127.0.0.1",
  50. "CONTENT_SUPPLY_DB_PORT=3307",
  51. "CONTENT_SUPPLY_DB_NAME=content-deconstruction-supply",
  52. "CONTENT_SUPPLY_DB_USER=content_rw",
  53. "CONTENT_SUPPLY_DB_" + "PASS" + "WORD=dummy_password",
  54. ]
  55. ),
  56. encoding="utf-8",
  57. )
  58. config = ContentSupplyDbConfig.from_env(env_file=env_file)
  59. assert config.host == "127.0.0.1"
  60. assert config.port == 3307
  61. assert config.database == "content-deconstruction-supply"
  62. assert config.user == "content_rw"
  63. def test_content_supply_db_config_requires_all_project_db_keys(tmp_path):
  64. env_file = tmp_path / ".env"
  65. env_file.write_text(
  66. "\n".join(
  67. [
  68. "CONTENT_SUPPLY_DB_HOST=127.0.0.1",
  69. "CONTENT_SUPPLY_DB_PORT=3307",
  70. "CONTENT_SUPPLY_DB_NAME=content-deconstruction-supply",
  71. "CONTENT_SUPPLY_DB_USER=content_rw",
  72. ]
  73. ),
  74. encoding="utf-8",
  75. )
  76. try:
  77. ContentSupplyDbConfig.from_env(env_file=env_file)
  78. except ValueError as exc:
  79. assert "CONTENT_SUPPLY_DB_PASSWORD" in str(exc)
  80. else:
  81. raise AssertionError("expected missing db env key to fail")
  82. def test_database_runtime_writes_source_context_with_db_schema_version():
  83. connection = FakeConnection()
  84. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  85. store.write_json(
  86. "run_001",
  87. "source_context.json",
  88. {
  89. "schema_version": "runtime_record.v1",
  90. "run_id": "run_001",
  91. "demand_content_id": "123",
  92. "ext_data": {
  93. "evidence_pack": {
  94. "pattern_source_system": "pg_pattern_v2",
  95. "source_kind": "pattern_itemset",
  96. "source_post_id": "post_001",
  97. "pattern_execution_id": 581,
  98. "mining_config_id": 2081,
  99. }
  100. },
  101. },
  102. )
  103. sql, params = connection.statements[-1]
  104. values = _insert_values(sql, params)
  105. assert "INSERT INTO `content_agent_source_contexts`" in sql
  106. assert values["schema_version"] == "content_agent.v1"
  107. assert values["run_id"] == "run_001"
  108. assert values["demand_content_id"] == 123
  109. assert json.loads(values["evidence_pack"])["source_post_id"] == "post_001"
  110. assert json.loads(values["source_context"])["schema_version"] == "runtime_record.v1"
  111. def test_database_runtime_derives_itemset_ids_from_seed_pack_itemsets():
  112. connection = FakeConnection()
  113. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  114. store.write_json(
  115. "run_001",
  116. "pattern_seed_pack.json",
  117. {
  118. "schema_version": "runtime_record.v1",
  119. "run_id": "run_001",
  120. "policy_run_id": "policy_run_001",
  121. "source_post_id": "post_001",
  122. "pattern_execution_id": 581,
  123. "itemsets": [{"itemset_id": 1608352}, {"itemset_id": 1608352}],
  124. "seed_terms": ["爱国情感", "人物故事"],
  125. },
  126. )
  127. sql, params = connection.statements[-1]
  128. values = _insert_values(sql, params)
  129. assert "INSERT INTO `content_agent_pattern_seed_packs`" in sql
  130. assert values["schema_version"] == "content_agent.v1"
  131. assert json.loads(values["itemset_ids"]) == [1608352]
  132. assert json.loads(values["pattern_seed_pack"])["schema_version"] == "runtime_record.v1"
  133. def test_database_runtime_appends_jsonl_with_raw_payload():
  134. connection = FakeConnection()
  135. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  136. store.append_jsonl(
  137. "run_001",
  138. "search_queries.jsonl",
  139. [
  140. {
  141. "record_schema_version": "runtime_record.v1",
  142. "run_id": "run_001",
  143. "policy_run_id": "policy_run_001",
  144. "search_query_id": "q_001",
  145. "search_query": "对比分析",
  146. "search_query_generation_method": "item_single",
  147. "pattern_seed_ref": {
  148. "source_field": "seed_terms",
  149. "source_index": 0,
  150. "seed_term": "对比分析",
  151. },
  152. "raw_payload": {"run_id": "run_001", "search_query_id": "q_001"},
  153. }
  154. ],
  155. )
  156. sql, params = connection.statements[-1]
  157. values = _insert_values(sql, params)
  158. assert "INSERT INTO `content_agent_queries`" in sql
  159. assert values["schema_version"] == "content_agent.v1"
  160. assert json.loads(values["pattern_seed_ref"])["seed_term"] == "对比分析"
  161. assert json.loads(values["raw_payload"])["search_query_id"] == "q_001"
  162. def test_database_runtime_upserts_content_media_records():
  163. connection = FakeConnection()
  164. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  165. store.append_jsonl(
  166. "run_001",
  167. "content_media_records.jsonl",
  168. [
  169. {
  170. "record_schema_version": "runtime_record.v1",
  171. "run_id": "run_001",
  172. "policy_run_id": "policy_run_001",
  173. "platform": "douyin",
  174. "platform_content_id": "content_001",
  175. "content_media_status": "oss_uploaded",
  176. "content_metadata_source": "douyin_keyword_search",
  177. "play_url": "https://source.example/video.mp4",
  178. "local_path": None,
  179. "oss_url": "https://res.example/video.mp4",
  180. "raw_payload": {
  181. "run_id": "run_001",
  182. "platform_content_id": "content_001",
  183. "oss_object_key": "crawler/video/content_001.mp4",
  184. },
  185. "created_at": "2026-06-15T00:00:00+00:00",
  186. }
  187. ],
  188. )
  189. sql, _ = connection.statements[-1]
  190. assert "INSERT INTO `content_agent_content_media_records`" in sql
  191. assert "ON DUPLICATE KEY UPDATE" in sql
  192. assert "`oss_url` = VALUES(`oss_url`)" in sql
  193. def test_database_runtime_upserts_failed_search_query_status():
  194. connection = FakeConnection()
  195. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  196. store.append_jsonl(
  197. "run_001",
  198. "search_queries.jsonl",
  199. [
  200. {
  201. "record_schema_version": "runtime_record.v1",
  202. "run_id": "run_001",
  203. "policy_run_id": "policy_run_001",
  204. "search_query_id": "q_001",
  205. "search_query": "接口失败",
  206. "search_query_generation_method": "item_single",
  207. "search_query_effect_status": "failed",
  208. "raw_payload": {
  209. "run_id": "run_001",
  210. "policy_run_id": "policy_run_001",
  211. "search_query_id": "q_001",
  212. "search_query_effect_status": "failed",
  213. "query_failure": {
  214. "status": "failed",
  215. "error_code": "PLATFORM_REQUEST_FAILED",
  216. },
  217. },
  218. }
  219. ],
  220. )
  221. sql, params = connection.statements[-1]
  222. values = _insert_values(sql, params)
  223. assert "INSERT INTO `content_agent_queries`" in sql
  224. assert "ON DUPLICATE KEY UPDATE" in sql
  225. assert values["search_query_effect_status"] == "failed"
  226. payload = json.loads(values["raw_payload"])
  227. assert payload["query_failure"]["error_code"] == "PLATFORM_REQUEST_FAILED"
  228. def test_database_runtime_preserves_llm_variant_payload_fields():
  229. connection = FakeConnection()
  230. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  231. store.append_jsonl(
  232. "run_001",
  233. "search_queries.jsonl",
  234. [
  235. {
  236. "record_schema_version": "runtime_record.v1",
  237. "run_id": "run_001",
  238. "policy_run_id": "policy_run_001",
  239. "search_query_id": "q_002",
  240. "search_query": "人物叙事素材",
  241. "search_query_generation_method": "llm_variant",
  242. "pattern_seed_ref": {
  243. "source_field": "seed_terms",
  244. "source_index": 0,
  245. "seed_term": "人物故事",
  246. },
  247. "llm_variant_of": "q_001",
  248. "raw_payload": {
  249. "run_id": "run_001",
  250. "policy_run_id": "policy_run_001",
  251. "search_query_id": "q_002",
  252. "search_query_generation_method": "llm_variant",
  253. "llm_variant_of": "q_001",
  254. "llm_input_evidence": {"seed_term": "人物故事"},
  255. "llm_prompt_version": "fake-query-prompt-v1",
  256. "llm_generation_model": "fake-query-model",
  257. },
  258. }
  259. ],
  260. )
  261. sql, params = connection.statements[-1]
  262. values = _insert_values(sql, params)
  263. assert "INSERT INTO `content_agent_queries`" in sql
  264. assert "llm_variant_of" not in values
  265. payload = json.loads(values["raw_payload"])
  266. assert payload["llm_variant_of"] == "q_001"
  267. assert payload["llm_input_evidence"]["seed_term"] == "人物故事"
  268. assert payload["llm_prompt_version"] == "fake-query-prompt-v1"
  269. assert payload["llm_generation_model"] == "fake-query-model"
  270. def test_database_runtime_upserts_pattern_recall_evidence():
  271. connection = FakeConnection()
  272. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  273. # V3 形状: Gemini 直读判定证据(recall_status=judged + evidence_summary 判定字段);
  274. # decode 时代 7 列已随 V2 数据归档删除,列过滤必须把混入的旧键挡在 DB 外。
  275. store.append_jsonl(
  276. "run_001",
  277. "pattern_recall_evidence.jsonl",
  278. [
  279. {
  280. "record_schema_version": "runtime_record.v1",
  281. "run_id": "run_001",
  282. "policy_run_id": "policy_run_001",
  283. "recall_evidence_id": "recall_001",
  284. "content_discovery_id": "content_001",
  285. "platform": "douyin",
  286. "platform_content_id": "7390000000000000000",
  287. "recall_status": "judged",
  288. "decode_status": "success",
  289. "matched_terms": ["爱国情感"],
  290. "evidence_summary": {
  291. "judge_status": "ok",
  292. "fit_senior_50plus": True,
  293. "relevance_score": 0.85,
  294. },
  295. "raw_payload": {
  296. "platform": "douyin",
  297. "judge_reason": "贴题且适合 50+",
  298. },
  299. }
  300. ],
  301. )
  302. sql, params = connection.statements[-1]
  303. values = _insert_values(sql, params)
  304. assert "INSERT INTO `content_agent_pattern_recall_evidence`" in sql
  305. assert "ON DUPLICATE KEY UPDATE" in sql
  306. assert values["schema_version"] == "content_agent.v1"
  307. assert values["recall_evidence_id"] == "recall_001"
  308. assert values["recall_status"] == "judged"
  309. assert "platform" not in values
  310. assert "decode_status" not in values
  311. assert "matched_terms" not in values
  312. assert json.loads(values["evidence_summary"])["fit_senior_50plus"] is True
  313. assert json.loads(values["raw_payload"])["platform"] == "douyin"
  314. def test_database_runtime_preserves_p5_rule_decision_fields():
  315. connection = FakeConnection()
  316. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  317. store.append_jsonl(
  318. "run_001",
  319. "rule_decisions.jsonl",
  320. [
  321. {
  322. "record_schema_version": "runtime_record.v1",
  323. "run_id": "run_001",
  324. "policy_run_id": "policy_run_001",
  325. "decision_id": "d_001",
  326. "policy_bundle_id": "policy_bundle_v1",
  327. "rule_pack_id": "douyin_content_discovery_rule_pack_v1",
  328. "rule_pack_version": "1.0.0",
  329. "strategy_version": "V1",
  330. "decision_target_type": "content",
  331. "decision_target_id": "content_001",
  332. "decision_action": "REJECT_CONTENT",
  333. "decision_reason_code": "pattern_recall_failed",
  334. "search_query_effect_status": "rule_blocked",
  335. "score": None,
  336. "triggered_blocking_rules": ["gate_pattern_recall_failed"],
  337. "scorecard": {"total_score": None, "score_missing": True},
  338. "decision_replay_data": {
  339. "policy_bundle_hash": "hash_001",
  340. "dispatch_id": "dispatch_content",
  341. "effect_mapping_id": "map_hard_gate_reject_rule_blocked",
  342. },
  343. "raw_payload": {
  344. "decision_id": "d_001",
  345. "search_query_effect_status": "rule_blocked",
  346. "decision_replay_data": {
  347. "policy_bundle_hash": "hash_001",
  348. "dispatch_id": "dispatch_content",
  349. },
  350. },
  351. }
  352. ],
  353. )
  354. sql, params = connection.statements[-1]
  355. values = _insert_values(sql, params)
  356. assert "INSERT INTO `content_agent_rule_decisions`" in sql
  357. assert values["search_query_effect_status"] == "rule_blocked"
  358. assert json.loads(values["triggered_blocking_rules"]) == ["gate_pattern_recall_failed"]
  359. assert json.loads(values["scorecard"])["score_missing"] is True
  360. assert json.loads(values["decision_replay_data"])["dispatch_id"] == "dispatch_content"
  361. assert json.loads(values["raw_payload"])["decision_replay_data"]["policy_bundle_hash"] == "hash_001"
  362. def test_database_runtime_preserves_v4_score_and_walk_json_contract():
  363. connection = FakeConnection()
  364. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  365. store.append_jsonl(
  366. "run_001",
  367. "rule_decisions.jsonl",
  368. [
  369. {
  370. "record_schema_version": "runtime_record.v1",
  371. "run_id": "run_001",
  372. "policy_run_id": "policy_run_001",
  373. "decision_id": "d_v4_001",
  374. "policy_bundle_id": "policy_bundle_v4",
  375. "rule_pack_id": "douyin_content_discovery_rule_pack_v4",
  376. "rule_pack_version": "4.0.0",
  377. "strategy_version": "V4",
  378. "decision_target_type": "content",
  379. "decision_target_id": "content_001",
  380. "decision_action": "ADD_TO_CONTENT_POOL",
  381. "decision_reason_code": "v4_query_and_platform_pass",
  382. "search_query_effect_status": "success",
  383. "score": 80,
  384. "scorecard": {
  385. "schema_version": "v4_scorecard.v1",
  386. "query_relevance_score": 82,
  387. "platform_performance_score": 78,
  388. "missing_observable_fields": ["view_count"],
  389. },
  390. "decision_replay_data": {
  391. "policy_bundle_hash": "hash_v4",
  392. "rule_pack_id": "douyin_content_discovery_rule_pack_v4",
  393. "rule_pack_version": "4.0.0",
  394. "dispatch_id": "dispatch_content_v4",
  395. "strategy_version": "V4",
  396. "allow_walk": True,
  397. "walk_gate_snapshot": {
  398. "query_relevance_score": 82,
  399. "platform_performance_score": 78,
  400. "score": 80,
  401. },
  402. },
  403. "raw_payload": {
  404. "decision_id": "d_v4_001",
  405. "v4_contract": {
  406. "query_relevance_score": 82,
  407. "platform_performance_score": 78,
  408. },
  409. },
  410. }
  411. ],
  412. )
  413. sql, params = connection.statements[-1]
  414. values = _insert_values(sql, params)
  415. assert "INSERT INTO `content_agent_rule_decisions`" in sql
  416. scorecard = json.loads(values["scorecard"])
  417. replay_data = json.loads(values["decision_replay_data"])
  418. assert scorecard["schema_version"] == "v4_scorecard.v1"
  419. assert scorecard["query_relevance_score"] == 82
  420. assert scorecard["platform_performance_score"] == 78
  421. assert scorecard["missing_observable_fields"] == ["view_count"]
  422. assert replay_data["allow_walk"] is True
  423. assert replay_data["walk_gate_snapshot"]["score"] == 80
  424. def test_database_runtime_preserves_p5_search_clue_aggregation_in_raw_payload():
  425. connection = FakeConnection()
  426. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  427. store.append_jsonl(
  428. "run_001",
  429. "search_clues.jsonl",
  430. [
  431. {
  432. "record_schema_version": "runtime_record.v1",
  433. "run_id": "run_001",
  434. "policy_run_id": "policy_run_001",
  435. "clue_id": "clue_001",
  436. "search_query_id": "q_001",
  437. "search_query": "爱国情感",
  438. "discovery_start_source": "pattern_seed",
  439. "previous_discovery_step": "search_query_generated",
  440. "result_count": 1,
  441. "pooled_content_count": 0,
  442. "review_content_count": 0,
  443. "pending_content_count": 0,
  444. "rejected_content_count": 1,
  445. "search_query_effect_status": "rule_blocked",
  446. "query_aggregation_id": "agg_query_rule_blocked",
  447. "walk_next_step": "stop_search_query",
  448. "raw_payload": {
  449. "clue_id": "clue_001",
  450. "query_aggregation_id": "agg_query_rule_blocked",
  451. "effect_status_counts": {"rule_blocked": 1},
  452. },
  453. }
  454. ],
  455. )
  456. sql, params = connection.statements[-1]
  457. values = _insert_values(sql, params)
  458. assert "INSERT INTO `content_agent_search_clues`" in sql
  459. assert "query_aggregation_id" not in values
  460. assert values["search_query_effect_status"] == "rule_blocked"
  461. assert json.loads(values["raw_payload"])["query_aggregation_id"] == "agg_query_rule_blocked"
  462. def test_database_runtime_writes_failed_search_clue_and_platform_query_failed_event():
  463. connection = FakeConnection()
  464. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  465. query_failure = {
  466. "search_query_id": "q_001",
  467. "search_query": "接口失败",
  468. "search_query_generation_method": "item_single",
  469. "status": "failed",
  470. "error_code": "PLATFORM_REQUEST_FAILED",
  471. "message": "platform query failed",
  472. "error_detail": {
  473. "business_code": "10000",
  474. "business_data": {"trace_id": "trace-q_001", "upstream_reason": "opaque"},
  475. "trace_refs": {"x-request-id": "req-q_001"},
  476. },
  477. }
  478. store.append_jsonl(
  479. "run_001",
  480. "search_clues.jsonl",
  481. [
  482. {
  483. "record_schema_version": "runtime_record.v1",
  484. "run_id": "run_001",
  485. "policy_run_id": "policy_run_001",
  486. "clue_id": "clue_001",
  487. "search_query_id": "q_001",
  488. "search_query": "接口失败",
  489. "discovery_start_source": "pattern_itemset",
  490. "previous_discovery_step": "pattern_search_query",
  491. "result_count": 0,
  492. "pooled_content_count": 0,
  493. "review_content_count": 0,
  494. "pending_content_count": 0,
  495. "rejected_content_count": 0,
  496. "search_query_effect_status": "failed",
  497. "query_aggregation_id": "platform_query_failure",
  498. "walk_next_step": "stop_search_query",
  499. "raw_payload": {
  500. "clue_id": "clue_001",
  501. "query_failure": query_failure,
  502. },
  503. }
  504. ],
  505. )
  506. store.append_jsonl(
  507. "run_001",
  508. "run_events.jsonl",
  509. [
  510. {
  511. "record_schema_version": "runtime_record.v1",
  512. "run_id": "run_001",
  513. "policy_run_id": "policy_run_001",
  514. "event_id": "evt_001",
  515. "event_type": "platform_query_failed",
  516. "status": "failed",
  517. "input_ref": "search_queries.jsonl:q_001",
  518. "output_ref": "search_clues.jsonl",
  519. "error_code": "PLATFORM_REQUEST_FAILED",
  520. "message": "platform query failed",
  521. "raw_payload": {
  522. "event_id": "evt_001",
  523. "query_failure": query_failure,
  524. },
  525. }
  526. ],
  527. )
  528. clue_sql, clue_params = connection.statements[-2]
  529. clue_values = _insert_values(clue_sql, clue_params)
  530. assert "INSERT INTO `content_agent_search_clues`" in clue_sql
  531. assert "ON DUPLICATE KEY UPDATE" in clue_sql
  532. assert clue_values["search_query_effect_status"] == "failed"
  533. assert clue_values["walk_next_step"] == "stop_search_query"
  534. clue_payload = json.loads(clue_values["raw_payload"])
  535. assert clue_payload["query_failure"]["search_query_id"] == "q_001"
  536. assert clue_payload["query_failure"]["error_detail"]["business_data"] == {
  537. "trace_id": "trace-q_001",
  538. "upstream_reason": "opaque",
  539. }
  540. event_sql, event_params = connection.statements[-1]
  541. event_values = _insert_values(event_sql, event_params)
  542. assert "INSERT INTO `content_agent_run_events`" in event_sql
  543. assert "ON DUPLICATE KEY UPDATE" in event_sql
  544. assert event_values["event_type"] == "platform_query_failed"
  545. assert event_values["status"] == "failed"
  546. assert event_values["error_code"] == "PLATFORM_REQUEST_FAILED"
  547. event_payload = json.loads(event_values["raw_payload"])
  548. assert event_payload["query_failure"]["error_detail"]["trace_refs"] == {
  549. "x-request-id": "req-q_001"
  550. }
  551. def test_database_runtime_writes_publish_jobs_db_only_records():
  552. connection = FakeConnection()
  553. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  554. store.write_publish_jobs(
  555. "run_001",
  556. "policy_run_001",
  557. [
  558. {
  559. "publish_job_id": "publish_job_001",
  560. "platform_content_id": "7390000000000000000",
  561. "job_status": "created",
  562. "trigger_mode": "manual_review",
  563. "request_payload": {
  564. "decision_id": "decision_001",
  565. "source_path_record_ids": ["path_001"],
  566. },
  567. "response_payload": {},
  568. }
  569. ],
  570. )
  571. sql, params = connection.statements[-1]
  572. values = _insert_values(sql, params)
  573. assert "INSERT INTO `content_agent_publish_jobs`" in sql
  574. assert "ON DUPLICATE KEY UPDATE" in sql
  575. assert values["schema_version"] == "content_agent.v1"
  576. assert values["run_id"] == "run_001"
  577. assert values["policy_run_id"] == "policy_run_001"
  578. assert values["publish_job_id"] == "publish_job_001"
  579. assert values["job_status"] == "created"
  580. assert values["trigger_mode"] == "manual_review"
  581. assert json.loads(values["request_payload"])["decision_id"] == "decision_001"
  582. def test_database_runtime_writes_author_assets():
  583. connection = FakeConnection()
  584. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  585. store.write_author_assets(
  586. [
  587. {
  588. "author_asset_id": "author_asset_001",
  589. "platform": "douyin",
  590. "platform_author_id": "author_001",
  591. "author_display_name": "作者一号",
  592. "asset_status": "active",
  593. "source_type": "runtime_author_work",
  594. "validation_status": "validated",
  595. "eligible_as_source": 1,
  596. "elderly_ratio": 0.72,
  597. "elderly_tgi": 138,
  598. "content_tags": ["人物故事"],
  599. "source_run_id": "run_001",
  600. "source_policy_run_id": "policy_run_001",
  601. "profile_snapshot": {"sample_count": 9},
  602. "evidence_refs": {"decision_ids": ["d_001"]},
  603. "raw_payload": {"author_asset_id": "author_asset_001"},
  604. }
  605. ]
  606. )
  607. sql, params = connection.statements[-1]
  608. values = _insert_values(sql, params)
  609. assert "INSERT INTO `content_agent_author_assets`" in sql
  610. assert "ON DUPLICATE KEY UPDATE" in sql
  611. assert values["schema_version"] == "content_agent.v1"
  612. assert values["author_asset_id"] == "author_asset_001"
  613. assert values["platform_author_id"] == "author_001"
  614. assert values["eligible_as_source"] == 1
  615. assert values["elderly_ratio"] == 0.72
  616. assert values["elderly_tgi"] == 138
  617. assert json.loads(values["content_tags"]) == ["人物故事"]
  618. assert json.loads(values["profile_snapshot"])["sample_count"] == 9
  619. assert json.loads(values["evidence_refs"])["decision_ids"] == ["d_001"]
  620. def test_database_runtime_writes_author_asset_roles():
  621. connection = FakeConnection()
  622. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  623. store.write_author_asset_roles(
  624. [
  625. {
  626. "author_asset_id": "author_asset_001",
  627. "role": "source_seed",
  628. "role_status": "active",
  629. "role_reason_code": "p7_author_asset_eligible",
  630. "assigned_by": "system",
  631. "source_run_id": "run_001",
  632. "raw_payload": {"role": "source_seed"},
  633. }
  634. ]
  635. )
  636. sql, params = connection.statements[-1]
  637. values = _insert_values(sql, params)
  638. assert "INSERT INTO `content_agent_author_asset_roles`" in sql
  639. assert "ON DUPLICATE KEY UPDATE" in sql
  640. assert values["schema_version"] == "content_agent.v1"
  641. assert values["author_asset_id"] == "author_asset_001"
  642. assert values["role"] == "source_seed"
  643. assert values["assigned_by"] == "system"
  644. assert json.loads(values["raw_payload"])["role"] == "source_seed"
  645. def test_database_runtime_writes_search_clue_assets():
  646. connection = FakeConnection()
  647. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  648. store.write_search_clue_assets(
  649. [
  650. {
  651. "search_clue_asset_id": "search_clue_asset_001",
  652. "platform": "douyin",
  653. "clue_type": "search_query",
  654. "normalized_clue_text": "银发旅行",
  655. "display_clue_text": "银发旅行",
  656. "promotion_status": "promoted",
  657. "reusable_priority": 1,
  658. "can_seed_next_run": 1,
  659. "first_seen_run_id": "run_001",
  660. "first_seen_policy_run_id": "policy_run_001",
  661. "summary_metrics": {"pooled_content_count": 1},
  662. "raw_payload": {"promotion_reason": "success_search_clue"},
  663. }
  664. ]
  665. )
  666. sql, params = connection.statements[-1]
  667. values = _insert_values(sql, params)
  668. assert "INSERT INTO `content_agent_search_clue_assets`" in sql
  669. assert "ON DUPLICATE KEY UPDATE" in sql
  670. assert values["schema_version"] == "content_agent.v1"
  671. assert values["search_clue_asset_id"] == "search_clue_asset_001"
  672. assert values["can_seed_next_run"] == 1
  673. assert json.loads(values["summary_metrics"])["pooled_content_count"] == 1
  674. def test_database_runtime_writes_search_clue_asset_evidence():
  675. connection = FakeConnection()
  676. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  677. store.write_search_clue_asset_evidence(
  678. [
  679. {
  680. "evidence_id": "search_clue_evidence_001",
  681. "search_clue_asset_id": "search_clue_asset_001",
  682. "run_id": "run_001",
  683. "policy_run_id": "policy_run_001",
  684. "clue_id": "clue_001",
  685. "search_query_id": "q_001",
  686. "pooled_content_count": 1,
  687. "source_path_record_ids": ["path_001"],
  688. "decision_ids": ["decision_001"],
  689. "performance_feedback_refs": [],
  690. "raw_payload": {"clue_id": "clue_001"},
  691. }
  692. ]
  693. )
  694. sql, params = connection.statements[-1]
  695. values = _insert_values(sql, params)
  696. assert "INSERT INTO `content_agent_search_clue_asset_evidence`" in sql
  697. assert "ON DUPLICATE KEY UPDATE" in sql
  698. assert values["schema_version"] == "content_agent.v1"
  699. assert values["run_id"] == "run_001"
  700. assert values["clue_id"] == "clue_001"
  701. assert json.loads(values["source_path_record_ids"]) == ["path_001"]
  702. assert json.loads(values["decision_ids"]) == ["decision_001"]
  703. def test_database_runtime_update_final_output_upserts_validation_status():
  704. connection = FakeConnection()
  705. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  706. store.update_json(
  707. "run_001",
  708. "final_output.json",
  709. {
  710. "schema_version": "runtime_record.v1",
  711. "run_id": "run_001",
  712. "policy_run_id": "policy_run_001",
  713. "summary": {"run_path_complete": True},
  714. "validation_status": "pass",
  715. },
  716. )
  717. sql, params = connection.statements[-1]
  718. values = _insert_values(sql, params)
  719. assert "INSERT INTO `content_agent_final_outputs`" in sql
  720. assert "ON DUPLICATE KEY UPDATE" in sql
  721. assert values["validation_status"] == "pass"
  722. assert json.loads(values["summary"])["run_path_complete"] is True
  723. def test_database_runtime_preserves_m5_explanation_payloads():
  724. connection = FakeConnection()
  725. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  726. explanation = {
  727. "schema_version": "v4_decision_explanation.v1",
  728. "query_relevance_score": 80,
  729. "platform_performance_score": 70,
  730. "score": 75,
  731. "allow_walk": True,
  732. "walk_gate_snapshot": {"score": 75},
  733. }
  734. v4_summary = {
  735. "schema_version": "v4_strategy_review_summary.v1",
  736. "score_buckets": {"pool": 1},
  737. "allow_walk_distribution": {"allowed": 1, "denied": 0, "missing": 0},
  738. "walk_gate_review": {"v4_gate_distribution": {"allowed": 1}},
  739. }
  740. store.update_json(
  741. "run_001",
  742. "final_output.json",
  743. {
  744. "schema_version": "runtime_record.v1",
  745. "run_id": "run_001",
  746. "policy_run_id": "policy_run_001",
  747. "summary": {},
  748. "validation_status": "pass",
  749. "decision_records": [
  750. {
  751. "decision_id": "decision_001",
  752. "v4_explanation": explanation,
  753. }
  754. ],
  755. },
  756. )
  757. store.update_json(
  758. "run_001",
  759. "strategy_review.json",
  760. {
  761. "schema_version": "runtime_record.v1",
  762. "run_id": "run_001",
  763. "policy_run_id": "policy_run_001",
  764. "review_id": "review_001",
  765. "review_status": "generated",
  766. "summary": {},
  767. "v4_summary": v4_summary,
  768. "effective_search_queries": [],
  769. "weak_search_queries": [],
  770. "top_reject_reasons": [],
  771. "productive_paths": [],
  772. "suggestions": [],
  773. "raw_payload": {"v4_summary": v4_summary},
  774. },
  775. )
  776. final_values = _insert_values(*connection.statements[-2])
  777. final_payload = json.loads(final_values["final_output"])
  778. assert final_payload["decision_records"][0]["v4_explanation"] == explanation
  779. review_values = _insert_values(*connection.statements[-1])
  780. review_payload = json.loads(review_values["raw_payload"])
  781. assert review_payload["v4_summary"] == v4_summary
  782. def test_database_runtime_reads_performance_feedback_payloads():
  783. connection = FakeConnection()
  784. connection.select_all_result = [
  785. {
  786. "raw_payload": json.dumps(
  787. {
  788. "run_id": "run_001",
  789. "policy_run_id": "policy_run_001",
  790. "feedback_id": "feedback_001",
  791. "completion_rate": 0.8,
  792. }
  793. )
  794. }
  795. ]
  796. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  797. rows = store.read_performance_feedback("run_001", "policy_run_001")
  798. sql, params = connection.statements[-1]
  799. assert "FROM `content_agent_performance_feedback`" in sql
  800. assert params == ["run_001", "policy_run_001"]
  801. assert rows[0]["feedback_id"] == "feedback_001"
  802. assert rows[0]["raw_payload"]["completion_rate"] == 0.8
  803. def test_database_runtime_read_jsonl_reconstructs_runtime_payload():
  804. connection = FakeConnection()
  805. connection.select_all_result = [
  806. {
  807. "raw_payload": json.dumps(
  808. {
  809. "record_schema_version": "runtime_record.v1",
  810. "run_id": "run_001",
  811. "policy_run_id": "policy_run_001",
  812. "search_query_id": "q_001",
  813. }
  814. )
  815. }
  816. ]
  817. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  818. rows = store.read_jsonl("run_001", "search_queries.jsonl")
  819. assert rows[0]["search_query_id"] == "q_001"
  820. assert rows[0]["raw_payload"]["search_query_id"] == "q_001"
  821. def test_database_runtime_read_jsonl_reconstructs_rule_decision_from_formal_columns():
  822. connection = FakeConnection()
  823. connection.select_all_result = [
  824. {
  825. "run_id": "run_001",
  826. "policy_run_id": "policy_run_001",
  827. "decision_id": "d_001",
  828. "policy_bundle_id": "bundle_001",
  829. "rule_pack_id": "rule_pack_001",
  830. "rule_pack_version": "4.0.0",
  831. "strategy_version": "V4",
  832. "decision_target_type": "content",
  833. "decision_target_id": "content_001",
  834. "decision_action": "ADD_TO_CONTENT_POOL",
  835. "decision_reason_code": "v4_query_and_platform_pass",
  836. "search_query_effect_status": "success",
  837. "score": Decimal("75.50"),
  838. "triggered_blocking_rules": json.dumps([]),
  839. "scorecard": json.dumps({"schema_version": "v4_scorecard.v1", "total_score": 75.5}),
  840. "source_evidence": json.dumps({"source_post_id": "post_001"}),
  841. "decision_replay_data": json.dumps({"allow_walk": True}),
  842. "raw_payload": json.dumps(
  843. {
  844. "record_schema_version": "runtime_record.v1",
  845. "strategy_id": "strategy_001",
  846. "policy_bundle_hash": "hash_001",
  847. }
  848. ),
  849. }
  850. ]
  851. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  852. rows = store.read_jsonl("run_001", "rule_decisions.jsonl")
  853. sql, params = connection.statements[-1]
  854. assert "FROM `content_agent_rule_decisions`" in sql
  855. assert "`scorecard`" in sql
  856. assert params == ["run_001"]
  857. row = rows[0]
  858. assert row["record_schema_version"] == "runtime_record.v1"
  859. assert row["decision_id"] == "d_001"
  860. assert row["score"] == 75.5
  861. assert row["scorecard"]["schema_version"] == "v4_scorecard.v1"
  862. assert row["source_evidence"]["source_post_id"] == "post_001"
  863. assert row["decision_replay_data"]["allow_walk"] is True
  864. assert row["raw_payload"] == {
  865. "record_schema_version": "runtime_record.v1",
  866. "strategy_id": "strategy_001",
  867. "policy_bundle_hash": "hash_001",
  868. }
  869. def test_database_runtime_rejects_forbidden_raw_payload_keys_in_lists():
  870. connection = FakeConnection()
  871. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  872. try:
  873. store.append_jsonl(
  874. "run_001",
  875. "search_queries.jsonl",
  876. [
  877. {
  878. "record_schema_version": "runtime_record.v1",
  879. "run_id": "run_001",
  880. "policy_run_id": "policy_run_001",
  881. "search_query_id": "q_001",
  882. "search_query": "对比分析",
  883. "raw_payload": {"items": [{"dsn": "should_not_be_stored"}]},
  884. }
  885. ],
  886. )
  887. except ValueError as exc:
  888. assert "forbidden key" in str(exc)
  889. else:
  890. raise AssertionError("expected forbidden raw_payload key to be rejected")
  891. def test_database_runtime_update_run_record_ignores_empty_sanitized_updates():
  892. connection = FakeConnection()
  893. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  894. store.update_run_record("run_001", {"unknown_field": "ignored", "status": None})
  895. assert connection.statements == []
  896. def test_database_runtime_update_run_record_persists_platform_failure_detail():
  897. connection = FakeConnection()
  898. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  899. store.update_run_record(
  900. "run_001",
  901. {
  902. "status": "failed",
  903. "error_code": "PLATFORM_REQUEST_FAILED",
  904. "error_detail": {
  905. "query_failures": [
  906. {
  907. "search_query_id": "q_001",
  908. "status": "failed",
  909. "error_code": "PLATFORM_REQUEST_FAILED",
  910. }
  911. ]
  912. },
  913. },
  914. )
  915. sql, params = connection.statements[-1]
  916. assert "UPDATE `content_agent_runs` SET" in sql
  917. assert "`error_detail` = %s" in sql
  918. assert json.loads(params[2])["query_failures"][0]["search_query_id"] == "q_001"
  919. assert params[-1] == "run_001"
  920. def test_database_runtime_upserts_media_pipeline_task():
  921. connection = FakeConnection()
  922. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  923. task = {
  924. "task_id": "mpt_001",
  925. "task_type": "oss_upload",
  926. "idempotency_key": "run_001:q_001_p01:douyin:c1:oss_upload",
  927. "run_id": "run_001",
  928. "policy_run_id": "policy_run_001",
  929. "batch_id": "q_001_p01_initial_top3",
  930. "platform": "douyin",
  931. "platform_content_id": "c1",
  932. "status": "queued",
  933. "attempt_count": 0,
  934. "lease_until": None,
  935. "next_retry_at": None,
  936. "input_payload": {"play_url_present": True},
  937. "result_payload": None,
  938. "error_summary": None,
  939. }
  940. stored = store.enqueue_media_pipeline_task(task)
  941. sql, params = connection.statements[-1]
  942. values = _insert_values(sql, params)
  943. assert stored == task
  944. assert "INSERT INTO `content_agent_media_pipeline_tasks`" in sql
  945. assert "ON DUPLICATE KEY UPDATE" in sql
  946. assert values["schema_version"] == "content_agent.v1"
  947. assert values["task_id"] == "mpt_001"
  948. assert json.loads(values["input_payload"])["play_url_present"] is True
  949. def test_database_runtime_updates_media_pipeline_task_status():
  950. connection = FakeConnection()
  951. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  952. store.update_media_pipeline_task(
  953. "mpt_001",
  954. {
  955. "status": "completed",
  956. "attempt_count": 1,
  957. "result_payload": {"oss_url_present": True},
  958. },
  959. )
  960. sql, params = connection.statements[-1]
  961. assert "UPDATE `content_agent_media_pipeline_tasks` SET" in sql
  962. assert "`status` = %s" in sql
  963. assert "`result_payload` = %s" in sql
  964. assert json.loads(params[2])["oss_url_present"] is True
  965. assert params[-1] == "mpt_001"
  966. def test_business_modules_do_not_import_or_name_database_tables():
  967. root = Path("content_agent/business_modules")
  968. text = "\n".join(path.read_text(encoding="utf-8") for path in root.rglob("*.py"))
  969. assert not re.search(
  970. r"pymysql|sqlalchemy|psycopg|sqlite3|SELECT |INSERT |UPDATE |DELETE |SHOW |CREATE |ALTER |content_agent_",
  971. text,
  972. )
  973. def _config():
  974. return ContentSupplyDbConfig(
  975. host="127.0.0.1",
  976. port=3306,
  977. user="content_rw",
  978. password="dummy_password",
  979. database="content-deconstruction-supply",
  980. )
  981. def _insert_values(sql, params):
  982. match = re.search(r"\((.*?)\) VALUES", sql)
  983. assert match, sql
  984. columns = [part.strip().strip("`") for part in match.group(1).split(",")]
  985. return dict(zip(columns, params))