| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127 |
- from __future__ import annotations
- from dataclasses import dataclass
- from typing import Any
- from langgraph.graph import END, START, StateGraph
- from content_agent.business_modules import (
- candidate_evidence,
- learning_review,
- platform_access,
- policy_version,
- result_source_lookup,
- rule_judgment,
- run_record,
- search_intent,
- source_seed,
- walk_strategy,
- )
- from content_agent.interfaces import PlatformSearchClient, PolicyBundleStore, RuntimeFileStore
- from content_agent.models import RunState
- @dataclass(frozen=True)
- class RunDependencies:
- runtime: RuntimeFileStore
- platform_client: PlatformSearchClient
- policy_store: PolicyBundleStore
- def build_run_graph(deps: RunDependencies):
- graph = StateGraph(RunState)
- def load_source(state: RunState) -> dict[str, Any]:
- result = source_seed.run(state["trace_id"], state.get("source"), deps.runtime)
- return {**result, "current_step": "load_source"}
- def plan_queries(state: RunState) -> dict[str, Any]:
- queries = search_intent.run(state["trace_id"], state["pattern_seed_pack"], deps.runtime)
- return {"queries": queries, "current_step": "plan_queries"}
- def search_platform(state: RunState) -> dict[str, Any]:
- results = platform_access.run(state["queries"], deps.platform_client)
- return {"platform_results": results, "current_step": "search_platform"}
- def build_candidates(state: RunState) -> dict[str, Any]:
- result = candidate_evidence.run(
- state["trace_id"],
- state["platform_results"],
- state["source_context"],
- deps.runtime,
- )
- return {**result, "current_step": "build_candidates"}
- def load_policy(state: RunState) -> dict[str, Any]:
- bundle = policy_version.run(state["policy_bundle_version"], deps.policy_store)
- return {"policy_bundle": bundle, "current_step": "load_policy"}
- def evaluate_rules(state: RunState) -> dict[str, Any]:
- decisions = rule_judgment.run(
- state["trace_id"],
- state["evidence_bundles"],
- state["policy_bundle"],
- deps.runtime,
- )
- return {"rule_decisions": decisions, "current_step": "evaluate_rules"}
- def plan_walk(state: RunState) -> dict[str, Any]:
- result = walk_strategy.run(
- state["pattern_seed_pack"],
- state["queries"],
- state["candidates"],
- state["rule_decisions"],
- )
- return {**result, "current_step": "plan_walk"}
- def record_run(state: RunState) -> dict[str, Any]:
- result = run_record.run(
- state["trace_id"],
- state["queries"],
- state["candidates"],
- state["rule_decisions"],
- state["source_edge_basis"],
- deps.runtime,
- )
- return {**result, "current_step": "record_run"}
- def commit_results(state: RunState) -> dict[str, Any]:
- final_output = result_source_lookup.run(
- state["trace_id"],
- state["candidates"],
- state["media_assets"],
- state["rule_decisions"],
- state["source_edges"],
- state["search_clues"],
- deps.runtime,
- )
- return {"final_output": final_output, "current_step": "commit_results"}
- def review_strategy(state: RunState) -> dict[str, Any]:
- review = learning_review.run(state["trace_id"], deps.runtime)
- return {"strategy_review": review, "current_step": "review_strategy", "status": "success"}
- graph.add_node("load_source", load_source)
- graph.add_node("plan_queries", plan_queries)
- graph.add_node("search_platform", search_platform)
- graph.add_node("build_candidates", build_candidates)
- graph.add_node("load_policy", load_policy)
- graph.add_node("evaluate_rules", evaluate_rules)
- graph.add_node("plan_walk", plan_walk)
- graph.add_node("record_run", record_run)
- graph.add_node("commit_results", commit_results)
- graph.add_node("review_strategy", review_strategy)
- graph.add_edge(START, "load_source")
- graph.add_edge("load_source", "plan_queries")
- graph.add_edge("plan_queries", "search_platform")
- graph.add_edge("search_platform", "build_candidates")
- graph.add_edge("build_candidates", "load_policy")
- graph.add_edge("load_policy", "evaluate_rules")
- graph.add_edge("evaluate_rules", "plan_walk")
- graph.add_edge("plan_walk", "record_run")
- graph.add_edge("record_run", "commit_results")
- graph.add_edge("commit_results", "review_strategy")
- graph.add_edge("review_strategy", END)
- return graph.compile()
|