result_source_lookup.py 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585
  1. from __future__ import annotations
  2. import hashlib
  3. from collections import Counter
  4. from typing import Any
  5. from content_agent.constants import RUNTIME_SCHEMA_VERSION
  6. from content_agent.business_modules.run_record.validation import compute_final_output_completeness
  7. from content_agent.interfaces import RuntimeFileStore
  8. def run(
  9. run_id: str,
  10. policy_run_id: str,
  11. policy_bundle: dict[str, Any],
  12. discovered_content_items: list[dict[str, Any]],
  13. content_media_records: list[dict[str, Any]],
  14. decisions: list[dict[str, Any]],
  15. source_path_records: list[dict[str, Any]],
  16. search_clues: list[dict[str, Any]],
  17. runtime: RuntimeFileStore,
  18. ) -> dict[str, Any]:
  19. decision_by_target_id = {decision["decision_target_id"]: decision for decision in decisions}
  20. media_by_platform_content_id = {
  21. media["platform_content_id"]: media for media in content_media_records
  22. }
  23. paths_by_content_id = _paths_by_content_id(source_path_records)
  24. content_assets = _build_content_assets(
  25. policy_run_id,
  26. discovered_content_items,
  27. decision_by_target_id,
  28. media_by_platform_content_id,
  29. paths_by_content_id,
  30. )
  31. review_records = _build_review_records(
  32. policy_run_id,
  33. discovered_content_items,
  34. decision_by_target_id,
  35. media_by_platform_content_id,
  36. paths_by_content_id,
  37. )
  38. technical_retry_records = _build_technical_retry_records(
  39. policy_run_id,
  40. discovered_content_items,
  41. decision_by_target_id,
  42. media_by_platform_content_id,
  43. paths_by_content_id,
  44. )
  45. reject_records = _build_reject_records(
  46. policy_run_id,
  47. discovered_content_items,
  48. decision_by_target_id,
  49. media_by_platform_content_id,
  50. paths_by_content_id,
  51. )
  52. author_assets, author_asset_rows, author_role_rows = _build_author_assets(
  53. run_id,
  54. policy_run_id,
  55. discovered_content_items,
  56. decision_by_target_id,
  57. paths_by_content_id,
  58. )
  59. action_counts = Counter(decision["decision_action"] for decision in decisions)
  60. effect_status_counts = Counter(
  61. decision.get("search_query_effect_status", "failed") for decision in decisions
  62. )
  63. final_output = {
  64. "schema_version": RUNTIME_SCHEMA_VERSION,
  65. "run_id": run_id,
  66. "policy_run_id": policy_run_id,
  67. "policy": {
  68. "policy_bundle_id": policy_bundle["policy_bundle_id"],
  69. "strategy_id": policy_bundle["strategy_id"],
  70. "strategy_version": policy_bundle["strategy_version"],
  71. "rule_pack_id": policy_bundle["rule_pack_id"],
  72. "rule_pack_version": policy_bundle["rule_pack_version"],
  73. "policy_bundle_hash": policy_bundle["policy_bundle_hash"],
  74. "strategy_source_ref": policy_bundle["strategy_source_ref"],
  75. "rule_pack_source_ref": policy_bundle["rule_pack_source_ref"],
  76. },
  77. "walk_strategy": {
  78. "walk_strategy_id": policy_bundle.get("walk_strategy_id"),
  79. "walk_strategy_version": policy_bundle.get("walk_strategy_version"),
  80. "walk_strategy_source_ref": policy_bundle.get("walk_strategy_source_ref"),
  81. },
  82. "dispatch": policy_bundle.get("dispatch"),
  83. "dispatch_id": policy_bundle.get("dispatch_id"),
  84. "runtime_status_contract": policy_bundle.get("runtime_status_contract", {}),
  85. "content_assets": content_assets,
  86. "author_assets": author_assets,
  87. "review_records": review_records,
  88. "technical_retry_records": technical_retry_records,
  89. "decision_records": [
  90. {
  91. "decision_id": decision["decision_id"],
  92. "policy_run_id": policy_run_id,
  93. "rule_pack_id": decision["rule_pack_id"],
  94. "rule_pack_version": decision["rule_pack_version"],
  95. "strategy_version": decision["strategy_version"],
  96. "decision_target_id": decision["decision_target_id"],
  97. "decision_action": decision["decision_action"],
  98. "decision_reason_code": decision["decision_reason_code"],
  99. "search_query_effect_status": decision["search_query_effect_status"],
  100. "score": decision.get("score"),
  101. "source_evidence_ref": _source_evidence_ref(decision),
  102. **_v4_explanation_field(decision),
  103. }
  104. for decision in decisions
  105. ],
  106. "search_clues": [
  107. {
  108. "search_query_id": clue["search_query_id"],
  109. "policy_run_id": policy_run_id,
  110. "final_asset_status": "clue_only",
  111. "search_query_effect_status": clue["search_query_effect_status"],
  112. }
  113. for clue in search_clues
  114. ],
  115. "reject_records": reject_records,
  116. "summary": {
  117. "search_query_count": len(search_clues),
  118. "pooled_content_count": action_counts["ADD_TO_CONTENT_POOL"],
  119. "review_content_count": action_counts["KEEP_CONTENT_FOR_REVIEW"],
  120. "pending_content_count": 0,
  121. "rejected_content_count": action_counts["REJECT_CONTENT"],
  122. "technical_retry_content_count": action_counts["TECHNICAL_RETRY_REQUIRED"],
  123. "effect_status_counts": {
  124. "success": effect_status_counts["success"],
  125. "pending": effect_status_counts["pending"],
  126. "failed": effect_status_counts["failed"],
  127. "rule_blocked": effect_status_counts["rule_blocked"],
  128. },
  129. "policy_bundle_hash": policy_bundle["policy_bundle_hash"],
  130. },
  131. }
  132. completeness = compute_final_output_completeness(
  133. final_output,
  134. decisions,
  135. source_path_records,
  136. )
  137. final_output["validation_status"] = completeness["validation_status"]
  138. final_output["summary"]["run_path_complete"] = completeness["run_path_complete"]
  139. final_output["summary"]["trace_complete"] = completeness["trace_complete"]
  140. final_output["summary"]["validation_findings_summary"] = completeness["findings_summary"]
  141. runtime.write_json(run_id, "final_output.json", final_output)
  142. runtime.write_publish_jobs(
  143. run_id,
  144. policy_run_id,
  145. _build_publish_jobs(run_id, policy_run_id, content_assets),
  146. )
  147. runtime.write_author_assets(author_asset_rows)
  148. runtime.write_author_asset_roles(author_role_rows)
  149. return final_output
  150. def _build_content_assets(
  151. policy_run_id: str,
  152. discovered_content_items: list[dict[str, Any]],
  153. decision_by_target_id: dict[str, dict[str, Any]],
  154. media_by_platform_content_id: dict[str, dict[str, Any]],
  155. paths_by_content_id: dict[str, list[str]],
  156. ) -> list[dict[str, Any]]:
  157. content_assets: list[dict[str, Any]] = []
  158. for item in discovered_content_items:
  159. platform_content_id = item["platform_content_id"]
  160. decision = decision_by_target_id[platform_content_id]
  161. if decision["decision_action"] != "ADD_TO_CONTENT_POOL":
  162. continue
  163. path_ids = paths_by_content_id[platform_content_id]
  164. content_assets.append(
  165. _with_v4_explanation(
  166. {
  167. "platform": item["platform"],
  168. "platform_content_id": platform_content_id,
  169. "policy_run_id": policy_run_id,
  170. "content_discovery_id": item["content_discovery_id"],
  171. "final_asset_status": "pooled",
  172. "decision_id": decision["decision_id"],
  173. "rule_pack_id": decision["rule_pack_id"],
  174. "rule_pack_version": decision["rule_pack_version"],
  175. "strategy_version": decision["strategy_version"],
  176. "source_path_record_ids": path_ids,
  177. "source_evidence_ref": _source_evidence_ref(decision, path_ids),
  178. "content_media_status": media_by_platform_content_id[
  179. platform_content_id
  180. ]["content_media_status"],
  181. "media_snapshot": _media_snapshot(media_by_platform_content_id.get(platform_content_id)),
  182. },
  183. decision,
  184. )
  185. )
  186. return content_assets
  187. def _build_review_records(
  188. policy_run_id: str,
  189. discovered_content_items: list[dict[str, Any]],
  190. decision_by_target_id: dict[str, dict[str, Any]],
  191. media_by_platform_content_id: dict[str, dict[str, Any]],
  192. paths_by_content_id: dict[str, list[str]],
  193. ) -> list[dict[str, Any]]:
  194. review_records: list[dict[str, Any]] = []
  195. for item in discovered_content_items:
  196. platform_content_id = item["platform_content_id"]
  197. decision = decision_by_target_id[platform_content_id]
  198. if decision["decision_action"] != "KEEP_CONTENT_FOR_REVIEW":
  199. continue
  200. path_ids = paths_by_content_id[platform_content_id]
  201. review_records.append(
  202. _with_v4_explanation(
  203. {
  204. "platform": item["platform"],
  205. "platform_content_id": platform_content_id,
  206. "policy_run_id": policy_run_id,
  207. "content_discovery_id": item["content_discovery_id"],
  208. "review_status": "pending_review",
  209. "final_asset_status": "review_only",
  210. "decision_id": decision["decision_id"],
  211. "rule_pack_id": decision["rule_pack_id"],
  212. "rule_pack_version": decision["rule_pack_version"],
  213. "strategy_version": decision["strategy_version"],
  214. "decision_reason_code": decision["decision_reason_code"],
  215. "source_path_record_ids": path_ids,
  216. "source_evidence_ref": _source_evidence_ref(decision, path_ids),
  217. "content_media_status": media_by_platform_content_id[
  218. platform_content_id
  219. ]["content_media_status"],
  220. "media_snapshot": _media_snapshot(media_by_platform_content_id.get(platform_content_id)),
  221. },
  222. decision,
  223. )
  224. )
  225. return review_records
  226. def _build_reject_records(
  227. policy_run_id: str,
  228. discovered_content_items: list[dict[str, Any]],
  229. decision_by_target_id: dict[str, dict[str, Any]],
  230. media_by_platform_content_id: dict[str, dict[str, Any]],
  231. paths_by_content_id: dict[str, list[str]],
  232. ) -> list[dict[str, Any]]:
  233. reject_records: list[dict[str, Any]] = []
  234. for item in discovered_content_items:
  235. platform_content_id = item["platform_content_id"]
  236. decision = decision_by_target_id[platform_content_id]
  237. if decision["decision_action"] != "REJECT_CONTENT":
  238. continue
  239. path_ids = paths_by_content_id.get(platform_content_id, [])
  240. reject_records.append(
  241. _with_v4_explanation(
  242. {
  243. "decision_target_id": platform_content_id,
  244. "policy_run_id": policy_run_id,
  245. "main_decision_reason_code": decision["decision_reason_code"],
  246. "decision_id": decision["decision_id"],
  247. "source_path_record_ids": path_ids,
  248. "media_snapshot": _media_snapshot(media_by_platform_content_id.get(platform_content_id)),
  249. "source_evidence_ref": _source_evidence_ref(decision, path_ids),
  250. },
  251. decision,
  252. )
  253. )
  254. return reject_records
  255. def _build_technical_retry_records(
  256. policy_run_id: str,
  257. discovered_content_items: list[dict[str, Any]],
  258. decision_by_target_id: dict[str, dict[str, Any]],
  259. media_by_platform_content_id: dict[str, dict[str, Any]],
  260. paths_by_content_id: dict[str, list[str]],
  261. ) -> list[dict[str, Any]]:
  262. retry_records: list[dict[str, Any]] = []
  263. for item in discovered_content_items:
  264. platform_content_id = item["platform_content_id"]
  265. decision = decision_by_target_id[platform_content_id]
  266. if decision["decision_action"] != "TECHNICAL_RETRY_REQUIRED":
  267. continue
  268. path_ids = paths_by_content_id.get(platform_content_id, [])
  269. retry_records.append(
  270. _with_v4_explanation(
  271. {
  272. "decision_target_id": platform_content_id,
  273. "policy_run_id": policy_run_id,
  274. "main_decision_reason_code": decision["decision_reason_code"],
  275. "decision_id": decision["decision_id"],
  276. "technical_retry_status": "retry_required",
  277. "source_path_record_ids": path_ids,
  278. "media_snapshot": _media_snapshot(media_by_platform_content_id.get(platform_content_id)),
  279. "source_evidence_ref": _source_evidence_ref(decision, path_ids),
  280. },
  281. decision,
  282. )
  283. )
  284. return retry_records
  285. def _media_snapshot(media: dict[str, Any] | None) -> dict[str, Any]:
  286. if not media:
  287. return {}
  288. raw_payload = media.get("raw_payload") if isinstance(media.get("raw_payload"), dict) else {}
  289. snapshot = {
  290. "content_media_status": media.get("content_media_status"),
  291. "oss_url": media.get("oss_url"),
  292. "local_path": media.get("local_path"),
  293. "play_url": media.get("play_url"),
  294. "oss_object_key": raw_payload.get("oss_object_key"),
  295. "upload_failure_reason": raw_payload.get("upload_failure_reason") or raw_payload.get("failure_reason"),
  296. }
  297. return {key: value for key, value in snapshot.items() if value is not None}
  298. def _source_evidence_ref(
  299. decision: dict[str, Any],
  300. source_path_record_ids: list[str] | None = None,
  301. ) -> dict[str, Any]:
  302. ref = {
  303. "decision_id": decision.get("decision_id"),
  304. "decision_target_id": decision.get("decision_target_id"),
  305. }
  306. if source_path_record_ids is not None:
  307. ref["source_path_record_ids"] = source_path_record_ids
  308. return {key: value for key, value in ref.items() if value is not None}
  309. def _with_v4_explanation(
  310. record: dict[str, Any],
  311. decision: dict[str, Any],
  312. ) -> dict[str, Any]:
  313. explanation = _v4_decision_explanation(decision)
  314. if explanation:
  315. record["v4_explanation"] = explanation
  316. return record
  317. def _v4_explanation_field(decision: dict[str, Any]) -> dict[str, Any]:
  318. explanation = _v4_decision_explanation(decision)
  319. return {"v4_explanation": explanation} if explanation else {}
  320. def _v4_decision_explanation(decision: dict[str, Any]) -> dict[str, Any]:
  321. scorecard = decision.get("scorecard") or {}
  322. if scorecard.get("schema_version") != "v4_scorecard.v1":
  323. return {}
  324. replay_data = decision.get("decision_replay_data") or {}
  325. explanation = {
  326. "schema_version": "v4_decision_explanation.v1",
  327. "scorecard_schema_version": scorecard.get("schema_version"),
  328. "query_relevance_score": scorecard.get("query_relevance_score"),
  329. "platform_performance_score": scorecard.get("platform_performance_score"),
  330. "score": decision.get("score"),
  331. "missing_observable_fields": scorecard.get("missing_observable_fields", []),
  332. "decision_action": decision.get("decision_action"),
  333. "decision_reason_code": decision.get("decision_reason_code"),
  334. "search_query_effect_status": decision.get("search_query_effect_status"),
  335. "allow_walk": replay_data.get("allow_walk"),
  336. "allow_walk_reason": replay_data.get("allow_walk_reason"),
  337. }
  338. for optional_field in [
  339. "score_missing",
  340. "failure_type",
  341. "exception_type",
  342. "error_message",
  343. "http_status_code",
  344. "retry_count",
  345. "final_status",
  346. ]:
  347. if optional_field in scorecard:
  348. explanation[optional_field] = scorecard[optional_field]
  349. elif optional_field in replay_data:
  350. explanation[optional_field] = replay_data[optional_field]
  351. return explanation
  352. def _build_publish_jobs(
  353. run_id: str,
  354. policy_run_id: str,
  355. content_assets: list[dict[str, Any]],
  356. ) -> list[dict[str, Any]]:
  357. jobs: list[dict[str, Any]] = []
  358. for asset in content_assets:
  359. publish_job_id = _stable_id(
  360. "publish_job",
  361. run_id,
  362. policy_run_id,
  363. asset["platform_content_id"],
  364. asset["decision_id"],
  365. )
  366. jobs.append(
  367. {
  368. "publish_job_id": publish_job_id,
  369. "platform_content_id": asset["platform_content_id"],
  370. "job_status": "created",
  371. "trigger_mode": "manual_review",
  372. "request_payload": {
  373. "run_id": run_id,
  374. "policy_run_id": policy_run_id,
  375. "content_asset": asset,
  376. "decision_id": asset["decision_id"],
  377. "source_path_record_ids": asset["source_path_record_ids"],
  378. },
  379. "response_payload": {},
  380. }
  381. )
  382. return jobs
  383. def _build_author_assets(
  384. run_id: str,
  385. policy_run_id: str,
  386. discovered_content_items: list[dict[str, Any]],
  387. decision_by_target_id: dict[str, dict[str, Any]],
  388. paths_by_content_id: dict[str, list[str]],
  389. ) -> tuple[list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]]]:
  390. by_author: dict[tuple[str, str], list[dict[str, Any]]] = {}
  391. for item in discovered_content_items:
  392. author_id = item.get("platform_author_id")
  393. if not author_id:
  394. continue
  395. platform = item.get("platform", "douyin")
  396. by_author.setdefault((platform, author_id), []).append(item)
  397. summaries: list[dict[str, Any]] = []
  398. asset_rows: list[dict[str, Any]] = []
  399. role_rows: list[dict[str, Any]] = []
  400. for (platform, author_id), items in sorted(by_author.items()):
  401. decisions = [
  402. decision_by_target_id[item["platform_content_id"]]
  403. for item in items
  404. if item["platform_content_id"] in decision_by_target_id
  405. ]
  406. qualified_decisions = [
  407. decision for decision in decisions if decision["decision_action"] == "ADD_TO_CONTENT_POOL"
  408. ]
  409. sample_count = len(items)
  410. qualified_count = len(qualified_decisions)
  411. qualified_ratio = qualified_count / sample_count if sample_count else 0
  412. # 优质老年作者入池门:纯 50+ 画像(来自 fifty_plus,仅抖音有)占比>30% 且 TGI>120。
  413. elderly_ratio, elderly_tgi = _best_fifty_plus(decisions)
  414. if not _author_asset_eligible(elderly_ratio, elderly_tgi):
  415. continue
  416. author_asset_id = _stable_id("author_asset", platform, author_id)
  417. source_path_record_ids = sorted(
  418. {
  419. path_id
  420. for item in items
  421. for path_id in paths_by_content_id.get(item["platform_content_id"], [])
  422. }
  423. )
  424. decision_ids = [decision["decision_id"] for decision in decisions]
  425. tags = sorted({tag for item in items for tag in item.get("tags", [])})
  426. display_name = next((item.get("author_display_name") for item in items if item.get("author_display_name")), "")
  427. source_type = (
  428. "new_discovery"
  429. if any(item.get("previous_discovery_step") == "author_work" for item in items)
  430. else "new_discovery"
  431. )
  432. evidence_refs = {
  433. "decision_ids": decision_ids,
  434. "content_discovery_ids": [item["content_discovery_id"] for item in items],
  435. "source_path_record_ids": source_path_record_ids,
  436. }
  437. profile_snapshot = {
  438. "sample_count": sample_count,
  439. "qualified_content_count": qualified_count,
  440. "qualified_content_ratio": qualified_ratio,
  441. "elderly_ratio": elderly_ratio,
  442. "elderly_tgi": elderly_tgi,
  443. }
  444. summary = {
  445. "author_asset_id": author_asset_id,
  446. "platform": platform,
  447. "platform_author_id": author_id,
  448. "author_display_name": display_name,
  449. "asset_status": "active",
  450. "roles": ["author_asset", "source_seed", "high_50plus_profile"],
  451. "eligible_as_source": True,
  452. "source_path_record_ids": source_path_record_ids,
  453. "decision_ids": decision_ids,
  454. "evidence_refs": evidence_refs,
  455. }
  456. summaries.append(summary)
  457. asset_rows.append(
  458. {
  459. "author_asset_id": author_asset_id,
  460. "platform": platform,
  461. "platform_author_id": author_id,
  462. "author_display_name": display_name,
  463. "author_profile_url": None,
  464. "asset_status": "active",
  465. "source_type": source_type,
  466. "validation_status": "rule_validated",
  467. "eligible_as_source": 1,
  468. "elderly_ratio": elderly_ratio,
  469. "elderly_tgi": elderly_tgi,
  470. "content_tags": tags,
  471. "source_run_id": run_id,
  472. "source_policy_run_id": policy_run_id,
  473. "profile_snapshot": profile_snapshot,
  474. "evidence_refs": evidence_refs,
  475. "raw_payload": {
  476. "final_output_summary": summary,
  477. "profile_snapshot": profile_snapshot,
  478. },
  479. }
  480. )
  481. role_rows.extend(
  482. {
  483. "author_asset_id": author_asset_id,
  484. "role": role,
  485. "role_status": "active",
  486. "role_reason_code": "elderly_50plus_profile_pass",
  487. "assigned_by": "system",
  488. "source_run_id": run_id,
  489. "raw_payload": {"evidence_refs": evidence_refs},
  490. }
  491. for role in summary["roles"]
  492. )
  493. return summaries, asset_rows, role_rows
  494. def _author_asset_eligible(elderly_ratio: float | None, elderly_tgi: float | None) -> bool:
  495. """优质老年作者入池门:纯 50+ 占比 > 30% 且 50+ TGI > 120(仅抖音有画像,天然只抖音入池)。"""
  496. return (
  497. elderly_ratio is not None
  498. and elderly_tgi is not None
  499. and elderly_ratio > 30
  500. and elderly_tgi > 120
  501. )
  502. def _best_fifty_plus(decisions: list[dict[str, Any]]) -> tuple[float | None, float | None]:
  503. """从作者各 decision 的 scorecard 取 50+ 画像(status=ok),按 50+占比取最强的一条。
  504. 返回 (band_percent=纯50+占比, band_tgi=纯50+TGI);无 ok 画像 → (None, None)。
  505. """
  506. best: tuple[float, float] | None = None
  507. for decision in decisions:
  508. scorecard = decision.get("scorecard")
  509. if not isinstance(scorecard, dict) or scorecard.get("fifty_plus_status") != "ok":
  510. continue
  511. components = scorecard.get("fifty_plus_components")
  512. if not isinstance(components, dict):
  513. continue
  514. pct = components.get("band_percent")
  515. tgi = components.get("band_tgi")
  516. if not isinstance(pct, (int, float)) or not isinstance(tgi, (int, float)):
  517. continue
  518. if best is None or pct > best[0]:
  519. best = (float(pct), float(tgi))
  520. return best if best is not None else (None, None)
  521. def _best_age_level(decisions: list[dict[str, Any]]) -> str:
  522. priority = {"strong": 3, "medium": 2, "weak": 1, "missing": 0}
  523. levels = [decision.get("age_50_plus_level", "missing") for decision in decisions]
  524. return max(levels or ["missing"], key=lambda level: priority.get(level, 0))
  525. def _paths_by_content_id(source_path_records: list[dict[str, Any]]) -> dict[str, list[str]]:
  526. pattern_path_records = [
  527. path for path in source_path_records if path["source_path_type"] == "pattern_to_search_query"
  528. ]
  529. by_search_query_id = {
  530. path["to_node_id"]: path["source_path_record_id"] for path in pattern_path_records
  531. }
  532. result: dict[str, list[str]] = {}
  533. for path in source_path_records:
  534. if path["source_path_type"] == "search_query_to_content":
  535. result.setdefault(path["to_node_id"], [])
  536. pattern_source_path_record_id = by_search_query_id.get(path["from_node_id"])
  537. if pattern_source_path_record_id:
  538. result[path["to_node_id"]].append(pattern_source_path_record_id)
  539. result[path["to_node_id"]].append(path["source_path_record_id"])
  540. elif path["source_path_type"] == "decision_to_asset":
  541. result.setdefault(path["to_node_id"], []).append(path["source_path_record_id"])
  542. return result
  543. def _stable_id(prefix: str, *parts: Any) -> str:
  544. raw = ":".join(str(part) for part in parts)
  545. return f"{prefix}_{hashlib.sha1(raw.encode('utf-8')).hexdigest()[:16]}"