"""流程图结构化抽取 + 校验 的 SSE 端点。 第三步「方案总结」之后,要把报告 / 各席位交付的层级结论抽成一张**流程图**(节点 + 连线)。 之前让总结智能体在同一对话线程里顺手吐 ``flow-json`` 很不稳:业务链 ``coordinator_prompt`` 常写着「输出结构化结论 / 按阶段产出报告」,模型被线程上下文锚定,要么不出 flow-json、要么 结构残缺(只剩两三个节点)。 这里改成**独立的一次性模型调用**(不走 lead_agent / 不带线程历史),两段式: ① 抽取(``model.astream``):把报告/席位交付按层级抽成 flow-json,**流式**把过程吐给前端实时显示; ② 校验/修复(``model.ainvoke``):把抽取结果连同源材料再过一遍模型,补全漏层、修正连边、 去掉阶段/汇总节点,产出干净的 flow-json;后端再做容错解析 + 结构校验,返回最终图数据。 SSE 帧(``data: {json}\n\n``): - ``{"event":"extract_delta","content":"…"}`` 抽取流式增量(前端把 JSON 实时显示出来) - ``{"event":"extract_done"}`` 抽取结束,进入校验 - ``{"event":"validating"}`` 校验中 - ``{"event":"final","graph":{title,nodes,edges}}`` 校验后的最终图(前端据此绘制) - ``{"event":"error","detail":"…"}`` """ from __future__ import annotations import json import logging import re from collections.abc import AsyncIterator from typing import Any from fastapi import APIRouter, Body, Request from fastapi.responses import StreamingResponse from langchain_core.messages import HumanMessage, SystemMessage from pydantic import BaseModel, Field from app.gateway.roundtable_model_fallback import fallback_model_chain, reason_text from app.gateway.routers._roundtable_seed import ( STRUCTURE_AGENT_ID, ensure_roundtable_functional_agents, ) from deerflow.agents.roundtable_orchestrator.extraction_spec import ( build_structure_directive, chain_level_rank, is_single_parent_chain, ) from deerflow.config import get_app_config from deerflow.config.agents_config import load_agent_config, load_agent_soul from deerflow.models.factory import create_chat_model from deerflow.skills.storage import get_or_new_skill_storage logger = logging.getLogger(__name__) router = APIRouter(prefix="/api/flow-extract", tags=["flow-extract"]) # ── schemas ────────────────────────────────────────────────────────────── class SeatDelivery(BaseModel): sender: str = "" content: str = "" class FlowExtractRequest(BaseModel): # 数据来源:有报告优先用报告,否则用各席位交付。 report_md: str | None = Field(default=None) seat_deliveries: list[SeatDelivery] = Field(default_factory=list) # 层级骨架:优先 coordinator_prompt 里写明的层级,否则用 chain_levels。 chain_levels: list[str] = Field(default_factory=list) coordinator_prompt: str | None = Field(default=None) # 深链接业务码(rwfx→6BF 等)。非空 → 按该业务的「业务链抽取规范」**强制**注入完整层级 # (覆盖默认/示例层级假设),根治"无论什么业务都按目的→行为体抽"。普通会商为空 → 不注入。 business_code: str | None = Field(default=None) # 根节点 label 强制值(xdfx:本次「行动名称」)。非空 → 根节点 n1 的 label 必须用它,不要用任务名当根。 root_label: str | None = Field(default=None) model: str | None = Field(default=None) # Phase 2 问答修改:带上当前图 + 用户指令 → 进 edit 模式(按指令改图,其余不变)。 current_graph: dict[str, Any] | None = Field(default=None) instruction: str | None = Field(default=None) class FlowIntentRequest(BaseModel): """第三步输入框意图识别:判断用户想 重做报告 / 改流程图 / 重画流程图 / 提问 / 模糊待澄清。""" message: str = "" has_report: bool = False has_flow: bool = False # 当前是否可入库(结构化抽取智能体已判定有适配本场景的入库技能、入库按钮已出现且未入库)。 # 为真时才允许识别出 ingest 意图——否则「帮我入库」回退为 qa/unclear。 can_ingest: bool = False model: str | None = Field(default=None) # 意图枚举(与前端路由一一对应)。ingest 仅在 can_ingest=true 时可用。 _INTENTS = {"regenerate_report", "modify_flow", "regenerate_flow", "ingest", "qa", "unclear"} # ── flow-json 容错解析(Python 版,对齐前端 lib/flow-graph.ts)────────────── _FENCE = "```flow-json" _VALID_TYPES = {"start", "end", "stage", "task", "decision", "risk", "deliverable", "role"} def _as_str(v: Any) -> str: return v.strip() if isinstance(v, str) else "" def _loose_nodes(body: str) -> list[dict[str, Any]]: """正则逐字段抽取节点(strict json.loads 失败时回退)。容忍裸换行 / 尾逗号。""" # 取 "nodes" 后的 [ ... ](到与之配对的 ] 或文本结尾)。 key = body.find('"nodes"') if key < 0: return [] arr = body.find("[", key) if arr < 0: return [] # 平衡括号切出每个顶层 {…} slices: list[str] = [] depth = 0 start = -1 in_str = False esc = False for i in range(arr + 1, len(body)): ch = body[i] if in_str: if esc: esc = False elif ch == "\\": esc = True elif ch == '"': in_str = False continue if ch == '"': in_str = True elif ch == "{": if depth == 0: start = i depth += 1 elif ch == "}": depth -= 1 if depth == 0 and start >= 0: slices.append(body[start : i + 1]) start = -1 elif ch == "]" and depth == 0: break out: list[dict[str, Any]] = [] for sl in slices: out.append(_loose_one(sl)) return out def _loose_edges(body: str) -> list[dict[str, Any]]: """非法 JSON 回退路径下,从顶层 `"edges":[ {from,to,label} … ]` 容错抽边。""" key = body.find('"edges"') if key < 0: return [] arr = body.find("[", key) if arr < 0: return [] depth = 0 start = -1 in_str = False esc = False slices: list[str] = [] for i in range(arr + 1, len(body)): ch = body[i] if in_str: if esc: esc = False elif ch == "\\": esc = True elif ch == '"': in_str = False continue if ch == '"': in_str = True elif ch == "{": if depth == 0: start = i depth += 1 elif ch == "}": depth -= 1 if depth == 0 and start >= 0: slices.append(body[start : i + 1]) start = -1 elif ch == "]" and depth == 0: break def pick(slice_: str, *keys: str) -> str: for k in keys: m = re.search(r'"' + k + r'"\s*:\s*"([^"]*)"', slice_) if m: return m.group(1) return "" out: list[dict[str, Any]] = [] for sl in slices: frm = pick(sl, "from", "source") to = pick(sl, "to", "target") if frm and to: e: dict[str, Any] = {"from": frm, "to": to} lbl = pick(sl, "label") if lbl: e["label"] = lbl out.append(e) return out def _loose_one(slice_: str) -> dict[str, Any]: def pick(*keys: str) -> str | None: for k in keys: m = re.search(r'"' + k + r'"\s*:\s*"((?:\\.|[^"\\])*)"', slice_) if m: raw = m.group(1) try: return json.loads('"' + raw.replace("\n", "\\n").replace("\r", "") + '"') except Exception: # noqa: BLE001 return raw.replace("\\n", "\n").replace('\\"', '"') return None obj: dict[str, Any] = {} for k, dst in (("id", "id"), ("type", "type"), ("detail", "detail")): v = pick(k) if v: obj[dst] = v label = pick("label", "name") if label: obj["label"] = label group = pick("group", "level", "tier", "layer") if group: obj["group"] = group fm = re.search(r'"(?:from|prev)"\s*:\s*\[([\s\S]*?)\]', slice_) if fm: inner = fm.group(1) ids: list[Any] = [{"id": x} for x in re.findall(r'\{[^{}]*?"id"\s*:\s*"([^"]+)"[^{}]*?\}', inner)] if not ids: ids = re.findall(r'"([^"]+)"', inner) if ids: obj["from"] = ids return obj def _norm_node(raw: dict[str, Any]) -> dict[str, Any] | None: nid = _as_str(raw.get("id")) label = _as_str(raw.get("label")) or _as_str(raw.get("name")) if not nid or not label: return None rtype = _as_str(raw.get("type")).lower() ntype = rtype if rtype in _VALID_TYPES else "task" node: dict[str, Any] = {"id": nid, "label": label, "type": ntype} detail = _as_str(raw.get("detail")) or _as_str(raw.get("description")) if detail: node["detail"] = detail group = _as_str(raw.get("group")) or _as_str(raw.get("level")) or _as_str(raw.get("tier")) or _as_str(raw.get("layer")) if group: node["group"] = group return node def _edges_from(raw: dict[str, Any], nid: str) -> list[dict[str, Any]]: frm = raw.get("from") or raw.get("prev") if not isinstance(frm, list): return [] out: list[dict[str, Any]] = [] for e in frm: if isinstance(e, str) and e.strip(): out.append({"from": e.strip(), "to": nid}) elif isinstance(e, dict): fid = _as_str(e.get("id")) or _as_str(e.get("from")) if fid: edge = {"from": fid, "to": nid} lbl = _as_str(e.get("label")) if lbl: edge["label"] = lbl out.append(edge) return out def parse_flow_graph(text: str) -> dict[str, Any] | None: """从(可能带 fence / 非法 JSON 的)文本里解析出 ``{title, nodes, edges}``。无节点→None。""" start = text.rfind(_FENCE) body = text[start + len(_FENCE) :] if start >= 0 else text close = body.find("```") if close >= 0: body = body[:close] body = body.strip() title = "" nodes_raw: list[dict[str, Any]] = [] obj: Any = None for candidate in (body, re.sub(r",(\s*[}\]])", r"\1", body)): # 原文 + 去尾逗号 try: obj = json.loads(candidate) break except Exception: # noqa: BLE001 obj = None if isinstance(obj, dict): title = _as_str(obj.get("title")) raw_list = obj.get("nodes") nodes_raw = raw_list if isinstance(raw_list, list) else [] else: m = re.search(r'"title"\s*:\s*"([^"]*)"', body) if m: title = m.group(1) nodes_raw = _loose_nodes(body) nodes: list[dict[str, Any]] = [] edges: list[dict[str, Any]] = [] seen: set[str] = set() for raw in nodes_raw: if not isinstance(raw, dict): continue node = _norm_node(raw) if not node or node["id"] in seen: continue seen.add(node["id"]) nodes.append(node) edges.extend(_edges_from(raw, node["id"])) if not nodes: return None # 兼容**独立 `edges` 数组**:模型可能把连线写成节点上的 `from`,也可能写成顶层 # `edges:[{from,to,label}]`(或 source/target 别名)。后者若不处理,节点全无连线(bug)。 raw_edges: Any = obj.get("edges") if isinstance(obj, dict) else None if not isinstance(raw_edges, list): raw_edges = _loose_edges(body) for e in raw_edges or []: if not isinstance(e, dict): continue frm = _as_str(e.get("from")) or _as_str(e.get("source")) to = _as_str(e.get("to")) or _as_str(e.get("target")) if not frm or not to: continue edge: dict[str, Any] = {"from": frm, "to": to} lbl = _as_str(e.get("label")) if lbl: edge["label"] = lbl edges.append(edge) # 只保留两端都存在的边 + 去重。 valid_edges: list[dict[str, Any]] = [] keys: set[str] = set() for e in edges: if e["from"] not in seen or e["to"] not in seen or e["from"] == e["to"]: continue k = f'{e["from"]}->{e["to"]}' if k in keys: continue keys.add(k) valid_edges.append(e) # 连线选取(两全:尊重模型 from + 弱模型兜底): # · 模型把连线**连全了**(除根节点外每个节点都有入边)→ 直接用模型的 from/edges; # · 模型**漏连**(≥2 个节点没入边,最常见就是漏 from 导致平铺)→ 按 group 层级关系 # **确定性推导补全**,保证不平铺、连线质量与模型能力无关。 incoming = {e["to"] for e in valid_edges} orphans = [n for n in nodes if n["id"] not in incoming] derived = _derive_layered_edges(nodes) if valid_edges and len(orphans) <= 1: final_edges = valid_edges # 模型连全(仅根无入边)→ 用模型连边 elif derived is not None: final_edges = derived # 模型漏连 → 按层级推导补全(永不平铺) else: final_edges = valid_edges return {"title": title, "nodes": nodes, "edges": final_edges} def _derive_layered_edges(nodes: list[dict[str, Any]]) -> list[dict[str, Any]] | None: """按节点的 `group`(层级)确定性推导连线,**不依赖模型给的 from/edges**。 规则(对应「任务 → 各目的 → 各自对应的行为体 → … → 叙事」这种分层链): · 按 group 首次出现顺序排出层 L0,L1,…(L0 通常是「任务」单节点); · 相邻层 a→b 连边:|a|==1 → a 扇出到 b 全部;|b|==1 → a 全部汇入 b; |a|==|b| → 按下标一一对应(目的i→行为体i);否则该层对不规则 → 整体放弃推导(返回 None)。 返回 None 表示「无法可靠推导」(无 group / 只有一层 / 存在不规则层对),调用方回退模型连边。 """ # 按 group 分层,保持每层内节点的原始顺序 + 层的首次出现顺序。 # 加固:个别节点漏 group(弱模型可能漏标)→ 按发射顺序**就近归入前一个层级**(carry-forward), # 不再因一个漏标就整体放弃;只有连开头节点都没 group / 全程无任何 group 才放弃。 order: list[str] = [] layers: dict[str, list[str]] = {} last_group = "" any_group = False for n in nodes: g = _as_str(n.get("group")) if g: any_group = True last_group = g else: g = last_group # 漏 group → 沿用前一节点的层级 if not g: return None # 开头节点就没 group,无法锚定层级 → 回退模型连边 if g not in layers: layers[g] = [] order.append(g) layers[g].append(n["id"]) if not any_group or len(order) < 2: return None # 全程无 group / 只有一层 → 回退模型连边 out: list[dict[str, Any]] = [] for k in range(len(order) - 1): a = layers[order[k]] b = layers[order[k + 1]] if len(a) == 1: out.extend({"from": a[0], "to": nid} for nid in b) elif len(b) == 1: out.extend({"from": nid, "to": b[0]} for nid in a) elif len(a) == len(b): out.extend({"from": a[i], "to": b[i]} for i in range(len(a))) else: # 不规则层对(两边都 >1 且不等,多因模型漏/多了某层条目):**不放弃、不平铺**, # 按比例把下层每个节点连到上层对应位置,保证每个节点都有入边(连上 > 连得完美)。 for j, nid in enumerate(b): out.append({"from": a[j * len(a) // len(b)], "to": nid}) return out def _collapse_to_single_parent(graph: dict[str, Any], business_code: str | None) -> dict[str, Any]: """**纯树**业务链(rwfx)确定性兜底:每个非根节点最多保留一个父节点(一条入边)。 rwfx 全链是一棵树(除 1:1 的预期行为↔驱动因素,其余皆 1:N),任何节点都只有一个父节点、 绝无多对一/汇聚。但弱模型常给某节点(典型是「叙事」)连多个上层父节点 → 画出汇聚/乱线。 这里按「最近上一层」原则把每个节点的入边收敛到一条:在该节点的多个父节点里,优先保留 group 处于它**紧邻上一层**的那个(按业务链层级序 ``chain_level_rank``);同层多个则保留模型给出的 第一个;都判断不了就保留第一个。根节点(无入边)不变。**只删多余入边、不新增**,不影响节点。 """ nodes = graph.get("nodes") or [] edges = graph.get("edges") or [] rank = chain_level_rank(business_code) # 层级名 → 自上而下序号 group_by_id: dict[str, str] = { n["id"]: _as_str(n.get("group")) for n in nodes if isinstance(n, dict) and n.get("id") } appear: dict[str, int] = {} for n in nodes: g = _as_str(n.get("group")) if isinstance(n, dict) else "" if g and g not in appear: appear[g] = len(appear) def grank(nid: str) -> int: g = group_by_id.get(nid, "") if g in rank: return rank[g] return 1000 + appear.get(g, 0) # 未在业务层级里的 group 排到最后,不抢「上一层」判定 by_to: dict[str, list[dict[str, Any]]] = {} order: list[str] = [] for e in edges: t = e.get("to") if t not in by_to: by_to[t] = [] order.append(t) by_to[t].append(e) kept: list[dict[str, Any]] = [] collapsed = 0 for t in order: ins = by_to[t] if len(ins) <= 1: kept.extend(ins) continue tr = grank(t) prev_layer = [e for e in ins if grank(_as_str(e.get("from"))) == tr - 1] if prev_layer: best = prev_layer[0] else: upper = [e for e in ins if grank(_as_str(e.get("from"))) < tr] best = max(upper, key=lambda e: grank(_as_str(e.get("from")))) if upper else ins[0] kept.append(best) collapsed += len(ins) - 1 if collapsed: logger.info("flow-extract: collapsed %d extra parent edge(s) to single-parent (tree chain %s)", collapsed, business_code) graph["edges"] = kept return graph # ── prompts ──────────────────────────────────────────────────────────────── _SHAPE = ( "flow-json 严格用下面这个形状(一个 JSON 对象,nodes 是节点数组;示例仅示意结构、勿照抄内容):\n" "```flow-json\n" '{"title":"流程图标题","nodes":[\n' ' {"id":"n1","label":"任务简短名","type":"start","group":"任务","detail":"任务目标"},\n' ' {"id":"n2","label":"目的核心短语","type":"task","group":"目的","detail":"名称+理由要点(2-4句,单行)","from":["n1"]},\n' ' {"id":"n3","label":"行为体核心短语","type":"role","group":"行为体","detail":"名称+理由要点","from":["n2"]}\n' "]}\n" "```\n" "字段:id(n1/n2…唯一)、label(≤12字,节点标题)、type(start/end/stage/task/decision/risk/deliverable/role 之一,仅决定配色)、" "group(该节点所在层级名称,如 任务/目的/行为体/预期行为/驱动因素/叙事,必填)、" "detail(该条目的名称+**完整理由/说明,原样取自报告、不要简化或精简理由**;只需写成一行——内部换行用「;」或「/」替代,不要内含真实换行)、" "from(前驱节点 id 列表,只引用前面已出现的节点)。" ) _STRUCTURE = ( "【画法·严格遵守】" "① 根节点 n1=本次任务本身(报告/任务总览里的任务名称),label 用任务简短名,group 填「任务」,不要用「任务」泛称当 label;" "② 从第二层起按层级逐层展开:每层在报告/交付里通常列了多条具体条目(如 目的层有 目的1/2/3、行为体层有各自对应行为体,后面还有 预期行为·驱动因素、叙事…)," "把每一条具体条目各画成一个节点,不要把整层合并成一个笼统节点;" "③ 每个节点 group 填它所在层级名称(取自报告里该层级真实名称);" "④ **连线**:(a)group 标对层级;(b)**每个非根节点都必须填 from**,连到它在上一层级里**所属/对应**的那个节点;" "层间关系(哪层一对一、哪层一对多)**以本轮给出的业务链规范为准**:**一对一**的相邻层按位次配对(下层第i条连上层第i条);" "**一对多**的相邻层里,上层一条可对应下层多条——下层每条连到它**所属**的那一条上层节点,**绝不能把多条上层节点汇聚连到同一个下层节点**(那是把一对多画反成多对一);" "(c)最常见错误是漏 from 导致平铺,或把一对多画成多对一,务必每个都填对;系统还会按层级校正/补全连线兜底;" "⑤ 务必把每层每一条都输出成节点、不要漏层、不要只输出开头一两个节点就停(典型约13-16个节点);" "⑥ 只画这些层级实体节点,不要把『综合分析/核心洞察/执行路线/建议/结论』及『阶段①②③④』等汇总段落/阶段名画成节点。" ) def _material(req: FlowExtractRequest) -> str: """抽取的数据来源:**严格以《方案总结报告》及其文件内容为准**;各席位原始交付仅作兜底。 报告是各席位交付经总控综合提炼、**已做取舍**后的最终结论(如从 10 个候选目的里锁定 3 个)。 因此抽取**只按报告最终呈现的条目**逐条抽全——既不要漏(报告列了 3 条别只抽 1 条),也不要多 (别去席位产物里把报告已淘汰的候选也抽进来)。各席位原始交付**仅在报告某层整层缺失 / 报告明确 说有 N 条却没逐条列出时**才用来补全那一层;与报告冲突时一律以报告为准。 """ # xdfx:把本次「行动名称」作为材料抬头喂进去(根节点就是它),让抽取贴着这条行动展开。 action = (req.root_label or "").strip() prefix = f"【本次行动】{action}\n\n" if action else "" report = (req.report_md or "").strip() seats = [s for s in req.seat_deliveries if (s.content or "").strip()] parts: list[str] = [] if report: parts.append( "===== 方案总结报告全文(**主数据来源,以它为准**:按它最终呈现的条目逐条抽取、勿改写、勿增删)=====\n" + report + "\n===== 报告结束 =====" ) if seats: body = "\n\n".join(f"## 席位{i + 1}:{s.sender or '(未具名席位)'}\n{s.content.strip()}" for i, s in enumerate(seats)) if report: parts.append( "===== 各席位原始交付产物(**兜底参考,仅报告不足时才用**)=====\n" "用途:**只有当报告某个层级整层缺失、或报告明确说有 N 条却没逐条列出时**,才回到这里把那一层补全。\n" "⚠️ **绝不要用这里去扩充报告已经取舍过的层级**——报告里被淘汰 / 未采用的候选项(典型:席位列了 10 个候选目的," "但报告最终只锁定 3 个),**一律以报告最终采用的为准,不要把被淘汰的候选也抽进来**。与报告冲突时以报告为准。\n" + body + "\n===== 交付结束 =====" ) else: # 没有报告时,席位产物才是主来源(此时按席位最终采用的结论抽,候选池同样别全抽)。 parts.append("===== 各席位交付产物(数据来源,逐条抽取、勿改写;按其最终采用的结论抽,不要把被淘汰的候选都抽进来)=====\n" + body + "\n===== 交付结束 =====") if not parts: return "" return prefix + "\n\n".join(parts) def _hierarchy_hint(req: FlowExtractRequest) -> str: # 深链接业务规范**最高优先**:按 business_code 取写实/通用规范,覆盖任何默认层级假设。 spec = build_structure_directive(req.business_code) if spec: cp = (req.coordinator_prompt or "").strip() if cp: spec += ( "\n(若「业务链总控编排提示」对层级另有更细要求,可在**不违背上述完整层级**的前提下采纳:\n" + cp + ")" ) return spec if (req.coordinator_prompt or "").strip(): return "层级顺序优先采用下面「业务链总控编排提示」里写明的层级链(如它写的 任务→目的→行为体→…):\n" + req.coordinator_prompt.strip() levels = [s.strip() for s in req.chain_levels if s.strip()] if levels: return "层级顺序(自上而下):任务 → " + " → ".join(levels) + "。" return "层级顺序:以「任务」为根,按报告/交付的逻辑结构合理分层展开。" def _root_label_hint(req: FlowExtractRequest) -> str: """根节点 label 强制提示:xdfx 用「行动名称」当根节点(而非任务名)。无 root_label → 空。""" label = (req.root_label or "").strip() if not label: return "" return f"\n【根节点强制】整张图的根节点 n1 的 label **必须**用本次行动名称:「{label}」——不要用任务名/其它当根节点。" # 兜底系统提示:当 roundtable-structure 的 SOUL 缺失时用它(与 SOUL 同款契约)。 _DEFAULT_SYSTEM = ( "你是「结构化抽取智能体」。职责:读懂《方案总结报告》/ 各席位交付内容,分析其层级结构," "抽取成可交互流程图的结构化数据(flow-json);也按指令改图、校验修复。" "请**先用简短文字(2-4句)说明分析过程**(识别的层级链、每层条目数、分支对应)让用户看到进展," "**再**在回复末尾输出一个完整的 ```flow-json 代码块(除这段简短说明外不写长篇正文)。全程简体中文。\n" + _SHAPE + "\n" + _STRUCTURE ) def _system_prompt() -> str: """结构化抽取的系统提示 = roundtable-structure 智能体的 SOUL(可在智能体管理页编辑、 未来挂入库技能);SOUL 缺失则回退内置 _DEFAULT_SYSTEM。""" soul = load_agent_soul(STRUCTURE_AGENT_ID) return soul.strip() if soul and soul.strip() else _DEFAULT_SYSTEM def _resolve_model(payload_model: str | None) -> str | None: """模型优先级:前端所选(payload) > 智能体配置的 model > None(用 config 默认模型)。""" if payload_model: return payload_model try: cfg = load_agent_config(STRUCTURE_AGENT_ID) if cfg and cfg.model: return cfg.model except Exception: # noqa: BLE001 pass return None def _extract_messages(req: FlowExtractRequest) -> list[Any]: user = ( _material(req) + "\n\n" + _hierarchy_hint(req) + _root_label_hint(req) + "\n\n【按报告逐条抽全,不漏也不多(最重要)】**严格以《方案总结报告》最终呈现的条目为准**:" + "报告某层级用表格 / 列表 / ①②③ 逐条列了几条,就出几个节点、**一条都不能漏**(不要只抽被标「确定 / 第1名」的那条而漏掉同列的其余);" + "某上层节点名下列了几条下层条目就出几个节点,**不许把多条合并 / 概括成一条**;" + "⚠️ 但**不要无中生有**:报告已明确取舍、最终只采用 N 条的层级(如「从 10 个候选目的中锁定 3 个」「最终选定方案」)," + "**只抽报告最终采用的那 N 条**,**绝不要去席位产物或报告的「候选池 / 被淘汰项」里把没被采用的也抽出来凑数**;" + "报告把某层多条目压缩成一句话、没逐条列出时(如只说「12 条预期行为」却没列),才回《各席位原始交付产物》把那一层补全。" + "不要只照搬报告结尾「方案总纲 / 最终框架」里被精简的主干,要回到正文(如「各席位观点与交付」)把报告列出的每层条目抽全。" + "\n\n请分析以上数据来源,按你的职责设定输出一个完整的 ```flow-json 代码块。" ) return [SystemMessage(content=_system_prompt()), HumanMessage(content=user)] def _edit_messages(req: FlowExtractRequest) -> list[Any]: """Phase 2 问答修改:按用户指令改当前流程图(增/删/改节点或连线),其余保持不变。""" cur = json.dumps(req.current_graph or {}, ensure_ascii=False) user = ( _material(req) + "\n\n===== 当前流程图(flow-json,待按指令修改)=====\n" + cur + "\n===== 当前流程图结束 =====\n\n用户修改指令:\n" + (req.instruction or "").strip() + "\n\n请按指令修改这张流程图:未被指令涉及的部分保持原样、节点 id 尽量沿用、删节点连带删边," + "输出**修改后的完整** ```flow-json 代码块。" ) return [SystemMessage(content=_system_prompt()), HumanMessage(content=user)] def _validate_messages(req: FlowExtractRequest, extracted: str) -> list[Any]: # 深链接业务规范:校验阶段也带上完整层级 + 根节点强制,确保按规范补齐缺层、根节点用行动名。 spec = build_structure_directive(req.business_code) spec_block = (spec + _root_label_hint(req) + "\n\n") if (spec or req.root_label) else "" user = ( spec_block + _material(req) + "\n\n===== 已抽取的 flow-json(待校验修复)=====\n" + extracted + "\n===== 待校验结束 =====\n\n请对照数据来源,**逐个节点检查上面这份 flow-json 是不是缺字段**、并修复它:" + "\n**校验要点(务必遵守)**:" + "① **连接关系 `from`(最重要,必查必补)**:**逐个**检查上面每个节点——除根任务外," + "**只要某个节点缺了 `from`(或 `from` 为空),就必须给它补上正确的 `from`**:连到它在上一层级里**所属/对应**的那个节点。" + "层间关系(哪层一对一、哪层一对多)**以上面给出的业务链规范为准**:一对一层级按位次配对(下层第i条连上层第i条);" + "**一对多**层级里下层每条连到它**所属**的那一条上层节点,**绝不能把多条上层节点汇聚到同一个下层节点**(那是把一对多画反成多对一,是本次要重点纠正的错)。" + "缺 `from` / 把一对多画成多对一 是最常见错漏,**绝不能让任何非根节点悬空、也不要张冠李戴**;" + "② **一对一(1:1)层级必须严格 1:1、绝不多对一(按业务链规范,如预期行为↔驱动因素;本次重点纠错)**:" + "这两层节点数**必须完全相等**、按位次一一配对;**每个下层节点的 `from` 只能写一条上层节点**," + "**绝不允许多条上层节点共用 / 汇聚到同一个下层节点**(多对一),也不允许同一条上层节点被多个下层节点指向" + "(即 1:1 两层里任何节点都不能冒出 / 汇入多条连线)。若发现下层节点偏少(被合并)导致只能多对一," + "**必须按上层逐条拆分 / 补出等量的下层节点**(每条上层各配一个,内容综合提炼),使两层等量后再逐一配对;" + "③ **逐层条目数对齐《方案总结报告》(不漏也不多,必查)**:以报告为准——报告某层列了几条就该有几个节点。" + "漏抽的补齐(报告列了 3 条却只抽了 1 条 → 补到 3;某行为体名下多条预期行为被并成一句 → 拆全);" + "**但绝不要新增报告未采用的候选**——报告已取舍、最终只采用 N 条的层级(如从 10 个候选目的里锁定 3 个)," + "**只保留报告最终采用的那 N 条;若图里混进了报告未采用 / 被淘汰的候选(从席位候选池里抽出来的),把它们删掉**。" + "仅当报告把某层压缩成一句话、没逐条列时,才回《各席位原始交付产物》补全那一层;" + "④ **不要漏报告列出的层 / 条目**;一对多的下层(某上层对应多条下层)不要被错误合并成一条;" + "⑤ 每个节点的 `group`(层级名)标对;一对一层级内条目按相同位次对齐,一对多层级下层每条紧挨其所属上层条目分组;" + "⑥ 每个节点的 `detail` **原样保留、不要简化或精简理由**——理由要完整,不要删减内容;" + "⑦ 没有汇总/阶段编号节点;只在确有错漏时才动,没问题就原样规范化输出。" + "\n输出修复后的完整 ```flow-json 代码块(含全部节点 + 每个非根节点的 from)。" ) return [SystemMessage(content=_system_prompt()), HumanMessage(content=user)] # ── endpoint ───────────────────────────────────────────────────────────── def _frame(data: dict[str, Any]) -> str: return f"data: {json.dumps(data, ensure_ascii=False, default=str)}\n\n" def _chunk_parts(chunk: Any) -> tuple[str, str]: """从流式 chunk 取 (正文增量, 思考增量)。思考走 additional_kwargs.reasoning_content (DeepSeek/兼容 OpenAI 推理模型的单列字段),正文走 content。""" content = chunk.content if isinstance(getattr(chunk, "content", None), str) else "" ak = getattr(chunk, "additional_kwargs", None) or {} reasoning = "" if isinstance(ak, dict): r = ak.get("reasoning_content") or ak.get("reasoning") if isinstance(r, str): reasoning = r return content, reasoning @router.post("/stream") async def extract_flow_stream(request: Request, payload: FlowExtractRequest = Body(...)) -> StreamingResponse: """流式抽取流程图 JSON + 流式校验修复(抽取/校验过程含思考都实时吐前端)→ 返回最终结构化图。""" ensure_roundtable_functional_agents() # 自愈:确保 roundtable-structure 的 SOUL/config 存在 resolved_model = _resolve_model(payload.model) async def gen() -> AsyncIterator[str]: try: # 带 instruction + current_graph → Phase 2 问答修改(按指令改当前图);否则首次抽取。 is_edit = bool((payload.instruction or "").strip()) and payload.current_graph is not None if not _material(payload).strip() and not is_edit: yield _frame({"event": "error", "detail": "缺少数据来源(报告或各席位交付均为空)"}) return # ① 抽取 / 编辑(正文 + 思考(reasoning) 拼成 `…正文` 增量逐帧吐,前端实时显示) # 模型自动容错:结构化抽取非常重要——某模型 astream 抛错 / 产出空时,按 config.yaml # models[] 顺序切到下一个模型重抽,并发 model_switch 帧让前端把抽取面板已累计的半截内容重置。 round1_messages = _edit_messages(payload) if is_edit else _extract_messages(payload) chain = fallback_model_chain(resolved_model) or [resolved_model] extracted = "" used_model = chain[0] prev_model: str | None = None prev_reason: str | None = None for idx, candidate in enumerate(chain): if idx > 0: yield _frame( { "event": "model_switch", "agent_name": "结构化抽取", "failed_model": prev_model or "", "next_model": candidate or "", "reason": reason_text(prev_reason), } ) used_model = candidate extracted = "" think_open = False try: model = create_chat_model(name=candidate, thinking_enabled=True) async for chunk in model.astream(round1_messages): c, r = _chunk_parts(chunk) out = "" if r: if not think_open: out += "" think_open = True out += r if c: if think_open: out += "" think_open = False out += c extracted += c if out: yield _frame({"event": "extract_delta", "content": out}) if think_open: yield _frame({"event": "extract_delta", "content": ""}) except Exception as exc: # noqa: BLE001 — 该模型抽取失败:换下一个模型重抽 prev_model, prev_reason = candidate, "exception" logger.warning("flow-extract round1 model %s failed: %s", candidate, exc) if idx + 1 < len(chain): continue raise if extracted.strip(): break # 产出为空:换下一个模型(已是最后一个则继续走后续解析兜底)。 prev_model, prev_reason = candidate, "empty" if idx + 1 >= len(chain): break yield _frame({"event": "extract_done"}) # ② 校验 / 修复(同样流式 + 思考实时显示) yield _frame({"event": "validating"}) validated = "" graph: dict[str, Any] | None = None try: vmodel = create_chat_model(name=used_model, thinking_enabled=True) think_open = False async for chunk in vmodel.astream(_validate_messages(payload, extracted)): c, r = _chunk_parts(chunk) out = "" if r: if not think_open: out += "" think_open = True out += r if c: if think_open: out += "" think_open = False out += c validated += c if out: yield _frame({"event": "validate_delta", "content": out}) if think_open: yield _frame({"event": "validate_delta", "content": ""}) graph = parse_flow_graph(validated) except Exception: # noqa: BLE001 — 校验失败不致命,回退抽取结果 logger.warning("flow-extract validation failed, fall back to extraction", exc_info=True) graph = None extracted_graph = parse_flow_graph(extracted) # 防回归兜底:校验**绝不能让结构变差**。若校验后的图节点或连线比抽取结果更少 # (模型校验时把连线删没了、或漏了节点——实测会发生),直接回退用抽取结果, # 保住所有节点与连线。校验只在「不丢东西」时才采用。 if graph is not None and extracted_graph is not None: if len(graph.get("nodes", [])) < len(extracted_graph.get("nodes", [])) or len( graph.get("edges", []) ) < len(extracted_graph.get("edges", [])): logger.info( "flow-extract: validation reduced nodes/edges (%d/%d → %d/%d), keep extraction", len(extracted_graph.get("nodes", [])), len(extracted_graph.get("edges", [])), len(graph.get("nodes", [])), len(graph.get("edges", [])), ) graph = extracted_graph if graph is None: graph = extracted_graph if graph is None: yield _frame({"event": "error", "detail": "未能从模型输出解析出流程图节点"}) return # 纯树业务链(rwfx)确定性兜底:把任何节点的多父入边收敛到单一父节点,杜绝弱模型 # 违反 1:N 画出的汇聚/乱线(典型是叙事连了一片驱动因素)。非纯树业务(xdfx 有汇聚)跳过。 if is_single_parent_chain(payload.business_code): graph = _collapse_to_single_parent(graph, payload.business_code) yield _frame({"event": "final", "graph": graph}) except Exception as exc: # noqa: BLE001 logger.exception("flow-extract stream failed") yield _frame({"event": "error", "detail": str(exc)}) return StreamingResponse( gen(), media_type="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "Connection": "keep-alive"}, ) # ── 意图识别(一次性 JSON,参考 ai_writing/intent_parser 范式)───────────────── def _parse_intent_json(text: str) -> dict[str, Any] | None: """从模型输出里容错解析 {intent, reason}。strict → 去围栏/去尾逗号 → 平衡括号。""" candidates: list[str] = [] t = text.strip() candidates.append(t) m = re.search(r"```(?:json)?\s*([\s\S]+?)\s*```", t) if m: candidates.append(m.group(1).strip()) s = t.find("{") e = t.rfind("}") if s >= 0 and e > s: candidates.append(t[s : e + 1]) for c in candidates: for cand in (c, re.sub(r",(\s*[}\]])", r"\1", c)): try: obj = json.loads(cand) if isinstance(obj, dict): return obj except Exception: # noqa: BLE001 continue return None _INTENT_SYS = ( "你是圆桌「方案总结」第三步对话的**意图识别器**,也是个能帮用户的助手。第三步已产出:方案总结报告(Markdown)、" "结构化流程图,可能还可把成果入库到知识库。用户在对话框发了一条消息,请**按语义理解他真正想干什么**" "(不要只做关键词匹配),只输出一个 JSON 对象(不要解释、不要 markdown 代码块、用英文双引号):\n" '{"intent":"<枚举>","reason":"<一句话依据>"}\n' "intent 取值:\n" "- regenerate_report:想**重做/重写整份**方案总结报告(如「重新生成报告」「报告再写一版」「重新分析一下」);\n" "- modify_flow:想**修改当前流程图**的节点/连线(如「把目的1改成X」「删掉第3个行为体」「加个验收节点」「这一层调整下」);\n" "- regenerate_flow:想**重新生成整张**流程图(如「重新画流程图」「流程图重来」「重新结构化输出」);\n" "- ingest:想把本次成果**入库/沉淀/存进知识库**(如「帮我入库」「入库吧」「把方案存进知识库」「沉淀一下」);\n" "- qa:**提问 / 要求解释、介绍、概述、答疑**,只需用文字回答、不改任何产物(如「介绍一下这报告的内容」" "「这方案主要讲了啥」「为什么选 OTA 平台」「风险最高的是哪条」「帮我看看有没有遗漏」「总结成几句话」);\n" "- unclear:**实在**判断不出他想做哪件事(如只说「改一下」「不对」「重来」「再说」且无任何指向)。\n" "判断要点:\n" "1. **优先识别出明确意图,不要动不动就 unclear**——只要能合理判断就给具体意图;尤其『介绍/讲讲/解释/说明/" "为什么/是什么/有没有/帮我看/概述/总结一下』这类是**对内容的提问** → 一律 qa(用文字答),不要误判成重做报告。\n" "2. 含『节点/连线/这一层/这个目的/行为体/流程图…的 改/删/加/换』 → modify_flow(前提是已有流程图)。\n" "3. 明确说『报告』要重写/重做才是 regenerate_report;只是问报告内容是 qa。\n" "4. **入库类意图只有在下文上下文标注「可入库=是」时才允许返回 ingest**;若标注不可入库,用户提入库就归 qa" "(用文字说明当前不可入库/无适配技能)。\n" "示例:\n" '「帮我介绍一下这报告中的内容」→ {"intent":"qa","reason":"要求介绍报告内容,属提问答疑"}\n' '「这方案为什么这么定」→ {"intent":"qa","reason":"询问方案理由"}\n' '「帮我入库」(可入库=是) → {"intent":"ingest","reason":"明确要求入库且当前可入库"}\n' '「再帮我生成一次报告」→ {"intent":"regenerate_report","reason":"明确要求重做报告"}\n' '「把目的1的名称改短点」→ {"intent":"modify_flow","reason":"修改流程图某节点"}\n' '「流程图重新画一张」→ {"intent":"regenerate_flow","reason":"重画整张流程图"}\n' '「改一下」→ {"intent":"unclear","reason":"没说改什么"}' ) @router.post("/intent") async def classify_flow_intent(request: Request, payload: FlowIntentRequest = Body(...)) -> dict[str, str]: """识别第三步输入框消息的意图(一次模型调用)。失败 / 解析不出 → unclear,让前端发起澄清。""" msg = (payload.message or "").strip() if not msg: return {"intent": "unclear", "reason": "空消息"} ctx = ( f"当前上下文:{'已生成报告' if payload.has_report else '尚无报告'};" f"{'已生成流程图' if payload.has_flow else '尚无流程图'};" f"可入库={'是' if payload.can_ingest else '否'}。" ) try: model = create_chat_model(name=_resolve_model(payload.model), thinking_enabled=False) res = await model.ainvoke( [SystemMessage(content=_INTENT_SYS), HumanMessage(content=ctx + "\n\n用户消息:" + msg + "\n\n请输出 JSON:")] ) text = res.content if isinstance(res.content, str) else _as_str(getattr(res, "text", "")) obj = _parse_intent_json(text) or {} intent = _as_str(obj.get("intent")).lower() if intent not in _INTENTS: intent = "unclear" # 不可入库时模型若误判 ingest → 降级为 qa(用文字说明当前不可入库),不要触发入库。 if intent == "ingest" and not payload.can_ingest: intent = "qa" return {"intent": intent, "reason": _as_str(obj.get("reason"))} except Exception: # noqa: BLE001 — 分类失败不致命,回退澄清让用户自己定 logger.warning("flow intent classify failed", exc_info=True) return {"intent": "unclear", "reason": "意图识别失败"} # ── 入库意图识别(「稍后入库」后用户在输入框的自由文字)────────────────────────── class IngestIntentRequest(BaseModel): """第三步「稍后入库」待命时,识别输入框消息是否要开始入库 / 补充内容后入库 / 其他。""" message: str = "" model: str | None = Field(default=None) _INGEST_INTENTS = {"start_ingest", "add_data", "other", "unclear"} _INGEST_INTENT_SYS = ( "你是圆桌「方案总结」第三步**入库意图识别器**。流程图已生成、并且用户此前选择了「稍后入库」," "现在他在输入框发了一条消息。请判断他的**真实意图**,只输出一个 JSON 对象(不要解释、不要 markdown 代码块、用英文双引号):\n" '{"intent":"<枚举>","reason":"<一句话依据>"}\n' "intent 只能取其一:\n" "- start_ingest:现在就把当前方案入库(如「入库吧」「开始入库」「现在入库」「把方案存进知识库」);\n" "- add_data:想先补充/追加一些内容再入库(如「补一段背景再入库」「加上风险说明后入库」「把这段也一起入库」);\n" "- other:与入库无关,是改报告/改流程图/提问等其它诉求(如「报告再写细点」「流程图加个验收节点」「为什么选这个方案」);\n" "- unclear:和入库相关但说不清是现在入库还是要补内容(如只说「入库」却又像没说完、「那个再说」之类模糊表达)。\n" "示例:\n" '「现在入库」→ {"intent":"start_ingest","reason":"明确要求立即入库"}\n' '「我再补一段实施风险,然后入库」→ {"intent":"add_data","reason":"先补内容再入库"}\n' '「报告执行部分写细一点」→ {"intent":"other","reason":"是改报告,与入库无关"}\n' '「入库……嗯」→ {"intent":"unclear","reason":"提到入库但表达不完整"}' ) @router.post("/ingest-intent") async def classify_ingest_intent(request: Request, payload: IngestIntentRequest = Body(...)) -> dict[str, str]: """识别「稍后入库」待命时输入框消息的入库意图(一次模型调用)。失败/解析不出 → unclear,前端发起协助。""" msg = (payload.message or "").strip() if not msg: return {"intent": "unclear", "reason": "空消息"} try: model = create_chat_model(name=_resolve_model(payload.model), thinking_enabled=False) res = await model.ainvoke( [SystemMessage(content=_INGEST_INTENT_SYS), HumanMessage(content="用户消息:" + msg + "\n\n请输出 JSON:")] ) text = res.content if isinstance(res.content, str) else _as_str(getattr(res, "text", "")) obj = _parse_intent_json(text) or {} intent = _as_str(obj.get("intent")).lower() if intent not in _INGEST_INTENTS: intent = "unclear" return {"intent": intent, "reason": _as_str(obj.get("reason"))} except Exception: # noqa: BLE001 — 分类失败不致命,回退协助让用户自己定 logger.warning("ingest intent classify failed", exc_info=True) return {"intent": "unclear", "reason": "意图识别失败"} # ── 入库适用性判定(结构化抽取智能体「自己看技能是否适合本场景」)────────────────── # # 需求:入库不再由前端硬编码「有没有含『入库』关键字的技能」判定,而是让**结构化输出 # 智能体**读它**挂载的入库技能**(roundtable-structure 的 config.yaml ``skills``)+ 当前 # 业务场景(如 rwfx-A=六步法),**自行判断**是否有适合本场景的入库技能。判定为可入库才 # 由前端显示「入库」按钮。真正调用技能在 ``POST /api/multi-agent/run/stream``(role=ingest)。 class IngestEligibilityRequest(BaseModel): """入库适用性判定入参:业务场景 + taskId,供智能体对照挂载的入库技能判断。""" task_id: str | None = Field(default=None) business_code: str | None = Field(default=None, description="业务码,如 6BF/3Q/7BF 或 rwfx-A") chain_name: str | None = Field(default=None, description="业务链条名称") coordinator_prompt: str | None = Field(default=None, description="业务链总控编排提示(含场景线索,如六步法)") report_title: str | None = Field(default=None) model: str | None = Field(default=None) def _structure_skill_candidates() -> list[dict[str, str]]: """读取挂在 roundtable-structure 上、且**已启用**的技能(候选入库技能)。 返回 ``[{name, description}]``。失败 / 未配置 → 空列表(= 不可入库)。 技能元数据取自磁盘 SKILL.md(name + description),足够智能体判断适用性。 """ try: cfg = load_agent_config(STRUCTURE_AGENT_ID) except Exception: # noqa: BLE001 return [] attached = [s for s in (getattr(cfg, "skills", None) or []) if isinstance(s, str) and s.strip()] if not attached: return [] attached_set = {s.strip().lower() for s in attached} try: storage = get_or_new_skill_storage(app_config=get_app_config()) skills = storage.load_skills(enabled_only=True) except Exception: # noqa: BLE001 logger.warning("ingest eligibility: load_skills failed", exc_info=True) return [] out: list[dict[str, str]] = [] for sk in skills: if (sk.name or "").strip().lower() in attached_set: out.append({"name": sk.name, "description": sk.description or ""}) return out _INGEST_ELIGIBILITY_SYS = ( "你是圆桌「方案总结」第三步的**结构化抽取智能体**,正在判断「本次会商成果**能否入库**」。" "下面给你两样东西:① 当前业务场景(业务码 / 业务链名称 / 总控编排提示等线索);" "② 你**挂载的候选入库技能**清单(每个有 name 与 description)。\n" "请判断:这些候选技能里,**是否有适合把当前这种业务场景的会商成果入库的技能**。判断依据:\n" "- 技能名称 / 描述是否表明它就是做「入库 / 沉淀 / 写入知识库」的;\n" "- 它是否对应当前业务场景——例如六步法类业务(如 rwfx-A)的入库技能,其名称通常**含「六步法」**字样;\n" "- 若某技能明显是为本场景设计的入库技能 → 可入库;若候选里压根没有入库类技能、或都与本场景不符 → 不可入库。\n" "只输出一个 JSON 对象(不要解释、不要 markdown 代码块、用英文双引号):\n" '{"can_ingest":,"skill":"<选中的技能 name,不可入库时为空字符串>","reason":"<一句话依据>"}' ) @router.post("/ingest-eligibility") async def classify_ingest_eligibility(request: Request, payload: IngestEligibilityRequest = Body(...)) -> dict[str, Any]: """让结构化抽取智能体判断当前业务场景**是否有合适的入库技能**(决定前端是否显示入库按钮)。 候选 = 挂在 roundtable-structure 上且已启用的技能。无候选 → 直接不可入库(不调模型)。 有候选 → 一次模型调用判断适用性,返回 ``{can_ingest, skill, reason}``。失败 → 不可入库(安全回退)。 """ ensure_roundtable_functional_agents() # 自愈:确保 roundtable-structure 存在 candidates = _structure_skill_candidates() if not candidates: return {"can_ingest": False, "skill": "", "reason": "结构化抽取智能体未挂载入库技能"} scenario_lines = [] if (payload.business_code or "").strip(): scenario_lines.append(f"业务码:{payload.business_code.strip()}") if (payload.chain_name or "").strip(): scenario_lines.append(f"业务链名称:{payload.chain_name.strip()}") if (payload.report_title or "").strip(): scenario_lines.append(f"方案标题:{payload.report_title.strip()}") if (payload.coordinator_prompt or "").strip(): scenario_lines.append("总控编排提示(节选):" + payload.coordinator_prompt.strip()[:1200]) scenario = "\n".join(scenario_lines) or "(未提供额外业务场景线索)" skills_block = "\n".join(f"- name: {c['name']}\n description: {c['description']}" for c in candidates) valid_names = {c["name"] for c in candidates} user = ( "【当前业务场景】\n" + scenario + "\n\n【候选入库技能】\n" + skills_block + "\n\n请输出 JSON:" ) try: model = create_chat_model(name=_resolve_model(payload.model), thinking_enabled=False) res = await model.ainvoke([SystemMessage(content=_INGEST_ELIGIBILITY_SYS), HumanMessage(content=user)]) text = res.content if isinstance(res.content, str) else _as_str(getattr(res, "text", "")) obj = _parse_intent_json(text) or {} skill = _as_str(obj.get("skill")) can = bool(obj.get("can_ingest")) and skill in valid_names # 模型说可入库但没给/给错技能名 → 单候选时回退到唯一候选,否则判不可入库(避免乱选)。 if bool(obj.get("can_ingest")) and skill not in valid_names: if len(candidates) == 1: skill = candidates[0]["name"] can = True else: can = False return { "can_ingest": can, "skill": skill if can else "", "reason": _as_str(obj.get("reason")), } except Exception: # noqa: BLE001 — 判定失败不致命,安全回退为不可入库 logger.warning("ingest eligibility classify failed", exc_info=True) return {"can_ingest": False, "skill": "", "reason": "入库适用性判定失败"}