test_runtime_files.py 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628
  1. import json
  2. from pathlib import Path
  3. from content_agent.errors import ErrorCode
  4. from content_agent.integrations.runtime_files import LocalRuntimeFileStore
  5. from content_agent.run_service import RunService
  6. from content_agent.schemas import RunStartRequest
  7. from tests.p1_helpers import FakeQueryVariantClient, REAL_SOURCE_FIXTURE
  8. def _start_mock_run(tmp_path, **kwargs):
  9. service = RunService(
  10. runtime_root=tmp_path / "runtime" / "v1",
  11. query_variant_client=FakeQueryVariantClient(),
  12. )
  13. kwargs.setdefault("source", str(REAL_SOURCE_FIXTURE))
  14. state = service.start_run(RunStartRequest(platform_mode="mock", **kwargs))
  15. assert state["status"] == "success"
  16. return service, state["run_id"]
  17. def test_runtime_files_are_parseable_and_consistent(tmp_path):
  18. service, run_id = _start_mock_run(tmp_path)
  19. run_dir = service.runtime.run_dir(run_id)
  20. for json_file in [
  21. "source_context.json",
  22. "pattern_seed_pack.json",
  23. "final_output.json",
  24. "strategy_review.json",
  25. ]:
  26. data = json.loads((run_dir / json_file).read_text(encoding="utf-8"))
  27. assert data["schema_version"] == "runtime_record.v1"
  28. for jsonl_file in [
  29. "search_queries.jsonl",
  30. "discovered_content_items.jsonl",
  31. "content_media_records.jsonl",
  32. "pattern_recall_evidence.jsonl",
  33. "rule_decisions.jsonl",
  34. "walk_actions.jsonl",
  35. "run_events.jsonl",
  36. "source_path_records.jsonl",
  37. "search_clues.jsonl",
  38. ]:
  39. for line in (run_dir / jsonl_file).read_text(encoding="utf-8").splitlines():
  40. json.loads(line)
  41. items = service.read_jsonl(run_id, "discovered_content_items.jsonl")
  42. decisions = service.read_jsonl(run_id, "rule_decisions.jsonl")
  43. paths = service.read_jsonl(run_id, "source_path_records.jsonl")
  44. final_output = service.read_json(run_id, "final_output.json")
  45. source_context = service.read_json(run_id, "source_context.json")
  46. policy_run_id = service.read_json(run_id, "pattern_seed_pack.json")["policy_run_id"]
  47. content_ids = {item["platform_content_id"] for item in items}
  48. decision_target_ids = {decision["decision_target_id"] for decision in decisions}
  49. assert content_ids == decision_target_ids
  50. media_records = service.read_jsonl(run_id, "content_media_records.jsonl")
  51. recall_evidence = service.read_jsonl(run_id, "pattern_recall_evidence.jsonl")
  52. walk_actions = service.read_jsonl(run_id, "walk_actions.jsonl")
  53. for row in [*items, *media_records, *recall_evidence, *decisions, *walk_actions, *paths]:
  54. assert row["run_id"] == run_id
  55. assert row["policy_run_id"] == policy_run_id
  56. assert row["record_schema_version"] == "runtime_record.v1"
  57. if "run_id" in row["raw_payload"]:
  58. assert row["raw_payload"]["run_id"] == run_id
  59. assert {row["walk_status"] for row in walk_actions} <= {
  60. "success",
  61. "pending",
  62. "failed",
  63. "skipped",
  64. "rule_blocked",
  65. }
  66. assert all(row["walk_action_id"].startswith("wa_") for row in walk_actions)
  67. for media_record in media_records:
  68. assert media_record["platform"] == "douyin"
  69. assert final_output["policy_run_id"] == policy_run_id
  70. assert final_output["policy"]["policy_bundle_id"] == "douyin_policy_bundle_v1"
  71. assert final_output["policy"]["rule_pack_source_ref"]["file"].endswith(
  72. "douyin_rule_packs.v1.json"
  73. )
  74. assert final_output["walk_strategy"]["walk_strategy_version"] == "V1.0"
  75. assert "policy_run_id" not in source_context
  76. assert {decision["decision_action"] for decision in decisions} <= {
  77. "ADD_TO_CONTENT_POOL",
  78. "KEEP_CONTENT_FOR_REVIEW",
  79. "REJECT_CONTENT",
  80. }
  81. assert {decision["search_query_effect_status"] for decision in decisions} <= {
  82. "success",
  83. "pending",
  84. "failed",
  85. "rule_blocked",
  86. }
  87. search_clues = service.read_jsonl(run_id, "search_clues.jsonl")
  88. assert {clue["search_query_effect_status"] for clue in search_clues} <= {
  89. "success",
  90. "pending",
  91. "failed",
  92. "rule_blocked",
  93. }
  94. path_ids = {path["source_path_record_id"] for path in paths}
  95. for asset in final_output["content_assets"]:
  96. assert set(asset["source_path_record_ids"]) <= path_ids
  97. evidence_pack = source_context["ext_data"]["evidence_pack"]
  98. assert evidence_pack["pattern_source_system"] == "pg_pattern_v2"
  99. assert evidence_pack["pattern_execution_id"] == 581
  100. assert evidence_pack["mining_config_id"] == 2082
  101. assert evidence_pack["itemset_ids"] == [1608352]
  102. validation = service.validate_run(run_id)
  103. assert validation["status"] == "pass"
  104. def test_content_media_records_replace_same_video_row(tmp_path):
  105. runtime = LocalRuntimeFileStore(tmp_path / "runtime")
  106. runtime.prepare_run("run_001")
  107. base = {
  108. "record_schema_version": "runtime_record.v1",
  109. "run_id": "run_001",
  110. "policy_run_id": "policy_run_001",
  111. "platform": "douyin",
  112. "platform_content_id": "content_001",
  113. "content_media_status": "metadata_only",
  114. "content_metadata_source": "douyin_keyword_search",
  115. "play_url": "https://source.example/video.mp4",
  116. "local_path": None,
  117. "oss_url": None,
  118. "raw_payload": {"run_id": "run_001", "platform_content_id": "content_001"},
  119. "created_at": "2026-06-15T00:00:00+00:00",
  120. }
  121. runtime.append_jsonl("run_001", "content_media_records.jsonl", [base])
  122. runtime.append_jsonl(
  123. "run_001",
  124. "content_media_records.jsonl",
  125. [
  126. {
  127. **base,
  128. "content_media_status": "oss_uploaded",
  129. "oss_url": "https://res.example/video.mp4",
  130. "raw_payload": {
  131. **base["raw_payload"],
  132. "oss_object_key": "crawler/video/content_001.mp4",
  133. },
  134. }
  135. ],
  136. )
  137. rows = runtime.read_jsonl("run_001", "content_media_records.jsonl")
  138. assert len(rows) == 1
  139. assert rows[0]["content_media_status"] == "oss_uploaded"
  140. assert rows[0]["oss_url"] == "https://res.example/video.mp4"
  141. def test_all_platform_query_failure_writes_failed_query_runtime_records(tmp_path):
  142. service = RunService(
  143. runtime_root=tmp_path / "runtime" / "v1",
  144. query_variant_client=FakeQueryVariantClient(),
  145. )
  146. service._platform_client = lambda platform, platform_mode: _AllFailurePlatformClient()
  147. state = service.start_run(
  148. RunStartRequest(platform_mode="real", source=str(REAL_SOURCE_FIXTURE))
  149. )
  150. assert state["status"] == "failed"
  151. assert state["error_code"] == ErrorCode.PLATFORM_REQUEST_FAILED.value
  152. search_queries = service.read_jsonl(state["run_id"], "search_queries.jsonl")
  153. search_clues = service.read_jsonl(state["run_id"], "search_clues.jsonl")
  154. run_events = service.read_jsonl(state["run_id"], "run_events.jsonl")
  155. assert len(search_queries) == len(state["error_detail"]["query_failures"])
  156. assert {query["search_query_effect_status"] for query in search_queries} == {"failed"}
  157. assert all(query["raw_payload"]["query_failure"]["status"] == "failed" for query in search_queries)
  158. assert {clue["search_query_effect_status"] for clue in search_clues} == {"failed"}
  159. assert {clue["result_count"] for clue in search_clues} == {0}
  160. assert {clue["query_aggregation_id"] for clue in search_clues} == {"platform_query_failure"}
  161. platform_failures = [
  162. event for event in run_events if event["event_type"] == "platform_query_failed"
  163. ]
  164. assert len(platform_failures) == len(search_queries)
  165. assert {event["error_code"] for event in platform_failures} == {
  166. ErrorCode.PLATFORM_REQUEST_FAILED.value
  167. }
  168. def test_runtime_validation_catches_summary_drift(tmp_path):
  169. service, run_id = _start_mock_run(tmp_path)
  170. final_output_path = service.runtime.run_dir(run_id) / "final_output.json"
  171. final_output = json.loads(final_output_path.read_text(encoding="utf-8"))
  172. final_output["summary"]["pooled_content_count"] = 99
  173. final_output_path.write_text(
  174. json.dumps(final_output, ensure_ascii=False, indent=2) + "\n",
  175. encoding="utf-8",
  176. )
  177. validation = service.validate_run(run_id)
  178. assert validation["status"] == "fail"
  179. assert any(finding["check_id"] == "summary_mismatch" for finding in validation["findings"])
  180. def test_local_runtime_replaces_pattern_recall_evidence_by_id(tmp_path):
  181. runtime = LocalRuntimeFileStore(tmp_path / "runtime")
  182. runtime.prepare_run("run_001")
  183. runtime.append_jsonl(
  184. "run_001",
  185. "pattern_recall_evidence.jsonl",
  186. [
  187. {
  188. "run_id": "run_001",
  189. "policy_run_id": "policy_001",
  190. "recall_evidence_id": "recall_001",
  191. "recall_status": "pending",
  192. }
  193. ],
  194. )
  195. runtime.append_jsonl(
  196. "run_001",
  197. "pattern_recall_evidence.jsonl",
  198. [
  199. {
  200. "run_id": "run_001",
  201. "policy_run_id": "policy_001",
  202. "recall_evidence_id": "recall_001",
  203. "recall_status": "matched",
  204. }
  205. ],
  206. )
  207. rows = runtime.read_jsonl("run_001", "pattern_recall_evidence.jsonl")
  208. assert len(rows) == 1
  209. assert rows[0]["recall_status"] == "matched"
  210. def test_runtime_validation_catches_missing_policy_run_id(tmp_path):
  211. service, run_id = _start_mock_run(tmp_path)
  212. decisions_path = service.runtime.run_dir(run_id) / "rule_decisions.jsonl"
  213. decisions = [
  214. json.loads(line)
  215. for line in decisions_path.read_text(encoding="utf-8").splitlines()
  216. if line.strip()
  217. ]
  218. decisions[0].pop("policy_run_id")
  219. decisions_path.write_text(
  220. "".join(
  221. json.dumps(decision, ensure_ascii=False, separators=(",", ":")) + "\n"
  222. for decision in decisions
  223. ),
  224. encoding="utf-8",
  225. )
  226. validation = service.validate_run(run_id)
  227. assert validation["status"] == "fail"
  228. assert any(finding["check_id"] == "policy_run_id_mismatch" for finding in validation["findings"])
  229. def test_runtime_validation_catches_missing_record_schema_version(tmp_path):
  230. service, run_id = _start_mock_run(tmp_path)
  231. queries_path = service.runtime.run_dir(run_id) / "search_queries.jsonl"
  232. queries = [
  233. json.loads(line)
  234. for line in queries_path.read_text(encoding="utf-8").splitlines()
  235. if line.strip()
  236. ]
  237. queries[0].pop("record_schema_version")
  238. queries_path.write_text(
  239. "".join(
  240. json.dumps(query, ensure_ascii=False, separators=(",", ":")) + "\n"
  241. for query in queries
  242. ),
  243. encoding="utf-8",
  244. )
  245. validation = service.validate_run(run_id)
  246. assert validation["status"] == "fail"
  247. assert any(
  248. finding["check_id"] == "record_schema_version_missing"
  249. for finding in validation["findings"]
  250. )
  251. class _AllFailurePlatformClient:
  252. def search(self, search_query: dict) -> list[dict]:
  253. raise RuntimeError("platform unavailable")
  254. def test_runtime_validation_catches_missing_raw_payload(tmp_path):
  255. service, run_id = _start_mock_run(tmp_path)
  256. media_path = service.runtime.run_dir(run_id) / "content_media_records.jsonl"
  257. media_records = [
  258. json.loads(line)
  259. for line in media_path.read_text(encoding="utf-8").splitlines()
  260. if line.strip()
  261. ]
  262. media_records[0].pop("raw_payload")
  263. media_path.write_text(
  264. "".join(
  265. json.dumps(media_record, ensure_ascii=False, separators=(",", ":")) + "\n"
  266. for media_record in media_records
  267. ),
  268. encoding="utf-8",
  269. )
  270. validation = service.validate_run(run_id)
  271. assert validation["status"] == "fail"
  272. assert any(finding["check_id"] == "raw_payload_missing" for finding in validation["findings"])
  273. def test_runtime_validation_catches_forbidden_raw_payload_key(tmp_path):
  274. service, run_id = _start_mock_run(tmp_path)
  275. media_path = service.runtime.run_dir(run_id) / "content_media_records.jsonl"
  276. media_records = [
  277. json.loads(line)
  278. for line in media_path.read_text(encoding="utf-8").splitlines()
  279. if line.strip()
  280. ]
  281. media_records[0]["raw_payload"]["secret"] = "should_not_be_stored"
  282. media_path.write_text(
  283. "".join(
  284. json.dumps(media_record, ensure_ascii=False, separators=(",", ":")) + "\n"
  285. for media_record in media_records
  286. ),
  287. encoding="utf-8",
  288. )
  289. validation = service.validate_run(run_id)
  290. assert validation["status"] == "fail"
  291. assert any(
  292. finding["check_id"] == "raw_payload_forbidden_key"
  293. for finding in validation["findings"]
  294. )
  295. def test_runtime_validation_catches_missing_pattern_recall_evidence(tmp_path):
  296. service, run_id = _start_mock_run(tmp_path)
  297. items_path = service.runtime.run_dir(run_id) / "discovered_content_items.jsonl"
  298. items = [
  299. json.loads(line)
  300. for line in items_path.read_text(encoding="utf-8").splitlines()
  301. if line.strip()
  302. ]
  303. items[0]["pattern_match_result"].pop("pattern_recall_evidence_id", None)
  304. items_path.write_text(
  305. "".join(
  306. json.dumps(item, ensure_ascii=False, separators=(",", ":")) + "\n"
  307. for item in items
  308. ),
  309. encoding="utf-8",
  310. )
  311. validation = service.validate_run(run_id)
  312. assert validation["status"] == "fail"
  313. assert any(
  314. finding["check_id"] == "pattern_recall_evidence_missing"
  315. for finding in validation["findings"]
  316. )
  317. def test_runtime_validation_allows_missing_decode_case_ids_in_source_evidence(tmp_path):
  318. service, run_id = _start_mock_run(tmp_path)
  319. run_dir = service.runtime.run_dir(run_id)
  320. decisions_path = run_dir / "rule_decisions.jsonl"
  321. decisions = [
  322. json.loads(line)
  323. for line in decisions_path.read_text(encoding="utf-8").splitlines()
  324. if line.strip()
  325. ]
  326. for decision in decisions:
  327. decision["source_evidence"].pop("decode_case_ids", None)
  328. decisions_path.write_text(
  329. "".join(
  330. json.dumps(decision, ensure_ascii=False, separators=(",", ":")) + "\n"
  331. for decision in decisions
  332. ),
  333. encoding="utf-8",
  334. )
  335. final_output_path = run_dir / "final_output.json"
  336. final_output = json.loads(final_output_path.read_text(encoding="utf-8"))
  337. for section in ["content_assets", "reject_records", "decision_records"]:
  338. for row in final_output.get(section, []):
  339. row.get("source_evidence", {}).pop("decode_case_ids", None)
  340. final_output_path.write_text(
  341. json.dumps(final_output, ensure_ascii=False, indent=2) + "\n",
  342. encoding="utf-8",
  343. )
  344. validation = service.validate_run(run_id)
  345. assert validation["status"] == "pass"
  346. def test_runtime_validation_catches_missing_final_decision_record(tmp_path):
  347. service, run_id = _start_mock_run(tmp_path)
  348. final_output_path = service.runtime.run_dir(run_id) / "final_output.json"
  349. final_output = json.loads(final_output_path.read_text(encoding="utf-8"))
  350. final_output["decision_records"] = final_output["decision_records"][:-1]
  351. final_output_path.write_text(
  352. json.dumps(final_output, ensure_ascii=False, indent=2) + "\n",
  353. encoding="utf-8",
  354. )
  355. validation = service.validate_run(run_id)
  356. assert validation["status"] == "fail"
  357. assert any(finding["check_id"] == "final_decision_missing" for finding in validation["findings"])
  358. def test_runtime_validation_catches_platform_content_id_source_pollution(tmp_path):
  359. service, run_id = _start_mock_run(tmp_path)
  360. decisions_path = service.runtime.run_dir(run_id) / "rule_decisions.jsonl"
  361. decisions = [
  362. json.loads(line)
  363. for line in decisions_path.read_text(encoding="utf-8").splitlines()
  364. if line.strip()
  365. ]
  366. source_evidence = decisions[0]["source_evidence"]
  367. source_evidence["source_post_id"] = source_evidence["discovered_platform_content_id"]
  368. decisions_path.write_text(
  369. "".join(
  370. json.dumps(decision, ensure_ascii=False, separators=(",", ":")) + "\n"
  371. for decision in decisions
  372. ),
  373. encoding="utf-8",
  374. )
  375. validation = service.validate_run(run_id)
  376. assert validation["status"] == "fail"
  377. assert any(
  378. finding["check_id"] == "source_evidence_content_pollution"
  379. for finding in validation["findings"]
  380. )
  381. def test_runtime_validation_catches_reject_source_path_break(tmp_path):
  382. service, run_id = _start_mock_run(tmp_path)
  383. paths_path = service.runtime.run_dir(run_id) / "source_path_records.jsonl"
  384. paths = [
  385. json.loads(line)
  386. for line in paths_path.read_text(encoding="utf-8").splitlines()
  387. if line.strip()
  388. ]
  389. paths = [
  390. path
  391. for path in paths
  392. if not (
  393. path.get("source_path_type") == "search_query_to_content"
  394. and path.get("to_node_id") == "7390000000000000099"
  395. )
  396. ]
  397. paths_path.write_text(
  398. "".join(json.dumps(path, ensure_ascii=False, separators=(",", ":")) + "\n" for path in paths),
  399. encoding="utf-8",
  400. )
  401. validation = service.validate_run(run_id)
  402. assert validation["status"] == "fail"
  403. assert any(finding["check_id"] == "source_path_broken" for finding in validation["findings"])
  404. def test_real_source_fixture_keeps_upstream_evidence_pack(tmp_path):
  405. source_path = Path("tests/fixtures/real_case_source/source_context.json")
  406. service, run_id = _start_mock_run(tmp_path, source=str(source_path))
  407. source_context = service.read_json(run_id, "source_context.json")
  408. evidence_pack = source_context["ext_data"]["evidence_pack"]
  409. decisions = service.read_jsonl(run_id, "rule_decisions.jsonl")
  410. assert evidence_pack["pattern_source_system"] == "pg_pattern_v2"
  411. assert evidence_pack["source_certainty"] == "db_validated"
  412. assert evidence_pack["validation_status"] == "passed"
  413. assert evidence_pack["source_post_id"] == "51978710"
  414. assert evidence_pack["pattern_execution_id"] == 581
  415. assert evidence_pack["mining_config_id"] == 2082
  416. assert evidence_pack["itemset_ids"] == [1608352]
  417. assert evidence_pack["upstream_run_id"] == "f405f129-3341-4f4a-98e6-fd3f73632adb"
  418. assert evidence_pack["support"] == 0.0045734552921411235
  419. assert evidence_pack["absolute_support"] == 49
  420. assert evidence_pack["decode_case_ids"] == []
  421. assert decisions[0]["source_evidence"]["source_certainty"] == "db_validated"
  422. assert (
  423. decisions[0]["source_evidence"]["discovered_platform_content_id"]
  424. not in evidence_pack["matched_post_ids"]
  425. )
  426. assert service.validate_run(run_id)["status"] == "pass"
  427. def test_demand_content_json_array_source_is_adapted_to_source_context(tmp_path):
  428. source_path = tmp_path / "demand_content.json"
  429. source_path.write_text(
  430. json.dumps(
  431. [
  432. {
  433. "id": 123,
  434. "merge_leve2": "PG Pattern V2 需求测试",
  435. "name": "爱国情感,人物故事",
  436. "suggestion": None,
  437. "score": 1.0,
  438. "dt": "20260604",
  439. "ext_data": {
  440. "evidence_pack": {
  441. "source_kind": "pattern_itemset",
  442. "pattern_source_system": "pg_pattern_v2",
  443. "case_id_type": "post_id",
  444. "source_post_id": "51978710",
  445. "pattern_execution_id": 581,
  446. "mining_config_id": 2082,
  447. "itemset_ids": [1608352],
  448. "itemset_items": [{"itemset_id": 1608352}],
  449. "category_bindings": [{"category_id": 76006}],
  450. "element_bindings": [{"category_id": 76006}],
  451. "support": 0.0045734552921411235,
  452. "absolute_support": 49,
  453. "matched_post_ids": ["51978710"],
  454. "video_ids": ["51978710"],
  455. "case_ids": ["51978710"],
  456. "decode_case_ids": [],
  457. "seed_terms": ["爱国情感", "人物故事"],
  458. "source_certainty": "db_validated",
  459. "validation_status": "passed",
  460. }
  461. },
  462. }
  463. ],
  464. ensure_ascii=False,
  465. ),
  466. encoding="utf-8",
  467. )
  468. service, run_id = _start_mock_run(tmp_path, source=str(source_path))
  469. source_context = service.read_json(run_id, "source_context.json")
  470. evidence_pack = source_context["ext_data"]["evidence_pack"]
  471. assert source_context["demand_content_id"] == "123"
  472. assert evidence_pack["pattern_source_system"] == "pg_pattern_v2"
  473. assert evidence_pack["pattern_execution_id"] == 581
  474. assert evidence_pack["mining_config_id"] == 2082
  475. assert evidence_pack["itemset_ids"] == [1608352]
  476. assert service.validate_run(run_id)["status"] == "pass"
  477. def test_old_mysql_source_system_is_rejected(tmp_path):
  478. source_path = tmp_path / "source_context.json"
  479. source_path.write_text(
  480. json.dumps(
  481. {
  482. "run_id": "old_run",
  483. "demand_content_id": "old",
  484. "merge_leve2": "历史样例",
  485. "name": "旧 MySQL 样例",
  486. "ext_data": {
  487. "evidence_pack": {
  488. "source_kind": "pattern_itemset",
  489. "pattern_source_system": "mysql_topic_pattern",
  490. "case_id_type": "post_id",
  491. "source_post_id": "51978710",
  492. "pattern_execution_id": 581,
  493. "mining_config_id": 2082,
  494. "itemset_ids": [1608352],
  495. "itemset_items": [{"itemset_id": 1608352}],
  496. "category_bindings": [{"category_id": 76006}],
  497. "support": 0.0045734552921411235,
  498. "absolute_support": 49,
  499. "matched_post_ids": ["51978710"],
  500. "video_ids": ["51978710"],
  501. "case_ids": ["51978710"],
  502. "decode_case_ids": [],
  503. "seed_terms": ["爱国情感", "人物故事"],
  504. "source_certainty": "db_validated",
  505. "validation_status": "passed",
  506. }
  507. },
  508. },
  509. ensure_ascii=False,
  510. ),
  511. encoding="utf-8",
  512. )
  513. service = RunService(runtime_root=tmp_path / "runtime" / "v1")
  514. state = service.start_run(RunStartRequest(platform_mode="mock", source=str(source_path)))
  515. assert state["status"] == "failed"
  516. assert state["error_code"] == ErrorCode.INVALID_SOURCE.value
  517. assert state["errors"] == ["invalid source"]
  518. def test_high_weight_element_source_passes_validation(tmp_path):
  519. # M12A:非 itemset 来源(高权重元素,itemset 三件套为空)应获准进入并跑完。
  520. source_path = Path("tests/fixtures/v2_no_itemset/high_weight_element_source_context.json")
  521. service, run_id = _start_mock_run(tmp_path, source=str(source_path))
  522. evidence_pack = service.read_json(run_id, "source_context.json")["ext_data"]["evidence_pack"]
  523. assert evidence_pack["source_kind"] == "high_weight_element"
  524. assert evidence_pack["itemset_ids"] == []
  525. assert service.validate_run(run_id)["status"] == "pass"
  526. def test_multi_source_with_empty_itemset_is_rejected(tmp_path):
  527. # M12A:多来源里含 itemset 项却缺 itemset_ids → 仍拒。
  528. payload = json.loads(
  529. Path("tests/fixtures/v2_no_itemset/multi_source_source_context.json").read_text(encoding="utf-8")
  530. )
  531. payload["ext_data"]["evidence_pack"]["itemset_ids"] = []
  532. source_path = tmp_path / "bad_multi_source.json"
  533. source_path.write_text(json.dumps(payload, ensure_ascii=False), encoding="utf-8")
  534. service = RunService(
  535. runtime_root=tmp_path / "runtime" / "v1",
  536. query_variant_client=FakeQueryVariantClient(),
  537. )
  538. state = service.start_run(RunStartRequest(platform_mode="mock", source=str(source_path)))
  539. assert state["status"] == "failed"
  540. assert state["error_code"] == ErrorCode.INVALID_SOURCE.value