精读旗舰 Agent:sre-rca
前面 D03-D13 把 State / 图 / 专家 / Critic / 闸门 / 记忆 / 成本 / 运行时一块块拆开讲了。今天(D14)是总演练——逐文件读真实源码:从 state.py 的每个字段、builder.py 怎么连边,到 10 个节点各自的实现,把这些零件在一个 agent 上完整看一遍怎么咬合。这是本教程贴真代码最多的一天。读懂它,D15 再横扫其它 20 个 agent 就都是它的简化/变体了。
sre-rca 做什么
apps/sre-rca-agent 是平台的旗舰范本——一个 SRE 故障根因分析(Root Cause Analysis)Agent。你给它一个告警/现象("service-A 大量超时"),它自动从 trace、指标、发布记录、日志四个视角取证,综合出 Top-3 最可能的根因,并给出带证据、带教学性的分析。
它是"满配"agent:10 个节点、4 专家并行、双层 Critic 三态重试、四层记忆全用上。读懂它,其它 agent 都是它的简化/变体。
{"trace_id":"t-9f2","user_question":"service-A 大量超时,DB 连接池打满"}→ 分诊命中关键词"连接池"归为
alert_type="DB_SLOW" → 4 专家取证(trace 发现下游 DB 慢、metric 发现 conn_pool_usage_pct 打满、deploy 发现 14:02 上过线、log 发现大量 timeout)输出:
{"conclusion":{"confidence":"高","hypotheses":[{"description":"14:02 v1.7 发布把连接池 50→5 导致耗尽","risk_level":"high"}],"evidence_ids":[...]}} + 一段给人看的 summary_for_human。/agents/rca/graph 端点(graph_info(),server.py:143)把每个节点的 role/model/purpose 以 JSON 硬编码——相当于官方的"节点职责表"。今天讲的每个节点都能在那里对照。RcaState:全图共享的数据结构(先看数据)
看 RcaState 全貌(sre_rca/state.py:15-52),我按节点分组标注:
class RcaState(TypedDict, total=False): # total=False:所有字段都可缺省
# ── 输入 ──
trace_id: str
service_name: str | None
user_question: str
# ── Triage 输出 ──
alert_type: str
need_full_dag: bool
# ── Recall 输出(历史 case)──
similar_cases: list[dict[str, Any]]
# ── 4 Specialist 输出(字段独立,并行写零冲突)──
trace_finding: dict[str, Any] | None
metric_finding: dict[str, Any] | None
deploy_finding: dict[str, Any] | None
log_finding: dict[str, Any] | None
# ── RAG / Synthesizer 输出 ──
sop_context: list[dict[str, Any]]
conclusion: dict[str, Any] | None
synth_skipped: bool
synth_skip_reason: str | None
# ── Critic 输出 ──
critique_passed: bool
critique_feedback: str
critique_layer: int # 1=代码层失败 / 2=LLM层失败 / 0=通过
# ── Reducer:累加(防并发竞态)──
retry_count: Annotated[int, add] # ← 唯一一个非覆盖语义的字段
# ── Writeback / 输出 ──
writeback_skipped: bool
summary_for_human: str
| 字段组 | 谁写 | 合并语义 |
|---|---|---|
trace_id/service_name/user_question | 入口 | 只读 |
*_finding(4 个) | 4 专家并行 | 各写各字段,默认覆盖,互不打架 |
conclusion | synthesizer | 覆盖(可能被重试多次覆写) |
critique_passed/feedback/layer | critic | 覆盖 |
retry_count | critic(每轮) | 累加 Annotated[int, add] |
findings),并行合并时就会互相覆盖/丢数据。做成 trace_finding / metric_finding / ... 四个独立字段,各写各的,默认"覆盖"语义天然零冲突。而 retry_count 不一样——它要跨多轮"累计",用普通覆盖的话,重试时"1 覆盖 1"永远是 1,永远触不了顶。所以它专门声明 Annotated[int, add](呼应 D03 reducer),让 LangGraph 合并时做加法而非覆盖。并行写→拆字段,累计值→加法 reducer,这两条是并发安全的核心。builder:10 节点怎么连成一张图(看控制流)
数据布局懂了,看 build_rca_graph(sre_rca/builder.py:42)怎么把节点连起来。这是整台机器的"接线图":
builder = StateGraph(RcaState) # builder.py:73 · 以 RcaState 为共享内存
# 节点注册(工厂类的要先调工厂拿到闭包节点)
builder.add_node("triage", triage_node)
builder.add_node("recall", build_recall_node(episodic)) # 闭包注入 EpisodicStore
builder.add_node("trace_sp", trace_sp_node)
builder.add_node("metric_sp", metric_sp_node)
builder.add_node("deploy_sp", deploy_sp_node)
builder.add_node("log_sp", log_sp_node)
builder.add_node("rag", build_rag_node(semantic))
builder.add_node("synthesizer", build_synthesizer_node())
builder.add_node("critic", build_critic_for_rca())
builder.add_node("writeback", build_writeback_node(episodic))
# 边:线性入口
builder.add_edge(START, "triage")
builder.add_edge("triage", "recall")
# 4 SP 并行扇出 + 扇入到 RAG
for sp in ("trace_sp", "metric_sp", "deploy_sp", "log_sp"):
builder.add_edge("recall", sp) # recall → 每个专家(扇出)
builder.add_edge(sp, "rag") # 每个专家 → rag(扇入)
builder.add_edge("rag", "synthesizer")
builder.add_edge("synthesizer", "critic")
for sp in (...)4 条扇出 + 4 条扇入用一个 for 循环加边,而不是手写 8 行。因为 4 专家结构对称——recall 之后同时启动、都汇入 rag。LangGraph 看到"多个节点指向同一个下游"就自动并行跑(asyncio)。build_xxx_node(dep)凡是名字带 build_ 的都是工厂:先调它、把依赖(episodic/semantic)闭包进去,返回的才是真正的节点函数。triage / 4 专家没有依赖,直接是节点。synthesizer → critic 之后不是普通边,而是三态条件边(builder.py:100-116):
critic_router = build_critic_router( # builder.py:100 · 来自 toolkit(D07/D08)
pass_dest="writeback", # 通过 → 去沉淀
retry_dest="synthesizer", # 不过没触顶 → 回炉重做
end_dest=END, # 不过且触顶 → 诚实收尾
max_retries=max_critic_retries, # 默认 2
)
builder.add_conditional_edges(
"critic", critic_router,
{"writeback": "writeback", "synthesizer": "synthesizer", END: END},
)
builder.add_edge("writeback", END)
max_retries 传进去,路由的判断逻辑(读 critique_passed、比 retry_count)复用 toolkit 的 build_critic_router。业务只声明"通过去哪、重试去哪、触顶去哪",怎么判交给 toolkit——这就是 D06 讲的"渐进抽象"落到一个真 agent 上的样子。triage — 规则分诊(0 成本,⚪不调 LLM)
入口节点。用关键词字典把 user_question 归到 5 类 alert_type(sre_rca/nodes/triage.py:8-35):
ALERT_KEYWORDS: dict[str, list[str]] = { # triage.py:8
"DB_SLOW": ["DB", "数据库", "连接池", "慢查询", "锁等待", "select", "update"],
"JVM": ["GC", "OOM", "JVM", "heap", "young gen", "old gen"],
"ERROR_SPIKE": ["error", "异常", "exception", "5xx", "失败"],
"GENERIC_SLOW":["慢", "卡", "延迟", "P99", "P95", "timeout"],
"HEALTH_QUERY":["健康", "状态", "看一下", "check", "巡检"],
}
def _classify(text: str) -> str: # triage.py:17
text_lower = text.lower()
# 优先级:DB_SLOW > JVM > ERROR_SPIKE > GENERIC_SLOW > HEALTH_QUERY
for alert_type, keywords in ALERT_KEYWORDS.items(): # ← 靠 dict 插入顺序定优先级!
for kw in keywords:
if kw.lower() in text_lower:
return alert_type
return "GENERIC_SLOW" # 默认值
async def triage_node(state: dict) -> dict: # triage.py:27
user_question = state.get("user_question") or ""
alert_type = _classify(user_question)
need_full_dag = alert_type != "HEALTH_QUERY" # 纯健康查询走简版,不必全 DAG
return {"alert_type": alert_type, "need_full_dag": need_full_dag}
5 类 alert_type真实取值是 DB_SLOW / JVM / ERROR_SPIKE / GENERIC_SLOW / HEALTH_QUERY(别记成 "latency" 之类)。它决定下游专家分析时的侧重点。need_full_dag只有"健康查询"这类轻问题走简版;其余都要跑完整 4 专家 DAG。一个便宜的"提前分流"。return "GENERIC_SLOW"兜底:什么关键词都没命中时归到"泛化慢",绝不让分诊失败。_classify 从上往下遍历,命中第一个就返回。所以 ALERT_KEYWORDS 里 DB_SLOW 写在最前面 = 优先级最高。一句话同时提"连接池"和"慢",会被判成 DB_SLOW 而非 GENERIC_SLOW。改动关键词表时,调整字典的书写顺序就等于调整优先级——这种"顺序即语义"的隐式约定容易被忽略。recall — 情景记忆召回(工厂闭包 + 异常吞掉)
recall 要访问 EpisodicStore(D09 的记忆库),但 LangGraph 的节点签名是固定的 async def node(state)——没法多传参数。解法是工厂闭包(sre_rca/nodes/recall.py:10-35):
def build_recall_node(episodic: EpisodicStore, top_k: int = 3): # recall.py:10 · 工厂
async def recall_node(state: dict) -> dict: # recall.py:16 · 真节点
query = f"trace_id={state.get('trace_id')} service={state.get('service_name')} q={state.get('user_question')}"
try:
cases = await episodic.recall(query, top_k=top_k) # 向量检索相似历史 case
except Exception as e:
return {"similar_cases": [], "_recall_error": str(e)} # ← 挂了不崩,返空
return {"similar_cases": [
{"id": r.id, "text": r.text, "score": score, "metadata": r.metadata}
for r, score in cases
]}
return recall_node # 返回闭包,交给 builder
build_recall_node(episodic)外层工厂拿依赖,内层 recall_node 通过闭包"记住"了 episodic。返回的闭包签名仍是标准 async def node(state)→dict——既满足 LangGraph 约定,又拿到了依赖。这是全框架节点的通用写法。query 拼接把 trace_id / service / 问题拼成一句话去做向量检索。检索的是"以前有没有处理过长得像的故障"。top_k=3只召回最相似的 3 条,避免把太多历史塞进后面的 prompt(省 token)。except Exception 把召回失败整个吞掉,返回 {"similar_cases": []}(外加一个 _recall_error 供排查)。为什么?因为 recall 只是"锦上添花"——有历史参考更好,没有也能靠 4 专家取证下结论。不能让一个可选环节的故障(比如向量库连不上)拖垮整次根因分析。这是 D08 韧性思想在具体节点里的体现:非关键路径的失败要"降级"而非"崩溃"。4 专家并行取证(工厂配置 + 工厂内部)
核心取证环节。4 个专家从 recall 出发并行跑、汇入 rag。每个专家文件才 ~30 行,因为都用 toolkit 的 make_specialist_node。先看 metric 专家的配置(sre_rca/specialists/metric_sp.py:13-30):
metric_sp_node = make_specialist_node( # metric_sp.py:13
SpecialistConfig(
name="metric_sp",
state_field="metric_finding", # ← 写进 state 的哪个字段(并行不冲突!)
system_prompt=_PROMPT, # 从 prompts/metric-specialist-v1.md 读
fetch_fn=lambda state: get_metrics( # 先抓数据(@safe_tool_result 包过,见 observability.py:38)
service=state.get("service_name") or "unknown", window_sec=60),
fallback_finding={ # 抓失败时的降级结果
"health": "unknown", "confidence": "低",
"evidence_ids": [], "hint": "指标获取失败"},
model_tier="sonnet", # 取证是难活,用 Sonnet
)
)
state_field="metric_finding"呼应 L02——每个专家写不同字段,这是"4 专家并行零冲突"的根。trace 专家写 trace_finding,以此类推。fetch_fn先拉真实数据(这里调 get_metrics,被 @safe_tool_result 包过所以不会抛异常,D11)。拉到才喂给 LLM 分析。model_tier="sonnet"取证要理解复杂数据,是"难活",用贵的 Sonnet;而 critic 那种"挑刺"是"简单活"用 Haiku。省钱三档的分配。差异化配置就这些。真正的执行逻辑在 toolkit 工厂 make_specialist_node(ai_trust_toolkit/specialists/factory.py:84)里——4 个专家共用这套:
async def inner_node(state: dict) -> dict: # factory.py:93
tool_result = await cfg.fetch_fn(state) # ① 拉数据
if not tool_result.get("ok"): # ← 边界:抓数据失败
finding = dict(cfg.fallback_finding)
finding["error"] = tool_result.get("error", "tool failed")
return {cfg.state_field: finding} # 直接返回降级 finding,不调 LLM(省钱)
raw = tool_result["data"]
llm = get_llm(cfg.model_tier, cache_system=cfg.cache_system) # ② 调模型
response = await llm.ainvoke([{"role":"system","content":cfg.system_prompt},
{"role":"user","content":user_prompt_builder(state, raw)}])
finding = response_parser(getattr(response,"content",str(response))) # ③ 解析
finding.setdefault("evidence_ids", []) # 保证有 evidence_ids 字段(给 Critic 用)
return {cfg.state_field: finding}
inner_node.__name__ = cfg.name
return wrap_with_fallback(inner_node, # factory.py:125 · ④ 再包一层节点级兜底
FallbackConfig(state_field=cfg.state_field,
fallback_finding=dict(cfg.fallback_finding, fallback=True))) # ← 打 fallback=True
if not ...get("ok")数据源就挂了(比如 Prometheus 连不上),直接返回降级 finding、连 LLM 都不调——既省钱又快。setdefault("evidence_ids", [])强制保证输出里有 evidence_ids,因为 D07 的 Critic 要靠它抓"编造 ID"。节点之间的隐式契约。wrap_with_fallback(..., fallback=True)最外层再包 D08 的兜底闸:哪怕 LLM 异常/解析崩了,也只返回带 fallback:True 的降级 finding,绝不把异常抛出去炸整张图。SpecialistConfig(~30 行声明),加减一个专家几乎零成本(D15 会看到 security-audit 有 5 个、risk-reviewer 只有 2 个)。把"变的"做成配置、"不变的"沉进工厂——这就是让 toolkit 通用的关键手法。fallback=True 标记至关重要——"降级不崩"和"证据不足就诚实收尾"靠这个标记咬合:降级是允许的,但不能拿降级占位结果去凑"证据充分"。rag + synthesizer(拼查询 + 综合 + 双重容错)
rag(sre_rca/nodes/rag.py:19-40)用 4 专家的线索拼查询、检索运维手册 SOP:
async def rag_node(state: dict) -> dict: # rag.py:19
query_parts = []
for f_name in ("trace_finding", "metric_finding", "deploy_finding", "log_finding"):
f = state.get(f_name)
if f and isinstance(f, dict) and not f.get("fallback"): # ← 跳过降级的假 finding
hint = f.get("hint") or f.get("hypothesis") or ""
if hint:
query_parts.append(str(hint))
if not query_parts:
return {"sop_context": []} # ← 边界:一个有效线索都没有,直接返空
query = " | ".join(query_parts)
try:
sops = await semantic.lookup(query, record_type="sop", top_k=top_k)
except Exception as e:
return {"sop_context": [], "_rag_error": str(e)} # 检索挂了也不崩
return {"sop_context": [{"id": r.id, "text": r.text, "score": score, "metadata": r.metadata}
for r, score in sops]}
synthesizer(sre_rca/nodes/synthesizer.py)是全图最"贵"的一步——调 Sonnet 综合出 Top-3。它被两层保护包着:
async def _inner_synthesizer(state: dict) -> dict: # synthesizer.py:15 · 真正调 LLM
llm = get_llm("sonnet")
user_msg = (f"trace_id: {state.get('trace_id')}\n... "
f"trace_finding: {json.dumps(state.get('trace_finding'), ensure_ascii=False)}\n"
f"... --- Critic 上一轮反馈(如有)---\n{state.get('critique_feedback', '(无)')}\n"
f"请按 system prompt 的 schema 输出 conclusion JSON。")
response = await llm.ainvoke([{"role":"system","content":_PROMPT},
{"role":"user","content":user_msg}])
return {"conclusion": _parse_conclusion(getattr(response,"content",str(response)))}
def build_synthesizer_node(): # synthesizer.py:49
return build_synth_safeguard( # ← D08 闸②:包一层"证据不足开关"
_inner_synthesizer,
min_real_specialists=1, # 至少 1 个 SP 有真实数据才调 LLM
)
not f.get("fallback")rag 拼查询时跳过降级的假 finding——用假线索去检索 SOP 只会拿到不相关的手册,不如不拼。critique_feedbacksynthesizer 会把 Critic 上一轮的意见拼进 prompt。所以"回炉重做"时它是带着"你哪里错了"重写,而不是瞎重试。build_synth_safeguard(min_real_specialists=1)外层护栏:真实专家结果不够(全是 fallback)时跳过 LLM,直接返回"信息不足"。省下最贵的 Sonnet 调用,也避免拿垃圾数据硬编结论。_parse_conclusion(synthesizer.py:60)是三级容错解析——和 D07 Critic 的解析器同款套路:先裸 JSON → 再 ```json 代码块 → 最后兜底一个 info_sufficient:False 的 dict。
```json 包起来、有时前面还啰嗦两句。如果只 json.loads() 一次,稍不规范就抛异常、整次综合报废。所以 _parse_conclusion 三级兜底,最后哪怕全解析失败也返回一个结构完整、info_sufficient:False 的 dict(把原文塞进 _raw_text 供排查)——宁可老实说"没解析出来",也不让流程崩。跟不确定的 LLM 打交道,这种防御性解析是标配。critic + writeback(质检 + 守门沉淀)
critic(sre_rca/nodes/critic.py:12-30)几乎不写业务——直接复用 D07 的 build_critic_node,只填配置:
def build_critic_for_rca(): # critic.py:12
prompt = (Path(__file__).resolve().parents[2] / "prompts" / "critic-v1.md").read_text("utf-8")
return build_critic_node(CriticConfig( # ← 来自 toolkit(D07 双层 Critic)
system_prompt=prompt,
llm_factory=lambda: get_llm("haiku", cache_system=True), # ← Haiku 省80% + prompt缓存再省
conclusion_field="conclusion",
soft_fail=True, # LLM 挂了当"不通过"而非崩图
))
get_llm("haiku", cache_system=True)Critic 是"挑刺"简单活用便宜的 Haiku;cache_system=True 让那段固定的 critic prompt 走 prompt 缓存,多次 invoke 省 ~90% input 费用。soft_fail=TrueHaiku 调用本身失败(超时/限流)时不抛异常,而是当成"审查不通过"走重试。宁可多审一轮,也不让整图崩。业务代码 ~10 行L1 四项检查、L2 短路、三级解析全在 toolkit 里。sre-rca 只提供 prompt 和几个开关——这就是 toolkit 的价值。writeback(sre_rca/nodes/writeback.py:20-59)先守门再沉淀:
async def writeback_node(state: dict) -> dict: # writeback.py:20
if not should_writeback(state): # ← D08 闸④:仅"质检过 && 信息充分"才写
return {"writeback_skipped": True, "summary_for_human": _format_summary(state)}
conclusion = state.get("conclusion") or {}
case_id = state.get("trace_id") or "case-unknown"
text = json.dumps({...hypotheses..., "summary": conclusion.get("summary_for_human")}, ensure_ascii=False)
try:
await episodic.writeback(case_id=case_id, text=text, metadata={...}) # 沉淀进记忆库
except Exception as e:
return {"writeback_skipped": True, "_writeback_error": str(e),
"summary_for_human": _format_summary(state)} # 写失败也不崩
return {"writeback_skipped": False, "summary_for_human": _format_summary(state)}
should_writeback(state)守门:只有 critique_passed 且 info_sufficient 才写进记忆。没通过质检的坏结论绝不入库,否则会污染 recall。_format_summary无论写不写记忆,都生成给人看的 summary_for_human(writeback.py:64)。信息不足时它返回"⚠️ 证据不足,建议人工排查"。should_writeback 守门保证"没通过质检的"不会污染记忆。这就是 Day 09 说的"越用越聪明又不学坏"——一头(writeback)把关、一头(recall)取用,两端咬合成闭环。prompt 也是工程文档
别忽略 prompts/*.md——它们是本仓最好的 prompt 工程教材。以 prompts/critic-v1.md 为例(前面 critic.py 就是读它),它不只是一段 prompt,而是包含:
- 明确的模型指定:文件开头写"用 Claude Haiku 4.5 跑(spec 模型分层策略)"——prompt 和模型档位绑定记录在案。
- 6 项检查(
critic-v1.md"6 项检查"小节):证据 ID 真实性 / 置信度校准 / 语言风格铁律 / 矛盾处理 / 动作安全性 / 教学性,每项都规定了判定标准。 - 代码层预检说明:明确写"JSON 合法/必填字段/类型这些硬错误走代码不走 LLM"(呼应 D07 的 L1),Critic 只做"模型才能判断的事"。
- 三态决策:PASS / PASS_WITH_NOTES / FAIL,各自的判定条件与反馈格式。
- few-shot 示例:给出 FAIL / PASS 的完整样例,稳住 LLM 输出格式。
synthesizer-v1.md 更狠:多步思考 + 语言铁律(禁绝对措辞、必引 ≥2 专家证据、教学性、失败开关)+ 一个完整的分析范例。
-v1 支持回滚和灰度。把 prompt 当一等工程产物对待,而不是散落在代码里的字符串,这是这个仓最值得学的工程习惯之一。今日小结 + 动手
🧠 今天你应该能回答
- RcaState 的三类字段为什么用三种合并语义?(单写覆盖 / 并行拆字段 / retry_count 累加 reducer)
- builder 里 4 专家的并行是怎么用一个 for 循环连出来的?三态条件边为什么复用 toolkit?
- triage 为什么有 LLM 版 prompt 却用规则实现?关键词优先级靠什么定?(dict 插入顺序)
- recall 为什么用工厂闭包?记忆挂了为什么返空不抛异常?
- 专家为什么用工厂、fetch 失败/LLM 异常各在哪层兜底?
fallback=True和 Critic 怎么咬合? - synthesizer 的证据不足开关、三级解析容错各解决什么?
- critic 用 Haiku+缓存、writeback 用 should_writeback 守门,怎么形成"越用越准"闭环?
- prompt 文件为什么是"工程文档"?版本号有什么用?
✋ 动手:逐文件读一遍(今天的核心作业)
# 数据结构 + 图骨架
cat apps/sre-rca-agent/sre_rca/state.py
sed -n '42,118p' apps/sre-rca-agent/sre_rca/builder.py
# 6 个 nodes
cat apps/sre-rca-agent/sre_rca/nodes/triage.py
cat apps/sre-rca-agent/sre_rca/nodes/recall.py
cat apps/sre-rca-agent/sre_rca/nodes/rag.py
cat apps/sre-rca-agent/sre_rca/nodes/synthesizer.py
cat apps/sre-rca-agent/sre_rca/nodes/critic.py
cat apps/sre-rca-agent/sre_rca/nodes/writeback.py
# 一个专家(配置)+ toolkit 工厂(内部)
cat apps/sre-rca-agent/sre_rca/specialists/metric_sp.py
sed -n '84,131p' packages/ai-trust-toolkit/src/ai_trust_toolkit/specialists/factory.py
# 数据源工具(@safe_tool_result)+ prompt 工程文档
sed -n '1,55p' apps/sre-rca-agent/sre_rca/tools/observability.py
sed -n '1,60p' apps/sre-rca-agent/prompts/critic-v1.md
# 跑测试 + 跑评测
uv run pytest apps/sre-rca-agent -q
uv run python apps/sre-rca-agent/scripts/run_eval.py