test_database_runtime.py 40 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106
  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. }
  473. store.append_jsonl(
  474. "run_001",
  475. "search_clues.jsonl",
  476. [
  477. {
  478. "record_schema_version": "runtime_record.v1",
  479. "run_id": "run_001",
  480. "policy_run_id": "policy_run_001",
  481. "clue_id": "clue_001",
  482. "search_query_id": "q_001",
  483. "search_query": "接口失败",
  484. "discovery_start_source": "pattern_itemset",
  485. "previous_discovery_step": "pattern_search_query",
  486. "result_count": 0,
  487. "pooled_content_count": 0,
  488. "review_content_count": 0,
  489. "pending_content_count": 0,
  490. "rejected_content_count": 0,
  491. "search_query_effect_status": "failed",
  492. "query_aggregation_id": "platform_query_failure",
  493. "walk_next_step": "stop_search_query",
  494. "raw_payload": {
  495. "clue_id": "clue_001",
  496. "query_failure": query_failure,
  497. },
  498. }
  499. ],
  500. )
  501. store.append_jsonl(
  502. "run_001",
  503. "run_events.jsonl",
  504. [
  505. {
  506. "record_schema_version": "runtime_record.v1",
  507. "run_id": "run_001",
  508. "policy_run_id": "policy_run_001",
  509. "event_id": "evt_001",
  510. "event_type": "platform_query_failed",
  511. "status": "failed",
  512. "input_ref": "search_queries.jsonl:q_001",
  513. "output_ref": "search_clues.jsonl",
  514. "error_code": "PLATFORM_REQUEST_FAILED",
  515. "message": "platform query failed",
  516. "raw_payload": {
  517. "event_id": "evt_001",
  518. "query_failure": query_failure,
  519. },
  520. }
  521. ],
  522. )
  523. clue_sql, clue_params = connection.statements[-2]
  524. clue_values = _insert_values(clue_sql, clue_params)
  525. assert "INSERT INTO `content_agent_search_clues`" in clue_sql
  526. assert "ON DUPLICATE KEY UPDATE" in clue_sql
  527. assert clue_values["search_query_effect_status"] == "failed"
  528. assert clue_values["walk_next_step"] == "stop_search_query"
  529. assert json.loads(clue_values["raw_payload"])["query_failure"]["search_query_id"] == "q_001"
  530. event_sql, event_params = connection.statements[-1]
  531. event_values = _insert_values(event_sql, event_params)
  532. assert "INSERT INTO `content_agent_run_events`" in event_sql
  533. assert "ON DUPLICATE KEY UPDATE" in event_sql
  534. assert event_values["event_type"] == "platform_query_failed"
  535. assert event_values["status"] == "failed"
  536. assert event_values["error_code"] == "PLATFORM_REQUEST_FAILED"
  537. def test_database_runtime_writes_publish_jobs_db_only_records():
  538. connection = FakeConnection()
  539. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  540. store.write_publish_jobs(
  541. "run_001",
  542. "policy_run_001",
  543. [
  544. {
  545. "publish_job_id": "publish_job_001",
  546. "platform_content_id": "7390000000000000000",
  547. "job_status": "created",
  548. "trigger_mode": "manual_review",
  549. "request_payload": {
  550. "decision_id": "decision_001",
  551. "source_path_record_ids": ["path_001"],
  552. },
  553. "response_payload": {},
  554. }
  555. ],
  556. )
  557. sql, params = connection.statements[-1]
  558. values = _insert_values(sql, params)
  559. assert "INSERT INTO `content_agent_publish_jobs`" in sql
  560. assert "ON DUPLICATE KEY UPDATE" in sql
  561. assert values["schema_version"] == "content_agent.v1"
  562. assert values["run_id"] == "run_001"
  563. assert values["policy_run_id"] == "policy_run_001"
  564. assert values["publish_job_id"] == "publish_job_001"
  565. assert values["job_status"] == "created"
  566. assert values["trigger_mode"] == "manual_review"
  567. assert json.loads(values["request_payload"])["decision_id"] == "decision_001"
  568. def test_database_runtime_writes_author_assets():
  569. connection = FakeConnection()
  570. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  571. store.write_author_assets(
  572. [
  573. {
  574. "author_asset_id": "author_asset_001",
  575. "platform": "douyin",
  576. "platform_author_id": "author_001",
  577. "author_display_name": "作者一号",
  578. "asset_status": "active",
  579. "source_type": "runtime_author_work",
  580. "validation_status": "validated",
  581. "eligible_as_source": 1,
  582. "elderly_ratio": 0.72,
  583. "elderly_tgi": 138,
  584. "content_tags": ["人物故事"],
  585. "source_run_id": "run_001",
  586. "source_policy_run_id": "policy_run_001",
  587. "profile_snapshot": {"sample_count": 9},
  588. "evidence_refs": {"decision_ids": ["d_001"]},
  589. "raw_payload": {"author_asset_id": "author_asset_001"},
  590. }
  591. ]
  592. )
  593. sql, params = connection.statements[-1]
  594. values = _insert_values(sql, params)
  595. assert "INSERT INTO `content_agent_author_assets`" in sql
  596. assert "ON DUPLICATE KEY UPDATE" in sql
  597. assert values["schema_version"] == "content_agent.v1"
  598. assert values["author_asset_id"] == "author_asset_001"
  599. assert values["platform_author_id"] == "author_001"
  600. assert values["eligible_as_source"] == 1
  601. assert values["elderly_ratio"] == 0.72
  602. assert values["elderly_tgi"] == 138
  603. assert json.loads(values["content_tags"]) == ["人物故事"]
  604. assert json.loads(values["profile_snapshot"])["sample_count"] == 9
  605. assert json.loads(values["evidence_refs"])["decision_ids"] == ["d_001"]
  606. def test_database_runtime_writes_author_asset_roles():
  607. connection = FakeConnection()
  608. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  609. store.write_author_asset_roles(
  610. [
  611. {
  612. "author_asset_id": "author_asset_001",
  613. "role": "source_seed",
  614. "role_status": "active",
  615. "role_reason_code": "p7_author_asset_eligible",
  616. "assigned_by": "system",
  617. "source_run_id": "run_001",
  618. "raw_payload": {"role": "source_seed"},
  619. }
  620. ]
  621. )
  622. sql, params = connection.statements[-1]
  623. values = _insert_values(sql, params)
  624. assert "INSERT INTO `content_agent_author_asset_roles`" in sql
  625. assert "ON DUPLICATE KEY UPDATE" in sql
  626. assert values["schema_version"] == "content_agent.v1"
  627. assert values["author_asset_id"] == "author_asset_001"
  628. assert values["role"] == "source_seed"
  629. assert values["assigned_by"] == "system"
  630. assert json.loads(values["raw_payload"])["role"] == "source_seed"
  631. def test_database_runtime_writes_search_clue_assets():
  632. connection = FakeConnection()
  633. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  634. store.write_search_clue_assets(
  635. [
  636. {
  637. "search_clue_asset_id": "search_clue_asset_001",
  638. "platform": "douyin",
  639. "clue_type": "search_query",
  640. "normalized_clue_text": "银发旅行",
  641. "display_clue_text": "银发旅行",
  642. "promotion_status": "promoted",
  643. "reusable_priority": 1,
  644. "can_seed_next_run": 1,
  645. "first_seen_run_id": "run_001",
  646. "first_seen_policy_run_id": "policy_run_001",
  647. "summary_metrics": {"pooled_content_count": 1},
  648. "raw_payload": {"promotion_reason": "success_search_clue"},
  649. }
  650. ]
  651. )
  652. sql, params = connection.statements[-1]
  653. values = _insert_values(sql, params)
  654. assert "INSERT INTO `content_agent_search_clue_assets`" in sql
  655. assert "ON DUPLICATE KEY UPDATE" in sql
  656. assert values["schema_version"] == "content_agent.v1"
  657. assert values["search_clue_asset_id"] == "search_clue_asset_001"
  658. assert values["can_seed_next_run"] == 1
  659. assert json.loads(values["summary_metrics"])["pooled_content_count"] == 1
  660. def test_database_runtime_writes_search_clue_asset_evidence():
  661. connection = FakeConnection()
  662. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  663. store.write_search_clue_asset_evidence(
  664. [
  665. {
  666. "evidence_id": "search_clue_evidence_001",
  667. "search_clue_asset_id": "search_clue_asset_001",
  668. "run_id": "run_001",
  669. "policy_run_id": "policy_run_001",
  670. "clue_id": "clue_001",
  671. "search_query_id": "q_001",
  672. "pooled_content_count": 1,
  673. "source_path_record_ids": ["path_001"],
  674. "decision_ids": ["decision_001"],
  675. "performance_feedback_refs": [],
  676. "raw_payload": {"clue_id": "clue_001"},
  677. }
  678. ]
  679. )
  680. sql, params = connection.statements[-1]
  681. values = _insert_values(sql, params)
  682. assert "INSERT INTO `content_agent_search_clue_asset_evidence`" in sql
  683. assert "ON DUPLICATE KEY UPDATE" in sql
  684. assert values["schema_version"] == "content_agent.v1"
  685. assert values["run_id"] == "run_001"
  686. assert values["clue_id"] == "clue_001"
  687. assert json.loads(values["source_path_record_ids"]) == ["path_001"]
  688. assert json.loads(values["decision_ids"]) == ["decision_001"]
  689. def test_database_runtime_update_final_output_upserts_validation_status():
  690. connection = FakeConnection()
  691. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  692. store.update_json(
  693. "run_001",
  694. "final_output.json",
  695. {
  696. "schema_version": "runtime_record.v1",
  697. "run_id": "run_001",
  698. "policy_run_id": "policy_run_001",
  699. "summary": {"run_path_complete": True},
  700. "validation_status": "pass",
  701. },
  702. )
  703. sql, params = connection.statements[-1]
  704. values = _insert_values(sql, params)
  705. assert "INSERT INTO `content_agent_final_outputs`" in sql
  706. assert "ON DUPLICATE KEY UPDATE" in sql
  707. assert values["validation_status"] == "pass"
  708. assert json.loads(values["summary"])["run_path_complete"] is True
  709. def test_database_runtime_preserves_m5_explanation_payloads():
  710. connection = FakeConnection()
  711. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  712. explanation = {
  713. "schema_version": "v4_decision_explanation.v1",
  714. "query_relevance_score": 80,
  715. "platform_performance_score": 70,
  716. "score": 75,
  717. "allow_walk": True,
  718. "walk_gate_snapshot": {"score": 75},
  719. }
  720. v4_summary = {
  721. "schema_version": "v4_strategy_review_summary.v1",
  722. "score_buckets": {"pool": 1},
  723. "allow_walk_distribution": {"allowed": 1, "denied": 0, "missing": 0},
  724. "walk_gate_review": {"v4_gate_distribution": {"allowed": 1}},
  725. }
  726. store.update_json(
  727. "run_001",
  728. "final_output.json",
  729. {
  730. "schema_version": "runtime_record.v1",
  731. "run_id": "run_001",
  732. "policy_run_id": "policy_run_001",
  733. "summary": {},
  734. "validation_status": "pass",
  735. "decision_records": [
  736. {
  737. "decision_id": "decision_001",
  738. "v4_explanation": explanation,
  739. }
  740. ],
  741. },
  742. )
  743. store.update_json(
  744. "run_001",
  745. "strategy_review.json",
  746. {
  747. "schema_version": "runtime_record.v1",
  748. "run_id": "run_001",
  749. "policy_run_id": "policy_run_001",
  750. "review_id": "review_001",
  751. "review_status": "generated",
  752. "summary": {},
  753. "v4_summary": v4_summary,
  754. "effective_search_queries": [],
  755. "weak_search_queries": [],
  756. "top_reject_reasons": [],
  757. "productive_paths": [],
  758. "suggestions": [],
  759. "raw_payload": {"v4_summary": v4_summary},
  760. },
  761. )
  762. final_values = _insert_values(*connection.statements[-2])
  763. final_payload = json.loads(final_values["final_output"])
  764. assert final_payload["decision_records"][0]["v4_explanation"] == explanation
  765. review_values = _insert_values(*connection.statements[-1])
  766. review_payload = json.loads(review_values["raw_payload"])
  767. assert review_payload["v4_summary"] == v4_summary
  768. def test_database_runtime_reads_performance_feedback_payloads():
  769. connection = FakeConnection()
  770. connection.select_all_result = [
  771. {
  772. "raw_payload": json.dumps(
  773. {
  774. "run_id": "run_001",
  775. "policy_run_id": "policy_run_001",
  776. "feedback_id": "feedback_001",
  777. "completion_rate": 0.8,
  778. }
  779. )
  780. }
  781. ]
  782. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  783. rows = store.read_performance_feedback("run_001", "policy_run_001")
  784. sql, params = connection.statements[-1]
  785. assert "FROM `content_agent_performance_feedback`" in sql
  786. assert params == ["run_001", "policy_run_001"]
  787. assert rows[0]["feedback_id"] == "feedback_001"
  788. assert rows[0]["raw_payload"]["completion_rate"] == 0.8
  789. def test_database_runtime_read_jsonl_reconstructs_runtime_payload():
  790. connection = FakeConnection()
  791. connection.select_all_result = [
  792. {
  793. "raw_payload": json.dumps(
  794. {
  795. "record_schema_version": "runtime_record.v1",
  796. "run_id": "run_001",
  797. "policy_run_id": "policy_run_001",
  798. "search_query_id": "q_001",
  799. }
  800. )
  801. }
  802. ]
  803. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  804. rows = store.read_jsonl("run_001", "search_queries.jsonl")
  805. assert rows[0]["search_query_id"] == "q_001"
  806. assert rows[0]["raw_payload"]["search_query_id"] == "q_001"
  807. def test_database_runtime_read_jsonl_reconstructs_rule_decision_from_formal_columns():
  808. connection = FakeConnection()
  809. connection.select_all_result = [
  810. {
  811. "run_id": "run_001",
  812. "policy_run_id": "policy_run_001",
  813. "decision_id": "d_001",
  814. "policy_bundle_id": "bundle_001",
  815. "rule_pack_id": "rule_pack_001",
  816. "rule_pack_version": "4.0.0",
  817. "strategy_version": "V4",
  818. "decision_target_type": "content",
  819. "decision_target_id": "content_001",
  820. "decision_action": "ADD_TO_CONTENT_POOL",
  821. "decision_reason_code": "v4_query_and_platform_pass",
  822. "search_query_effect_status": "success",
  823. "score": Decimal("75.50"),
  824. "triggered_blocking_rules": json.dumps([]),
  825. "scorecard": json.dumps({"schema_version": "v4_scorecard.v1", "total_score": 75.5}),
  826. "source_evidence": json.dumps({"source_post_id": "post_001"}),
  827. "decision_replay_data": json.dumps({"allow_walk": True}),
  828. "raw_payload": json.dumps(
  829. {
  830. "record_schema_version": "runtime_record.v1",
  831. "strategy_id": "strategy_001",
  832. "policy_bundle_hash": "hash_001",
  833. }
  834. ),
  835. }
  836. ]
  837. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  838. rows = store.read_jsonl("run_001", "rule_decisions.jsonl")
  839. sql, params = connection.statements[-1]
  840. assert "FROM `content_agent_rule_decisions`" in sql
  841. assert "`scorecard`" in sql
  842. assert params == ["run_001"]
  843. row = rows[0]
  844. assert row["record_schema_version"] == "runtime_record.v1"
  845. assert row["decision_id"] == "d_001"
  846. assert row["score"] == 75.5
  847. assert row["scorecard"]["schema_version"] == "v4_scorecard.v1"
  848. assert row["source_evidence"]["source_post_id"] == "post_001"
  849. assert row["decision_replay_data"]["allow_walk"] is True
  850. assert row["raw_payload"] == {
  851. "record_schema_version": "runtime_record.v1",
  852. "strategy_id": "strategy_001",
  853. "policy_bundle_hash": "hash_001",
  854. }
  855. def test_database_runtime_rejects_forbidden_raw_payload_keys_in_lists():
  856. connection = FakeConnection()
  857. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  858. try:
  859. store.append_jsonl(
  860. "run_001",
  861. "search_queries.jsonl",
  862. [
  863. {
  864. "record_schema_version": "runtime_record.v1",
  865. "run_id": "run_001",
  866. "policy_run_id": "policy_run_001",
  867. "search_query_id": "q_001",
  868. "search_query": "对比分析",
  869. "raw_payload": {"items": [{"dsn": "should_not_be_stored"}]},
  870. }
  871. ],
  872. )
  873. except ValueError as exc:
  874. assert "forbidden key" in str(exc)
  875. else:
  876. raise AssertionError("expected forbidden raw_payload key to be rejected")
  877. def test_database_runtime_update_run_record_ignores_empty_sanitized_updates():
  878. connection = FakeConnection()
  879. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  880. store.update_run_record("run_001", {"unknown_field": "ignored", "status": None})
  881. assert connection.statements == []
  882. def test_database_runtime_update_run_record_persists_platform_failure_detail():
  883. connection = FakeConnection()
  884. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  885. store.update_run_record(
  886. "run_001",
  887. {
  888. "status": "failed",
  889. "error_code": "PLATFORM_REQUEST_FAILED",
  890. "error_detail": {
  891. "query_failures": [
  892. {
  893. "search_query_id": "q_001",
  894. "status": "failed",
  895. "error_code": "PLATFORM_REQUEST_FAILED",
  896. }
  897. ]
  898. },
  899. },
  900. )
  901. sql, params = connection.statements[-1]
  902. assert "UPDATE `content_agent_runs` SET" in sql
  903. assert "`error_detail` = %s" in sql
  904. assert json.loads(params[2])["query_failures"][0]["search_query_id"] == "q_001"
  905. assert params[-1] == "run_001"
  906. def test_database_runtime_upserts_media_pipeline_task():
  907. connection = FakeConnection()
  908. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  909. task = {
  910. "task_id": "mpt_001",
  911. "task_type": "oss_upload",
  912. "idempotency_key": "run_001:q_001_p01:douyin:c1:oss_upload",
  913. "run_id": "run_001",
  914. "policy_run_id": "policy_run_001",
  915. "batch_id": "q_001_p01_initial_top3",
  916. "platform": "douyin",
  917. "platform_content_id": "c1",
  918. "status": "queued",
  919. "attempt_count": 0,
  920. "lease_until": None,
  921. "next_retry_at": None,
  922. "input_payload": {"play_url_present": True},
  923. "result_payload": None,
  924. "error_summary": None,
  925. }
  926. stored = store.enqueue_media_pipeline_task(task)
  927. sql, params = connection.statements[-1]
  928. values = _insert_values(sql, params)
  929. assert stored == task
  930. assert "INSERT INTO `content_agent_media_pipeline_tasks`" in sql
  931. assert "ON DUPLICATE KEY UPDATE" in sql
  932. assert values["schema_version"] == "content_agent.v1"
  933. assert values["task_id"] == "mpt_001"
  934. assert json.loads(values["input_payload"])["play_url_present"] is True
  935. def test_database_runtime_updates_media_pipeline_task_status():
  936. connection = FakeConnection()
  937. store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
  938. store.update_media_pipeline_task(
  939. "mpt_001",
  940. {
  941. "status": "completed",
  942. "attempt_count": 1,
  943. "result_payload": {"oss_url_present": True},
  944. },
  945. )
  946. sql, params = connection.statements[-1]
  947. assert "UPDATE `content_agent_media_pipeline_tasks` SET" in sql
  948. assert "`status` = %s" in sql
  949. assert "`result_payload` = %s" in sql
  950. assert json.loads(params[2])["oss_url_present"] is True
  951. assert params[-1] == "mpt_001"
  952. def test_business_modules_do_not_import_or_name_database_tables():
  953. root = Path("content_agent/business_modules")
  954. text = "\n".join(path.read_text(encoding="utf-8") for path in root.rglob("*.py"))
  955. assert not re.search(
  956. r"pymysql|sqlalchemy|psycopg|sqlite3|SELECT |INSERT |UPDATE |DELETE |SHOW |CREATE |ALTER |content_agent_",
  957. text,
  958. )
  959. def _config():
  960. return ContentSupplyDbConfig(
  961. host="127.0.0.1",
  962. port=3306,
  963. user="content_rw",
  964. password="dummy_password",
  965. database="content-deconstruction-supply",
  966. )
  967. def _insert_values(sql, params):
  968. match = re.search(r"\((.*?)\) VALUES", sql)
  969. assert match, sql
  970. columns = [part.strip().strip("`") for part in match.group(1).split(",")]
  971. return dict(zip(columns, params))