deerflow-code/offline-backend-20260512/backend/app/gateway/routers/flow_extract.py
2026-09-07 18:24:55 +08:00

1075 lines
57 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""流程图结构化抽取 + 校验 的 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) 拼成 `<think>…</think>正文` 增量逐帧吐,前端实时显示)
# 模型自动容错:结构化抽取非常重要——某模型 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>"
think_open = True
out += r
if c:
if think_open:
out += "</think>"
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": "</think>"})
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>"
think_open = True
out += r
if c:
if think_open:
out += "</think>"
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": "</think>"})
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":<true|false>,"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": "入库适用性判定失败"}