run_service.py 33 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816
  1. from __future__ import annotations
  2. import os
  3. import socket
  4. import threading
  5. from collections import Counter
  6. from datetime import datetime, timezone
  7. from pathlib import Path
  8. from typing import Any
  9. from uuid import uuid4
  10. from content_agent.constants import RUNTIME_SCHEMA_VERSION
  11. from content_agent.business_modules import run_record
  12. from content_agent.errors import ContentAgentError, ErrorCode, error_from_exception
  13. from content_agent.graph import RunDependencies, build_run_graph
  14. from content_agent.integrations.composite_runtime import CompositeRuntimeStore
  15. from content_agent.integrations.database_runtime import ContentSupplyDbConfig, DatabaseRuntimeStore
  16. from content_agent.integrations.demand_source import DemandSourceService
  17. from content_agent.integrations.douyin import CrawapiDouyinClient
  18. from content_agent.integrations.gemini_video import GeminiVideoClient as RealGeminiVideoClient
  19. from content_agent.integrations.qwen_video import QwenVideoClient
  20. from content_agent.integrations.kuaishou import CrawapiKuaishouClient
  21. from content_agent.integrations.mock_platform import MockPlatformClient
  22. from content_agent.integrations import oss_archive
  23. from content_agent.integrations.shipinhao import CrawapiShipinhaoClient
  24. from content_agent.integrations.policy_json import JsonPolicyBundleStore
  25. from content_agent.integrations.query_variant import (
  26. MissingQueryVariantClient,
  27. query_variant_client_from_env,
  28. )
  29. from content_agent.integrations.runtime_files import LocalRuntimeFileStore
  30. from content_agent.interfaces import (
  31. GeminiVideoClient,
  32. PlatformSearchClient,
  33. PolicyBundleStore,
  34. QueryVariantClient,
  35. RuntimeStore,
  36. )
  37. from content_agent.models import RunState
  38. from content_agent.schemas import RunStartRequest
  39. class RunService:
  40. def __init__(
  41. self,
  42. runtime_root: Path | str = Path("runtime/v1"),
  43. runtime: RuntimeStore | None = None,
  44. policy_store: PolicyBundleStore | None = None,
  45. demand_source: DemandSourceService | None = None,
  46. query_variant_client: QueryVariantClient | None = None,
  47. gemini_video_client: GeminiVideoClient | None = None,
  48. pattern_pg: Any | None = None,
  49. ) -> None:
  50. self.runtime = runtime or LocalRuntimeFileStore(runtime_root)
  51. self.policy_store = policy_store or JsonPolicyBundleStore(Path("."))
  52. self.demand_source = demand_source
  53. # M9C Gate 1:None → 真实 run 时惰性 from_env;测试可注入 fake。
  54. self.pattern_pg = pattern_pg
  55. self.query_variant_client = query_variant_client or MissingQueryVariantClient(
  56. "query variant client is not configured"
  57. )
  58. self._gemini_video_client = gemini_video_client or _DeterministicGeminiVideoClient()
  59. @classmethod
  60. def from_env(cls, runtime_root: Path | str = Path("runtime/v1")) -> "RunService":
  61. local_runtime = LocalRuntimeFileStore(runtime_root)
  62. env = _merged_project_env()
  63. query_variant_client = query_variant_client_from_env(env)
  64. gemini_video_client = _gemini_video_client_from_env(env)
  65. db_runtime_enabled = _env_enabled("CONTENT_AGENT_DB_RUNTIME_ENABLED")
  66. try:
  67. config = ContentSupplyDbConfig.from_env()
  68. except Exception as exc:
  69. if not db_runtime_enabled:
  70. return cls(
  71. runtime_root=runtime_root,
  72. runtime=local_runtime,
  73. query_variant_client=query_variant_client,
  74. gemini_video_client=gemini_video_client,
  75. )
  76. raise ContentAgentError(
  77. ErrorCode.DB_CONFIG_MISSING,
  78. "content supply db config is missing",
  79. {"exception_type": type(exc).__name__},
  80. ) from exc
  81. demand_source = DemandSourceService(config)
  82. if not db_runtime_enabled:
  83. return cls(
  84. runtime_root=runtime_root,
  85. runtime=local_runtime,
  86. demand_source=demand_source,
  87. query_variant_client=query_variant_client,
  88. gemini_video_client=gemini_video_client,
  89. )
  90. db_runtime = DatabaseRuntimeStore(config)
  91. return cls(
  92. runtime_root=runtime_root,
  93. runtime=CompositeRuntimeStore(db_runtime, local_runtime),
  94. demand_source=demand_source,
  95. query_variant_client=query_variant_client,
  96. gemini_video_client=gemini_video_client,
  97. )
  98. def start_run(self, request: RunStartRequest) -> RunState:
  99. run_id = request.run_id or f"v1_run_{uuid4().hex[:12]}"
  100. policy_run_id = f"policy_run_{run_id.removeprefix('v1_run_')}"
  101. source_ref = self._source_ref_from_request(request)
  102. initial_state: RunState = {
  103. "run_id": run_id,
  104. "policy_run_id": policy_run_id,
  105. "schema_version": RUNTIME_SCHEMA_VERSION,
  106. "platform": request.platform,
  107. "platform_mode": request.platform_mode,
  108. "source": request.source,
  109. "strategy_version": request.strategy_version,
  110. "current_step": "start",
  111. "status": "running",
  112. "errors": [],
  113. }
  114. try:
  115. self.runtime.prepare_run(run_id)
  116. self._create_run_record(run_id, request, source_ref)
  117. source = self._resolve_source(request)
  118. resolved_source_ref = self._source_ref_from_resolved_source(source, source_ref)
  119. if resolved_source_ref != source_ref:
  120. self.runtime.update_run_record(
  121. run_id,
  122. {
  123. "demand_content_id": _demand_content_id_from_source(source),
  124. "source_ref": resolved_source_ref,
  125. },
  126. )
  127. initial_state["source"] = source
  128. self._append_lifecycle_event(
  129. run_id,
  130. policy_run_id,
  131. event_id="lifecycle_start",
  132. event_type="run_started",
  133. status="running",
  134. message="run started",
  135. raw_payload={"source_ref": resolved_source_ref, "executor": _executor_identity()},
  136. )
  137. # 先验证/构造平台 client(拒绝非法平台),再过 Gate 1(不为非法平台浪费 PG 调用)。
  138. platform_client = self._platform_client(request.platform, request.platform_mode)
  139. gate1_reason = self._gate1_block_reason(request, source)
  140. if gate1_reason is not None:
  141. return self._gate1_blocked_state(initial_state, gate1_reason)
  142. deps = RunDependencies(
  143. runtime=self.runtime,
  144. platform_client=platform_client,
  145. policy_store=self.policy_store,
  146. query_variant_client=self.query_variant_client,
  147. gemini_video_client=self._gemini_video_client,
  148. )
  149. graph = build_run_graph(deps)
  150. state = graph.invoke(initial_state)
  151. self._record_success_metadata(state)
  152. self._trigger_post_run_oss_archive(state, request)
  153. return state
  154. except Exception as exc:
  155. error = self._classify_error(exc)
  156. failed_state = self._failed_state(initial_state, error)
  157. self._record_failure_metadata(
  158. failed_state,
  159. request,
  160. source_ref,
  161. error,
  162. )
  163. return failed_state
  164. def _source_ref_from_request(self, request: RunStartRequest) -> dict[str, Any]:
  165. if request.demand_content_id is not None:
  166. return {
  167. "source_type": "demand_content",
  168. "demand_content_id": request.demand_content_id,
  169. }
  170. if request.run_label:
  171. return {
  172. "source_type": "demand_content",
  173. "run_label": request.run_label,
  174. }
  175. if request.source:
  176. return {"source_type": "local_source", "source": request.source}
  177. return {
  178. "source_type": "demand_content_default",
  179. "selector": "pg_pattern_v2_passed_first",
  180. }
  181. def _resolve_source(self, request: RunStartRequest) -> str | dict[str, Any] | None:
  182. if request.demand_content_id is not None:
  183. if not self.demand_source:
  184. raise ContentAgentError(
  185. ErrorCode.DB_CONFIG_MISSING,
  186. "demand source db is not configured",
  187. {"selector": "demand_content_id"},
  188. )
  189. source = self.demand_source.get_by_id(request.demand_content_id)
  190. return source
  191. if request.run_label:
  192. if not self.demand_source:
  193. raise ContentAgentError(
  194. ErrorCode.DB_CONFIG_MISSING,
  195. "demand source db is not configured",
  196. {"selector": "run_label"},
  197. )
  198. source = self.demand_source.get_by_run_label(request.run_label)
  199. return source
  200. if request.source:
  201. return request.source
  202. if not self.demand_source:
  203. raise ContentAgentError(
  204. ErrorCode.DB_CONFIG_MISSING,
  205. "demand source db is not configured",
  206. {"selector": "default_pg_pattern_v2_passed"},
  207. )
  208. return self.demand_source.get_default_pg_pattern_source()
  209. def _source_ref_from_resolved_source(
  210. self,
  211. source: str | dict[str, Any] | None,
  212. source_ref: dict[str, Any],
  213. ) -> dict[str, Any]:
  214. if source_ref.get("source_type") != "demand_content_default" or not isinstance(
  215. source,
  216. dict,
  217. ):
  218. return source_ref
  219. demand_content_id = _demand_content_id_from_source(source)
  220. resolved = dict(source_ref)
  221. if demand_content_id is not None:
  222. resolved["demand_content_id"] = demand_content_id
  223. return resolved
  224. def _create_run_record(
  225. self,
  226. run_id: str,
  227. request: RunStartRequest,
  228. source_ref: dict[str, Any],
  229. ) -> None:
  230. self.runtime.create_run_record(
  231. {
  232. "run_id": run_id,
  233. "demand_content_id": request.demand_content_id,
  234. "run_label": request.run_label,
  235. "platform": request.platform,
  236. "platform_mode": request.platform_mode,
  237. "strategy_version": request.strategy_version,
  238. "status": "running",
  239. "current_step": "start",
  240. "source_ref": source_ref,
  241. "started_at": _utc_now(),
  242. }
  243. )
  244. def _record_success_metadata(self, state: RunState) -> None:
  245. final_status = "partial_success" if state.get("query_failures") else "success"
  246. state["status"] = final_status
  247. self.runtime.record_policy_run(_policy_run_record_from_state(state))
  248. validation = self.validate_run(state["run_id"])
  249. self._update_final_output_validation(state["run_id"], validation)
  250. self.runtime.update_run_record(
  251. state["run_id"],
  252. {
  253. "status": final_status,
  254. "current_step": state.get("current_step", "review_strategy"),
  255. "validation_status": validation["status"],
  256. "completed_at": _utc_now(),
  257. },
  258. )
  259. self._append_lifecycle_event(
  260. state["run_id"],
  261. state["policy_run_id"],
  262. event_id="lifecycle_success",
  263. event_type="run_succeeded",
  264. status=final_status,
  265. message="run succeeded"
  266. if final_status == "success"
  267. else "run partially succeeded",
  268. output_ref="final_output.json",
  269. raw_payload={
  270. "validation_status": validation["status"],
  271. "policy_bundle_id": state.get("policy_bundle_id"),
  272. "query_failures": state.get("query_failures", []),
  273. },
  274. )
  275. def _trigger_post_run_oss_archive(self, state: RunState, request: RunStartRequest) -> None:
  276. if request.platform_mode != "real":
  277. return
  278. run_id = state["run_id"]
  279. policy_run_id = state["policy_run_id"]
  280. try:
  281. records = self.runtime.read_jsonl(run_id, "content_media_records.jsonl")
  282. pending_count = sum(
  283. 1
  284. for row in records
  285. if row.get("content_media_status") == "oss_upload_pending" and row.get("play_url")
  286. )
  287. if pending_count <= 0:
  288. return
  289. self._append_lifecycle_event(
  290. run_id,
  291. policy_run_id,
  292. event_id="oss_archive_post_run_started",
  293. event_type="oss_archive_post_run",
  294. status="running",
  295. message="post-run oss archive started",
  296. raw_payload={"pending_due_count": pending_count},
  297. )
  298. except Exception:
  299. return
  300. def _archive() -> None:
  301. try:
  302. archived = oss_archive.archive_pending_for_run(self.runtime, run_id)
  303. status_counts = Counter(row.get("content_media_status") for row in archived)
  304. self._append_lifecycle_event(
  305. run_id,
  306. policy_run_id,
  307. event_id="oss_archive_post_run_completed",
  308. event_type="oss_archive_post_run",
  309. status="success",
  310. message="post-run oss archive completed",
  311. raw_payload={
  312. "pending_due_count": pending_count,
  313. "content_media_status_counts": dict(status_counts),
  314. },
  315. )
  316. except Exception as exc: # noqa: BLE001
  317. self._append_lifecycle_event(
  318. run_id,
  319. policy_run_id,
  320. event_id="oss_archive_post_run_failed",
  321. event_type="oss_archive_post_run",
  322. status="failed",
  323. message="post-run oss archive failed",
  324. error_code="OSS_ARCHIVE_POST_RUN_FAILED",
  325. raw_payload={
  326. "pending_due_count": pending_count,
  327. "exception_type": type(exc).__name__,
  328. "error": str(exc)[:500],
  329. },
  330. )
  331. self._start_background_thread(_archive, name=f"oss-archive-{run_id}")
  332. def _start_background_thread(self, target: Any, *, name: str) -> None:
  333. thread = threading.Thread(target=target, name=name, daemon=True)
  334. thread.start()
  335. def _update_final_output_validation(self, run_id: str, validation: dict[str, Any]) -> None:
  336. final_output = self.runtime.read_json(run_id, "final_output.json")
  337. validation_status = validation["status"]
  338. findings = validation.get("findings", [])
  339. final_output["validation_status"] = validation_status
  340. summary = final_output.setdefault("summary", {})
  341. summary["run_path_complete"] = validation_status == "pass"
  342. summary["trace_complete"] = validation_status == "pass"
  343. summary["validation_findings_summary"] = [
  344. finding.get("message", str(finding)) if isinstance(finding, dict) else str(finding)
  345. for finding in findings
  346. ]
  347. self.runtime.update_json(run_id, "final_output.json", final_output)
  348. def _record_failure_metadata(
  349. self,
  350. state: RunState,
  351. request: RunStartRequest,
  352. source_ref: dict[str, Any],
  353. error: ContentAgentError,
  354. ) -> None:
  355. run_id = state["run_id"]
  356. policy_run_id = state["policy_run_id"]
  357. try:
  358. self.runtime.update_run_record(
  359. run_id,
  360. {
  361. "status": "failed",
  362. "current_step": state.get("current_step", "failed"),
  363. "error_code": error.error_code.value,
  364. "error_message": error.message,
  365. "error_detail": error.detail,
  366. "completed_at": _utc_now(),
  367. },
  368. )
  369. self._record_platform_query_failure_details(run_id, policy_run_id, error)
  370. self._append_lifecycle_event(
  371. run_id,
  372. policy_run_id,
  373. event_id="lifecycle_failed",
  374. event_type="run_failed",
  375. status="failed",
  376. message=error.message,
  377. error_code=error.error_code.value,
  378. raw_payload={
  379. "error_detail": error.detail,
  380. "source_ref": source_ref,
  381. "platform_mode": request.platform_mode,
  382. },
  383. )
  384. except Exception:
  385. # Preserve the original run failure; DB failure here is already reflected by the state.
  386. return
  387. def _record_platform_query_failure_details(
  388. self,
  389. run_id: str,
  390. policy_run_id: str,
  391. error: ContentAgentError,
  392. ) -> None:
  393. query_failures = _query_failures_from_error(error)
  394. if not query_failures:
  395. return
  396. try:
  397. search_queries = self.runtime.read_jsonl(run_id, "search_queries.jsonl")
  398. except FileNotFoundError:
  399. return
  400. if not search_queries:
  401. return
  402. records = run_record.build_platform_query_failure_records(
  403. run_id,
  404. policy_run_id,
  405. search_queries,
  406. query_failures,
  407. )
  408. self.runtime.append_jsonl(run_id, "search_queries.jsonl", records["search_queries"])
  409. self.runtime.append_jsonl(run_id, "search_clues.jsonl", records["search_clues"])
  410. self.runtime.append_jsonl(run_id, "run_events.jsonl", records["run_events"])
  411. def _append_lifecycle_event(
  412. self,
  413. run_id: str,
  414. policy_run_id: str,
  415. *,
  416. event_id: str,
  417. event_type: str,
  418. status: str,
  419. message: str,
  420. input_ref: str | None = None,
  421. output_ref: str | None = None,
  422. error_code: str | None = None,
  423. raw_payload: dict[str, Any] | None = None,
  424. ) -> None:
  425. self.runtime.append_run_event_records(
  426. run_id,
  427. policy_run_id,
  428. [
  429. {
  430. "event_id": event_id,
  431. "event_type": event_type,
  432. "status": status,
  433. "input_ref": input_ref,
  434. "output_ref": output_ref,
  435. "error_code": error_code,
  436. "message": message,
  437. "raw_payload": raw_payload or {},
  438. "created_at": _utc_now(),
  439. }
  440. ],
  441. )
  442. def _gate1_block_reason(self, request: RunStartRequest, source: Any) -> str | None:
  443. """M9C Gate 1:仅非抖音真实 run;需求 itemset_items 至少含一个分类树终端元素才做。
  444. 无法判定(非 dict 源 / 缺 execution_id / 缺 category_id)→ 不阻断(交给后续)。
  445. PG 不可达 → has_terminal_element 抛 ContentAgentError(被 start_run except 归类),不静默放过。
  446. """
  447. if request.platform == "douyin" or request.platform_mode != "real":
  448. return None
  449. if not isinstance(source, dict):
  450. return None
  451. evidence_pack = (source.get("ext_data") or {}).get("evidence_pack") or {}
  452. execution_id = evidence_pack.get("pattern_execution_id")
  453. category_ids = [
  454. item.get("category_id")
  455. for item in evidence_pack.get("itemset_items", [])
  456. if isinstance(item, dict) and item.get("category_id") is not None
  457. ]
  458. if execution_id is None or not category_ids:
  459. return None
  460. client = self.pattern_pg
  461. if client is None:
  462. from content_agent.integrations.pattern_pg import PatternPgClient
  463. client = PatternPgClient.from_env()
  464. if client.has_terminal_element(execution_id, category_ids):
  465. return None
  466. return "gate1_no_terminal_element"
  467. def _gate1_blocked_state(self, initial_state: RunState, reason: str) -> RunState:
  468. run_id = initial_state["run_id"]
  469. policy_run_id = initial_state["policy_run_id"]
  470. blocked: RunState = {
  471. **initial_state,
  472. "current_step": "blocked_gate1",
  473. "status": "blocked",
  474. "error_code": reason,
  475. "error_message": "需求未来自分类树终端元素,非抖音 Gate 1 不放行",
  476. "errors": [],
  477. }
  478. try:
  479. self.runtime.update_run_record(
  480. run_id,
  481. {
  482. "status": "blocked",
  483. "current_step": "blocked_gate1",
  484. "error_code": reason,
  485. "completed_at": _utc_now(),
  486. },
  487. )
  488. self._append_lifecycle_event(
  489. run_id,
  490. policy_run_id,
  491. event_id="lifecycle_blocked_gate1",
  492. event_type="run_blocked",
  493. status="blocked",
  494. message=blocked["error_message"],
  495. error_code=reason,
  496. )
  497. except Exception:
  498. pass
  499. return blocked
  500. def _failed_state(self, initial_state: RunState, error: ContentAgentError) -> RunState:
  501. return {
  502. **initial_state,
  503. "current_step": "failed",
  504. "status": "failed",
  505. "error_code": error.error_code.value,
  506. "error_message": error.message,
  507. "error_detail": error.detail,
  508. "http_status_code": error.status_code,
  509. "errors": [error.message],
  510. }
  511. def _classify_error(self, exc: Exception) -> ContentAgentError:
  512. if isinstance(exc, ContentAgentError):
  513. return exc
  514. if isinstance(exc, FileNotFoundError):
  515. return ContentAgentError(
  516. ErrorCode.INVALID_SOURCE,
  517. "source file not found",
  518. {"exception_type": type(exc).__name__},
  519. status_code=400,
  520. )
  521. if isinstance(exc, ValueError) and "unknown strategy_version" in str(exc):
  522. return ContentAgentError(
  523. ErrorCode.POLICY_BUNDLE_NOT_FOUND,
  524. "policy bundle not found",
  525. {"exception_type": type(exc).__name__},
  526. status_code=400,
  527. )
  528. if isinstance(exc, ValueError) and "evidence_pack" in str(exc):
  529. return ContentAgentError(
  530. ErrorCode.INVALID_SOURCE,
  531. "invalid source",
  532. {"exception_type": type(exc).__name__},
  533. status_code=400,
  534. )
  535. return error_from_exception(exc, detail={"exception_type": type(exc).__name__})
  536. def _platform_client(self, platform: str, platform_mode: str) -> PlatformSearchClient:
  537. if platform_mode == "mock":
  538. return MockPlatformClient()
  539. if platform_mode == "real":
  540. real_clients = {
  541. "douyin": CrawapiDouyinClient.from_env,
  542. "kuaishou": CrawapiKuaishouClient.from_env,
  543. "shipinhao": CrawapiShipinhaoClient.from_env,
  544. }
  545. builder = real_clients.get(platform)
  546. if builder is None:
  547. raise ContentAgentError(
  548. ErrorCode.INVALID_REQUEST,
  549. "unsupported real platform",
  550. {"platform": platform},
  551. status_code=400,
  552. )
  553. try:
  554. return builder()
  555. except Exception as exc:
  556. raise ContentAgentError(
  557. ErrorCode.PLATFORM_CONFIG_MISSING,
  558. "platform config missing",
  559. {"exception_type": type(exc).__name__},
  560. ) from exc
  561. raise ContentAgentError(
  562. ErrorCode.INVALID_REQUEST,
  563. "unsupported platform_mode",
  564. {"platform_mode": platform_mode},
  565. status_code=400,
  566. )
  567. def get_summary(self, run_id: str) -> dict:
  568. final_output_exists = (self.runtime.run_dir(run_id) / "final_output.json").exists()
  569. status = self._summary_status(run_id, final_output_exists)
  570. validation = self.validate_run(run_id) if final_output_exists else {"status": "fail"}
  571. policy_run_id = None
  572. if final_output_exists:
  573. policy_run_id = self.runtime.read_json(run_id, "final_output.json").get("policy_run_id")
  574. return {
  575. "run_id": run_id,
  576. "policy_run_id": policy_run_id,
  577. "status": status,
  578. "current_step": "review_strategy" if final_output_exists else "unknown",
  579. "output_dir": str(self.runtime.run_dir(run_id)),
  580. "files": self.runtime.file_status(run_id),
  581. "validation_status": validation["status"],
  582. "errors": [],
  583. }
  584. def _summary_status(self, run_id: str, final_output_exists: bool) -> str:
  585. try:
  586. run_events = self.runtime.read_jsonl(run_id, "run_events.jsonl")
  587. lifecycle_events = [
  588. event for event in run_events if str(event.get("event_id", "")).startswith("lifecycle_")
  589. ]
  590. except Exception:
  591. run_events = []
  592. lifecycle_events = []
  593. if lifecycle_events:
  594. return lifecycle_events[-1].get("status") or (
  595. "success" if final_output_exists else "failed"
  596. )
  597. if final_output_exists and any(
  598. event.get("event_type") == "platform_query_failed" for event in run_events
  599. ):
  600. return "partial_success"
  601. return "success" if final_output_exists else "failed"
  602. def read_jsonl(self, run_id: str, filename: str) -> list[dict]:
  603. return self.runtime.read_jsonl(run_id, filename)
  604. def read_json(self, run_id: str, filename: str) -> dict:
  605. return self.runtime.read_json(run_id, filename)
  606. def strategy_review(self, run_id: str) -> dict:
  607. try:
  608. return self.runtime.read_json(run_id, "strategy_review.json")
  609. except FileNotFoundError:
  610. policy_run_id = None
  611. try:
  612. policy_run_id = self.runtime.read_json(run_id, "final_output.json").get(
  613. "policy_run_id"
  614. )
  615. except FileNotFoundError:
  616. pass
  617. return {
  618. "schema_version": RUNTIME_SCHEMA_VERSION,
  619. "run_id": run_id,
  620. "policy_run_id": policy_run_id,
  621. "review_status": "not_generated",
  622. }
  623. def regenerate_strategy_review(self, run_id: str) -> dict:
  624. from content_agent.business_modules.learning_review import run
  625. final_output = self.runtime.read_json(run_id, "final_output.json")
  626. timestamp = datetime.now(timezone.utc).strftime("%Y%m%d%H%M%S%f")
  627. review_id = f"review_{final_output['policy_run_id']}_{timestamp}"
  628. return run(run_id, final_output["policy_run_id"], self.runtime, review_id=review_id)
  629. def validate_run(self, run_id: str) -> dict:
  630. from content_agent.business_modules.run_record import validate_run
  631. return validate_run(run_id, self.runtime)
  632. def _policy_run_record_from_state(state: RunState) -> dict[str, Any]:
  633. policy_bundle = state.get("policy_bundle") or {}
  634. decisions = state.get("rule_decisions") or []
  635. return {
  636. "run_id": state["run_id"],
  637. "policy_run_id": state["policy_run_id"],
  638. "run_role": "primary",
  639. "policy_bundle_id": state.get("policy_bundle_id"),
  640. "rule_pack_id": policy_bundle.get("rule_pack_id"),
  641. "strategy_id": policy_bundle.get("strategy_id"),
  642. "strategy_version": state.get("strategy_version"),
  643. "rule_pack_version": policy_bundle.get("rule_pack_version"),
  644. "walk_strategy_version": policy_bundle.get("walk_strategy_version"),
  645. "policy_bundle_hash": policy_bundle.get("policy_bundle_hash"),
  646. "strategy_source_ref": policy_bundle.get("strategy_source_ref"),
  647. "rule_pack_source_ref": policy_bundle.get("rule_pack_source_ref"),
  648. "evidence_bundle_schema_version": policy_bundle.get("evidence_bundle_schema_version"),
  649. "runtime_record_schema_version": policy_bundle.get("runtime_record_schema_version"),
  650. "status": state.get("status", "success"),
  651. "metrics": (state.get("final_output") or {}).get("summary", {}),
  652. "decision_summary": _decision_summary(decisions),
  653. "raw_payload": {
  654. "current_step": state.get("current_step"),
  655. "search_query_count": len(state.get("search_queries") or []),
  656. "discovered_content_count": len(state.get("discovered_content_items") or []),
  657. "decision_count": len(decisions),
  658. "final_output_ref": "final_output.json",
  659. "dispatch": policy_bundle.get("dispatch"),
  660. "runtime_status_contract": policy_bundle.get("runtime_status_contract", {}),
  661. "effect_status_mapping": policy_bundle.get("effect_status_mapping", []),
  662. "query_effect_aggregation": policy_bundle.get("query_effect_aggregation", []),
  663. "decision_reason_codes": policy_bundle.get("decision_reason_codes", []),
  664. "policy_bundle_hash": policy_bundle.get("policy_bundle_hash"),
  665. },
  666. }
  667. def _decision_summary(decisions: list[dict[str, Any]]) -> dict[str, Any]:
  668. return {
  669. "decision_action_counts": dict(Counter(item.get("decision_action") for item in decisions)),
  670. "decision_reason_code_counts": dict(
  671. Counter(item.get("decision_reason_code") for item in decisions)
  672. ),
  673. "effect_status_counts": dict(
  674. Counter(item.get("search_query_effect_status") for item in decisions)
  675. ),
  676. }
  677. def _query_failures_from_error(error: ContentAgentError) -> list[dict[str, Any]]:
  678. query_failures = error.detail.get("query_failures") if isinstance(error.detail, dict) else None
  679. if not isinstance(query_failures, list):
  680. return []
  681. return [failure for failure in query_failures if isinstance(failure, dict)]
  682. def _gemini_video_client_from_env(env: dict[str, str]) -> GeminiVideoClient:
  683. # 视频理解 provider 开关:dashscope=通义千问(国内,video_url 直取 oss);否则 OpenRouter/Gemini(海外)。
  684. provider = (env.get("CONTENT_AGENT_VIDEO_LLM_PROVIDER") or "openrouter").strip().lower()
  685. if provider in {"dashscope", "qwen", "bailian", "tongyi"}:
  686. return QwenVideoClient.from_env(env)
  687. return RealGeminiVideoClient.from_env(env)
  688. class _DeterministicGeminiVideoClient:
  689. """mock/默认判定 client:固定返回 V4 高相关结果,供本地/smoke 无网跑通。"""
  690. def analyze(
  691. self,
  692. content: dict[str, Any],
  693. media: dict[str, Any],
  694. source_context: dict[str, Any],
  695. ) -> dict[str, Any]:
  696. return {
  697. "schema_version": "v4_gemini_query_relevance.v1",
  698. "query_text": "deterministic query",
  699. "query_relevance_score": 80,
  700. "query_relevance_reason": "deterministic stub relevance",
  701. "final_status": "ok",
  702. "retry_count": 0,
  703. }
  704. def _env_enabled(key: str) -> bool:
  705. value = _load_project_env().get(key)
  706. if os.environ.get(key) is not None:
  707. value = os.environ[key]
  708. return (value or "").lower() in {"1", "true", "yes", "on"}
  709. def _merged_project_env() -> dict[str, str]:
  710. return {
  711. **_load_project_env(),
  712. **dict(os.environ),
  713. }
  714. def _demand_content_id_from_source(source: str | dict[str, Any] | None) -> int | None:
  715. if not isinstance(source, dict):
  716. return None
  717. value = source.get("id") or source.get("demand_content_id")
  718. if value in (None, ""):
  719. return None
  720. try:
  721. return int(value)
  722. except (TypeError, ValueError):
  723. return None
  724. def _utc_now() -> str:
  725. return datetime.now(timezone.utc).isoformat()
  726. def _proc_start_time(pid: int) -> str | None:
  727. # Linux /proc/<pid>/stat 第 22 字段 starttime(自启动以来的时钟滴答),作为 PID 复用的指纹;
  728. # comm(第 2 字段)在括号里且可能含空格,故按最后一个 ')' 切分。非 Linux/取不到 → None。
  729. try:
  730. with open(f"/proc/{pid}/stat", encoding="utf-8") as handle:
  731. content = handle.read()
  732. return content[content.rindex(")") + 2:].split()[19]
  733. except Exception:
  734. return None
  735. def _executor_identity() -> dict[str, Any]:
  736. # 记录"谁在跑这条 run":pid + 主机 + 进程启动时刻。读状态时据此判进程是否还活着(无时间阈值)。
  737. pid = os.getpid()
  738. return {"pid": pid, "host": socket.gethostname(), "proc_start": _proc_start_time(pid)}
  739. def _load_project_env(env_file: Path | str = ".env") -> dict[str, str]:
  740. path = Path(env_file)
  741. if not path.exists():
  742. return {}
  743. result: dict[str, str] = {}
  744. for line in path.read_text(encoding="utf-8").splitlines():
  745. stripped = line.strip()
  746. if not stripped or stripped.startswith("#") or "=" not in stripped:
  747. continue
  748. key, value = stripped.split("=", 1)
  749. result[key.strip()] = value.strip().strip('"').strip("'")
  750. return result