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

1010 lines
51 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.

"""圆桌真网关(Phase 3b)—— 进程内驱动真实 leader / seat / report agent run。
实现 ``RoundtableGateway`` 协议(定义在 harness 的 roundtable_orchestrator.types),
但**进程内**跑真模型:复用 ``services.start_run``(scheduled task 同款)驱动一轮 agent,
跑完从 checkpointer 读最终 messages,用 ``multi_agent`` 的纯解析函数解析派活 / 澄清。
不走 HTTP、不需要 Request(用 ``SimpleNamespace(app=app, state=...)`` 当虚拟 request,
``set_current_user`` 由执行器设好)。
本网关位于 app 层(依赖 app.state + multi_agent helpers),实现 harness 定义的协议——
不违反 harness→app 边界(harness 永不 import app;app import harness 合法)。
⚠️ 需在能起完整后端栈 + 真模型的环境验证(见 scripts/probe_roundtable_job.py)。
"""
from __future__ import annotations
import asyncio
import logging
import os
import uuid
from collections.abc import Callable
from datetime import UTC, datetime
from types import SimpleNamespace
from typing import Any
from langgraph.checkpoint.base import empty_checkpoint
from app.gateway.roundtable_model_fallback import (
fallback_model_chain,
is_llm_error_message,
reason_text,
)
from app.gateway.roundtable_run_policy import roundtable_run_policy
from app.gateway.roundtable_seat_skills import append_seat_role_boundary, append_seat_skill_directive
from app.gateway.routers._roundtable_seed import (
COORDINATOR_AGENT_ID,
DASHBOARD_AGENT_ID,
REPORT_AGENT_ID,
SUMMARY_AGENT_ID,
ensure_roundtable_functional_agents,
)
from app.gateway.services import RUN_SCOPE_BACKGROUND_CHILD
from deerflow.agents.roundtable_orchestrator.delivery import classify_seat_delivery, visible_seat_text
from deerflow.agents.roundtable_orchestrator.report_json import extract_report_json
from deerflow.agents.roundtable_orchestrator.types import (
LeaderTurn,
ModelSwitch,
ReportResult,
SeatRef,
SeatTurn,
SummaryResult,
)
from deerflow.config.paths import get_paths
from deerflow.persistence.roundtable_diagnostics import DiagnosticsRecorder, get_default_store
from deerflow.runtime.serialization import serialize_channel_values
logger = logging.getLogger(__name__)
# 单轮 agent run(leader/seat/report/summary/dashboard)的**无活动(idle)超时**(秒)。
#
# 背景:后台编排是**串行 await 每一个 run**,且历史上**没有任何超时保护**——只要某一轮
# ``agent.astream`` 卡死(内网模型端点连接被打满 / sqlite 锁在挂载卷上挂死等环境相关原因),
# 整个后台作业就永久冻结(用户侧表现为「执行完第一个席位后一直不动」)。
#
# ⚠️ 关键:这里**不是总时长超时**,而是 **idle 超时**——只要这轮 run 还在持续产出
# (流式 token / 工具调用 / 进度事件,全部经 StreamBridge 外发),就一直等下去,哪怕这轮
# 真的要跑 20 分钟、半小时也绝不打断(慢但活着的生成不该被切)。只有当 run **连续
# ``_SEAT_IDLE_TIMEOUT_S`` 秒完全没有任何输出**(真卡死)时才取消它、抛错 → 顶层把 job 标
# error(而非永久卡死)+ 落一条 seat_turn_timeout 诊断。默认 1200s(20 分钟)——内网模型
# 很不稳定时单次「首 token / 工具间隔」可能拖很久,给足容忍度,宁可多等也别误杀慢生成;
# idle 只兜真正「连续 20 分钟啥都不发生」的死锁。
_SEAT_IDLE_TIMEOUT_S = int(os.getenv("ROUNDTABLE_SEAT_IDLE_TIMEOUT_SECONDS", "1200") or "1200")
# 可选的绝对总时长上限(秒):默认 0 = **关闭**(只靠 idle 超时,慢生成不受限)。
# 仅当你确实想给单轮设一个天花板时再设它(如某些场景要硬兜底)。
_SEAT_MAX_TOTAL_S = int(os.getenv("ROUNDTABLE_SEAT_MAX_TOTAL_SECONDS", "0") or "0")
# idle 检测的心跳间隔:每隔这么久检查一次「距上次活动过了多久」。
_IDLE_CHECK_INTERVAL_S = 15.0
_ROUNDTABLE_THREAD_METADATA: dict[str, Any] = {"thread_type": "roundtable", "system": True}
_REPORT_FILE = "方案总览.html"
_SUMMARY_FILE = "方案总结报告.md"
# 各角色的工具裁剪 / thinking / reasoning 全部取自单一数据源 roundtable_run_policy
# (与 multi_agent 的 _leader_run / _special_run 同源),不再在本文件硬编码——历史上
# 这里漏禁了 web_search,导致后台席位 / 报告能联网搜索而前端不能,现由 policy 统一。
# ── 后台作业「按用户限流」(2026-06-16)─────────────────────────────────────────
# 痛点:进 2+ 个后台挂起后,整个网关变慢——建立智能体连接、等待模型响应都慢。
# 根因:① checkpointer 用单连接 aiosqlite(AsyncSqliteSaver,内部一把 asyncio.Lock 串
# 行全进程的 checkpoint 读写);② 后台作业无并发上限——每个作业进程内驱动真实 leader/
# seat/report run,DAG/recommend 的 stage 还用 asyncio.gather 把席位真并行铺开。多个后台
# 作业同时把几十个 agent run 压到同一把 SQLite 锁 + 同一个离线模型端点上,前台实时会商
# 的 checkpoint 读写 / 模型调用只能排队 → 「干啥都慢」。
#
# 解法:给**每个用户**一个并发额度(默认 2),后台每个 agent 轮次(leader/seat/report/
# summary/dashboard)真正开跑前先取一个额度,跑完释放。同一用户的所有后台作业共享这 2 个
# 额度——随便挂多少个作业都能立刻建好,但同一瞬间真正在打模型/写 checkpoint 的后台轮次
# 至多 2 个,不会把共享资源打满。**前台实时路径(multi_agent._leader_run/_special_run、
# 普通聊天)不经过本文件的 `_run_agent_turn`,永不取这个额度,故前台不受限。**
_BACKGROUND_SLOTS_PER_USER = max(1, int(os.getenv("ROUNDTABLE_BACKGROUND_SLOTS_PER_USER", "2") or "2"))
# 按 user_id 缓存 Semaphore。网关全程跑在同一个事件循环上,普通 dict 取/建即可(无跨线程)。
_user_background_sems: dict[str, asyncio.Semaphore] = {}
def _user_background_semaphore(user_id: str) -> asyncio.Semaphore:
"""取该用户的后台并发额度信号量(懒建,每用户一把)。"""
sem = _user_background_sems.get(user_id)
if sem is None:
sem = asyncio.Semaphore(_BACKGROUND_SLOTS_PER_USER)
_user_background_sems[user_id] = sem
return sem
# ── 进程内底层操作 ───────────────────────────────────────────────────────────
async def _create_thread(app) -> str:
"""同进程建一个带圆桌 metadata 的 thread + 空 checkpoint(镜像 _create_thread_direct)。"""
checkpointer = app.state.checkpointer
thread_store = app.state.thread_store
thread_id = str(uuid.uuid4())
await thread_store.create(thread_id, assistant_id=None, metadata=dict(_ROUNDTABLE_THREAD_METADATA))
config = {"configurable": {"thread_id": thread_id, "checkpoint_ns": ""}}
await checkpointer.aput(
config,
empty_checkpoint(),
{"step": -1, "source": "input", "writes": None, "parents": {}, "created_at": datetime.now(UTC).isoformat()},
{},
)
return thread_id
async def _read_thread_messages(app, thread_id: str) -> list[dict[str, Any]]:
"""读 thread 最新 checkpoint 的 messages(已 serialize 成 dict 列表)。"""
checkpointer = app.state.checkpointer
config = {"configurable": {"thread_id": thread_id, "checkpoint_ns": ""}}
ct = await checkpointer.aget_tuple(config)
if ct is None:
return []
checkpoint = getattr(ct, "checkpoint", {}) or {}
values = serialize_channel_values(checkpoint.get("channel_values", {}) or {})
return values.get("messages") or []
def _latest_ai_message(messages: list[dict[str, Any]]) -> dict[str, Any] | None:
for m in reversed(messages):
if not isinstance(m, dict):
continue
if (m.get("type") or m.get("role")) in {"ai", "assistant"}:
return m
return None
def _message_text(message: dict[str, Any] | None) -> str:
if not message:
return ""
content = message.get("content")
if isinstance(content, str):
return content
if isinstance(content, list):
parts = [p.get("text", "") for p in content if isinstance(p, dict) and p.get("type") == "text"]
return "\n".join(parts)
return ""
def _is_broadcast_human_text(text: str) -> bool:
"""圆桌后广播进席位 thread 的 human,不算新一轮任务边界。"""
t = (text or "").strip()
return t.startswith("总控智能体收到") or t.startswith("子智能体")
def _latest_ai_text(messages: list[dict[str, Any]]) -> str:
"""取本轮**可见正文**(剥掉 ``<think>``)。
席位常以「调一堆工具」的 AIMessage 收尾(正文为空、只有 tool_calls,甚至 dangling),
直接取最后一条 AI 消息会拿到空串 → 报告就拿不到席位产物。这里向前找到真正有文字的那条。
思考-only 不算可见正文:一旦碰到最新一条有原文但剥 think 后为空的 AI,停止继承更早一轮。
后广播 human(``子智能体…`` / ``总控智能体收到``)不是轮次边界。
"""
for m in reversed(messages):
if not isinstance(m, dict):
continue
role = m.get("type") or m.get("role")
if role in {"human", "user"}:
if _is_broadcast_human_text(_message_text(m)):
continue
return ""
if role not in {"ai", "assistant"}:
continue
raw = _message_text(m)
kind = classify_seat_delivery(raw)
if kind == "valid":
return visible_seat_text(raw)
if kind == "thinking_only":
return ""
return ""
def _normalize_tool_calls(tool_calls: Any) -> list[dict[str, Any]]:
"""把 checkpoint 里的 tool_calls 归一化成 {name, args, id}(丢掉无名项)。"""
out: list[dict[str, Any]] = []
for tc in tool_calls or []:
if not isinstance(tc, dict):
continue
name = tc.get("name") or ""
if not name:
continue
out.append({"name": name, "args": tc.get("args") or {}, "id": tc.get("id") or ""})
return out
def _tool_calls_since_last_human(messages: list[dict[str, Any]]) -> list[dict[str, Any]]:
"""收集「最后一条 human 消息之后」各 AI 消息的 tool_calls(即**本轮**调用),保序。
席位 thread 跨轮累积历史,直接扫全量会把往轮的工具调用也算进来;以最后一条 human
(本轮 run 的输入)为界,只取其后的 AI 工具调用,对应前端「这一条气泡的步骤卡」。
"""
last_human = -1
for i, m in enumerate(messages):
if isinstance(m, dict) and (m.get("type") or m.get("role")) in {"human", "user"}:
last_human = i
out: list[dict[str, Any]] = []
for m in messages[last_human + 1 :]:
if not isinstance(m, dict):
continue
if (m.get("type") or m.get("role")) not in {"ai", "assistant"}:
continue
out.extend(_normalize_tool_calls(m.get("tool_calls")))
return out
async def _await_run_with_idle_guard(
app,
record,
*,
idle_timeout: float,
diag: DiagnosticsRecorder | None,
stage: str,
agent_name: str,
cycle: int | None,
max_total: float = 0.0,
) -> None:
"""等待一轮 run 完成,但用 **无活动(idle) 超时**而非总时长超时。
只要 run 仍在持续产出(StreamBridge 上有新事件:流式 token / 工具调用 / 进度),就一直
等下去——慢但活着的生成绝不打断。只有当 run **连续 ``idle_timeout`` 秒完全没有任何事件**
(真卡死)时才取消它并抛 ``RuntimeError``。``max_total`` > 0 时附加一个绝对总时长上限。
活动信号取自 StreamBridge(worker 每产一块都 ``bridge.publish``)。拿不到 bridge / run_id
时退化为「总时长超时(若配置)或直接等待」。
"""
task = getattr(record, "task", None)
if task is None:
return
run_id = getattr(record, "run_id", None)
bridge = getattr(app.state, "stream_bridge", None)
loop = asyncio.get_running_loop()
start = loop.time()
# 没有 bridge 可观察活动 → 退化:有 max_total 用之,否则直接等待(保持旧的「不卡死则一直等」)。
if bridge is None or run_id is None:
if max_total and max_total > 0:
await asyncio.wait_for(asyncio.shield(task), timeout=max_total)
else:
await task
return
from deerflow.runtime.stream_bridge.base import END_SENTINEL, HEARTBEAT_SENTINEL
last_activity = loop.time()
stuck = asyncio.Event()
async def _watch() -> None:
nonlocal last_activity
async for ev in bridge.subscribe(run_id, heartbeat_interval=_IDLE_CHECK_INTERVAL_S):
if ev is END_SENTINEL:
return # run 已收尾,正常退出(不算卡死)
now = loop.time()
if ev is HEARTBEAT_SENTINEL:
if now - last_activity >= idle_timeout:
stuck.set()
return
if max_total and max_total > 0 and now - start >= max_total:
stuck.set()
return
else:
last_activity = now # 任何真实事件 = 还活着
watch = asyncio.create_task(_watch())
try:
done, _pending = await asyncio.wait({task, watch}, return_when=asyncio.FIRST_COMPLETED)
if task in done:
return # run 正常结束(run_agent 内部已吞掉异常并落终态)
# watch 先完成:要么判定卡死,要么因 END 正常退出。
if not stuck.is_set():
await task # END 已出现,run 即将收尾,等它
return
# —— 判定卡死:取消该 run,落诊断,抛错 ——
idle_elapsed = loop.time() - last_activity
total = loop.time() - start
logger.error(
"roundtable %s turn '%s' stalled: no output for %.0fs (total %.0fs) on run %s",
stage, agent_name, idle_elapsed, total, run_id,
)
run_manager = getattr(app.state, "run_manager", None)
if run_manager is not None:
try:
await run_manager.cancel(run_id, action="interrupt")
except Exception: # noqa: BLE001
logger.warning("cancel stalled run %s failed", run_id, exc_info=True)
if diag is not None:
await diag.record(
stage=stage, level="error", event="seat_turn_timeout",
message=(
f"「{agent_name}」连续 {idle_elapsed:.0f}s 无任何输出、判定为卡死并已取消"
f"(idle 上限 {idle_timeout:.0f}s;本轮已运行 {total:.0f}s)。常见原因:内网模型端点连接"
"被打满 / sqlite checkpoint 在挂载卷上锁挂死。注意:慢但持续产出的生成不会被打断。"
),
detail={"idle_seconds": round(idle_elapsed, 1), "total_seconds": round(total, 1), "idle_timeout": idle_timeout},
agent_name=agent_name, cycle=cycle,
)
raise RuntimeError(
f"roundtable seat turn '{agent_name}' stalled (no output for {idle_timeout:.0f}s)"
) from None
finally:
if not watch.done():
watch.cancel()
try:
await watch
except BaseException: # noqa: BLE001 — 取消 watch 的异常无关紧要
pass
async def _run_agent_turn(
app,
*,
user_id: str,
thread_id: str,
agent_name: str,
message: str,
model: str | None,
excluded_tools: list[str],
skill_stop_names: list[str],
thinking_enabled: bool,
reasoning_effort: str,
subagent_enabled: bool = False,
diag: DiagnosticsRecorder | None = None,
stage: str = "step2_seat",
cycle: int | None = None,
) -> list[dict[str, Any]]:
"""进程内驱动 agent_name 在 thread_id 上跑一轮(发 message),阻塞到完成,返回最终 messages。
镜像 multi_agent ``_run_payload`` 的 body 形状:agent 通过 ``context.agent_name`` 选,
assistant_id 固定 lead_agent;关闭 memory recall / subagent / plan mode(临时研讨任务)。
并发受**该用户的后台额度**约束(``_user_background_semaphore``):真正开跑(start_run +
await run task)前先取额度,跑完释放。同一用户的所有后台作业共享这几个额度,避免 2+
后台挂起把单连接 SQLite checkpointer / 离线模型打满、拖垮前台实时会商。
"""
from app.gateway.routers.thread_runs import RunCreateRequest
from app.gateway.services import start_run
body = RunCreateRequest(
assistant_id="lead_agent",
input={"messages": [{"type": "human", "content": [{"type": "text", "text": message}], "additional_kwargs": {}}]},
metadata={"run_scope": RUN_SCOPE_BACKGROUND_CHILD},
config={"recursion_limit": 1000},
context={
"agent_name": agent_name,
"model_name": model or "",
"mode": "pro",
"reasoning_effort": reasoning_effort,
"thinking_enabled": thinking_enabled,
"is_plan_mode": False,
# ultra 席位才置真 → 解锁 task 子代理工具(席位 policy 未禁 task)。leader 恒 False。
"subagent_enabled": subagent_enabled,
"thread_id": thread_id,
"memory_recall_disabled": True,
# 与 multi_agent._run_payload 对齐:关 builtin 记忆注入,避免单例总控/席位
# 被上轮会话 memory.json 带偏(臆造席位 / 直接干活)。详见该处注释。
"memory_injection_enabled": False,
},
excluded_tools=excluded_tools,
skill_stop_names=skill_stop_names,
on_disconnect="continue",
multitask_strategy="reject",
)
# 虚拟 request:带上空 ``state``(auth 中间件在真请求上会写 state.user / state.auth;
# 进程内没有这层,下游凡是读 request.state.X 的代码必须能拿到 None 而不是 AttributeError)。
request = SimpleNamespace(app=app, state=SimpleNamespace(user=None, auth=None))
# 取该用户的后台额度后再开跑:start_run + await task 是最重的部分(建 agent、调模型、
# 落 checkpoint)。额度满时本轮在此 await 排队,等同用户某个在跑轮次释放额度后再开。
async with _user_background_semaphore(user_id):
try:
record = await start_run( # type: ignore[arg-type]
body,
thread_id,
request,
run_scope=RUN_SCOPE_BACKGROUND_CHILD,
)
except Exception as exc: # noqa: BLE001 — start_run 报错(409 抢占 / DB 抖动等)也要落诊断
if diag is not None:
await diag.record(
stage=stage, level="error", event="start_run_failed",
message=f"启动 agent run 失败({agent_name}):{exc}",
detail={"type": type(exc).__name__, "error": str(exc), "thread_id": thread_id},
agent_name=agent_name, cycle=cycle,
)
raise
# idle 超时守护:慢但持续产出的生成一直等;只有连续 _SEAT_IDLE_TIMEOUT_S 秒**毫无输出**
# (真卡死)才取消并抛错。绝不因「这轮总耗时长」而打断(除非显式配了 _SEAT_MAX_TOTAL_S)。
await _await_run_with_idle_guard(
app, record,
idle_timeout=_SEAT_IDLE_TIMEOUT_S,
max_total=_SEAT_MAX_TOTAL_S,
diag=diag, stage=stage, agent_name=agent_name, cycle=cycle,
)
return await _read_thread_messages(app, thread_id)
def _read_report_html(app, thread_id: str, user_id: str) -> str | None:
"""从 report thread 的 outputs 目录读出报告 HTML。
优先 ``方案总览.html``;找不到就扫目录里**任意 .html**(取体积最大的那个),
容忍报告 agent 用了别的文件名。
"""
try:
outputs_dir = get_paths().sandbox_outputs_dir(thread_id, user_id=user_id)
preferred = outputs_dir / _REPORT_FILE
if preferred.exists():
return preferred.read_text(encoding="utf-8")
if outputs_dir.exists():
htmls = sorted(
outputs_dir.glob("*.html"),
key=lambda p: p.stat().st_size,
reverse=True,
)
if htmls:
logger.info("report html: 方案总览.html not found, using %s", htmls[0].name)
return htmls[0].read_text(encoding="utf-8")
logger.warning("report html: no .html found in %s", outputs_dir)
except Exception: # noqa: BLE001
logger.warning("read report html failed for thread %s", thread_id, exc_info=True)
return None
def _read_summary_md(app, thread_id: str, user_id: str) -> str | None:
"""从 summary thread 的 outputs 目录读出总结报告 Markdown。
优先 ``方案总结报告.md``;找不到就扫目录里**任意 .md**(取体积最大的那个),
容忍总结 agent 用了别的文件名。镜像 ``_read_report_html``。
"""
try:
outputs_dir = get_paths().sandbox_outputs_dir(thread_id, user_id=user_id)
preferred = outputs_dir / _SUMMARY_FILE
if preferred.exists():
return preferred.read_text(encoding="utf-8")
if outputs_dir.exists():
mds = sorted(
outputs_dir.glob("*.md"),
key=lambda p: p.stat().st_size,
reverse=True,
)
if mds:
logger.info("summary md: 方案总结报告.md not found, using %s", mds[0].name)
return mds[0].read_text(encoding="utf-8")
logger.warning("summary md: no .md found in %s", outputs_dir)
except Exception: # noqa: BLE001
logger.warning("read summary md failed for thread %s", thread_id, exc_info=True)
return None
# Step3 产物自愈重试次数(首轮之外的纠偏轮数)。
# 实测痛点(2026-06-12,h:\xiugai\后台日志.txt):deepseek 报告 agent 说「let me write
# it all at once」后发出超长 write_file,单次输出被 max_tokens 截断 → 工具调用从未执行、
# run 以 success 收尾但 outputs 下没有 .html → step3.html 回退成开场白文字。
# 自愈:检测产物缺失/截断后,同线程追加纠偏消息(要求 write_file 分段 + append 续写)重试。
#
# 取值依据(2026-06-15 修复「后台 HTML 看板只生成一半、概率渲染失败」):模型单轮输出受
# ``config.yaml`` 的 ``max_tokens``(deepseek-chat = 4096)硬截。一份完整看板 HTML 几十 KB
# ≈ 上万 token,单轮写不完;旧值 2(最多 3 轮 ≈ 12k token)对长报告不够 → 落库半截 HTML →
# 渲染失败。调到 5(最多 6 轮 append ≈ 24k token,足够绝大多数看板)显著降低截断概率。
# 循环命中完整产物(``</html>`` / 非空 md / 合法 json)即 break,短报告不会白跑额外轮次。
_ARTIFACT_RETRY_LIMIT = 5
_REPORT_RETRY_MISSING = (
f"检测到 /mnt/user-data/outputs/{_REPORT_FILE} 还没有生成——你上一轮的 write_file 没有执行成功"
"(最常见原因:试图一次写入整份长 HTML,输出超长被截断,工具调用没有真正落盘)。\n"
"请立即重新写入,并**严格分段**:\n"
f"1. 先用 write_file 写入 /mnt/user-data/outputs/{_REPORT_FILE} 的开头部分(<!doctype html> 起,约 200-300 行);\n"
"2. 再用 write_file 的 append=true 参数分多次续写后续片段,每段同样控制在 300 行以内;\n"
"3. 直到文档以 </html> 收尾。\n"
"不要在对话里粘贴 HTML 源码,写完用一两句话总结即可。"
)
_REPORT_RETRY_TRUNCATED = (
f"检测到 /mnt/user-data/outputs/{_REPORT_FILE} 已写入但**不完整**(没有 </html> 收尾,疑似上一轮输出被截断)。\n"
"请用 write_file 的 append=true 参数从中断处**继续补全剩余内容**(不要重写已有部分),"
"每段控制在 300 行以内,分多次续写直到文档以 </html> 收尾。不要在对话里粘贴源码。"
)
_SUMMARY_RETRY_MISSING = (
f"检测到 /mnt/user-data/outputs/{_SUMMARY_FILE} 还没有生成——你上一轮的 write_file 没有执行成功"
"(可能是单次输出过长被截断)。请立即用 write_file 重新写入完整报告;若内容很长,"
"先写前半部分、再用 write_file 的 append=true 参数分段续写。写完后把报告正文也输出到对话回复里。"
)
_DASHBOARD_RETRY = (
"你刚才没有输出可用的大屏数据(缺少 ```report-json 代码块,或其中的 JSON 无法解析、"
"字段不符合契约)。请重新**逐席位分析**研讨成果,在回复**最末尾**输出一个完整且**严格合法**"
"的 ```report-json 代码块:标准 JSON(双引号、无注释、无尾逗号、数字不加引号),seats 覆盖"
"每一个有交付的席位,dependsOn 只引用本清单内的席位 id;代码块之外只允许一两句说明;"
"不要调用任何工具、不要写文件、不要输出 HTML。"
)
def _extract_report_json(text: str) -> str | None:
"""从模型回复里健壮抽取大屏 report-json,返回**归一化后的** JSON 字符串(失败 → None)。
旧实现是「找最后一个 ```report-json 围栏 + json.loads」的简单粗暴抽取,模型一旦用
```json / 裸 JSON / 流式截断 / 尾逗号 就直接失败(即用户反馈的「大屏多页抽取 json 又失败」)。
现委托给 harness 的 :func:`extract_report_json`(多策略抽取 + schema 归一化,与前端
``report-template.ts`` 同源),返回的 JSON 字段齐全、可被前端直接 JSON.parse。
"""
return extract_report_json(text)
def _html_complete(html: str | None) -> bool:
"""报告 HTML 是否完整可用(非空 + 有 </html> 收尾,截断检测)。"""
return bool(html and html.strip()) and "</html>" in html.lower()
def _build_seat_intro(agents: list[SeatRef]) -> str:
"""镜像 init 注入 coordinator 的「可调度席位清单」。"""
lines = [f"- agent_name: {a.agent_id}({a.name})" for a in agents]
return (
"【圆桌系统初始化】本次研讨可调度的席位清单如下,请严格按 agent_name 通过 "
"agent_orchestration 技能派活,不要派给不在清单内的 agent:\n" + "\n".join(lines)
)
# 席位执行模式 → (thinking_enabled, reasoning_effort, subagent_enabled) 派生参数。
# 与前端 `lib/seat-mode.ts::seatModeToRunParams` **逐字段对齐**(务必同步修改),供每席位
# 推理深度覆盖(self._seat_modes)在 run_seat 里把 SeatMode 解析成具体运行参数。未知/缺省档
# 落 None(→ 跟随作业级 seat_* / policy 默认)。
_SEAT_MODE_PARAMS: dict[str, tuple[bool | None, str | None, bool]] = {
"flash": (False, "minimal", False),
"thinking": (True, "low", False),
"pro": (True, "medium", False),
"ultra": (True, "high", True),
}
def _exception_reason(exc: BaseException) -> str:
"""把单轮 run 抛出的异常归类成模型容错原因码(用于换模型提示)。
idle 守护判定的卡死(``_await_run_with_idle_guard`` 抛 ``... stalled ...``)→ ``timeout``;
其余(start_run 失败 / DB 抖动 / 其它)→ ``exception``。
"""
text = str(exc).lower()
if "stall" in text or "idle" in text or "timeout" in text or "timed out" in text:
return "timeout"
return "exception"
def _seat_empty_reason(messages: list[dict[str, Any]]) -> str | None:
"""席位/取数轮校验:本轮无可见正文 → ``empty`` / ``thinking_only``(触发换模型)。"""
saw_thinking = False
for m in reversed(messages):
if not isinstance(m, dict):
continue
role = m.get("type") or m.get("role")
if role in {"human", "user"}:
if _is_broadcast_human_text(_message_text(m)):
continue
break
if role not in {"ai", "assistant"}:
continue
kind = classify_seat_delivery(_message_text(m))
if kind == "valid":
return None
if kind == "thinking_only":
saw_thinking = True
break
return "thinking_only" if saw_thinking else "empty"
def _leader_failed_reason(messages: list[dict[str, Any]]) -> str | None:
"""总控轮校验:既无正文、也无任何工具调用(派活/澄清)→ ``empty``(视为模型没干活)。"""
ai = _latest_ai_message(messages)
has_text = bool(_latest_ai_text(messages).strip())
has_tools = bool(_normalize_tool_calls((ai or {}).get("tool_calls")))
return None if (has_text or has_tools) else "empty"
# ── 网关实现 ─────────────────────────────────────────────────────────────────
class InProcessRoundtableGateway:
"""进程内真网关。executor 先调 ``init_threads`` 建好线程,再跑编排引擎。"""
def __init__(
self,
app,
*,
user_id: str,
agents: list[SeatRef],
model: str | None = None,
seat_thinking_enabled: bool | None = None,
seat_reasoning_effort: str | None = None,
seat_subagent_enabled: bool = False,
seat_skill_directives: dict[str, bool] | None = None,
seat_modes: dict[str, str] | None = None,
job_id: str | None = None,
draft_id: str | None = None,
task_id: str | None = None,
) -> None:
self._app = app
self._user_id = user_id
self._agents = list(agents)
self._model = model
# 诊断日志记录器(绑定本作业上下文)。用进程级默认 store(deps.py 启动时已登记)。
self._diag = DiagnosticsRecorder(
get_default_store(),
scope="background",
job_id=job_id,
draft_id=draft_id,
task_id=task_id,
user_id=user_id,
)
# 席位执行模式(沿用 Step2 选择):覆盖 seat policy 的 thinking/reasoning(None → policy 默认),
# subagent 仅 ultra 档为真。只作用 special 席位;leader 派活不受影响。
self._seat_thinking_enabled = seat_thinking_enabled
self._seat_reasoning_effort = seat_reasoning_effort
self._seat_subagent_enabled = seat_subagent_enabled
# 每席位「技能识别强化」开关映射 {agent_id: bool}(沿用 Step2 所选业务链条里每个席位
# 的配置)。run_seat 按 agent_id 查表,与各席位 agent 自身的 config.yaml 开关取 OR,
# 传给 append_seat_skill_directive。默认空 dict(仅各 agent 自身开关生效)。
self._seat_skill_directives = dict(seat_skill_directives or {})
# 每席位「推理深度覆盖」映射 {agent_id: SeatMode}。run_seat 命中即用该档(_SEAT_MODE_PARAMS)
# 派生的 thinking/reasoning/subagent 覆盖作业级 seat_*;未命中则跟随作业级。默认空 dict。
self._seat_modes = dict(seat_modes or {})
self._seat_intro = ""
self._first_leader = True
async def init_threads(self) -> tuple[dict[str, str], str]:
"""建 coordinator + 各席位 thread,返回 (thread_ids, coordinator_name)。
席位清单通过「首轮 leader 消息前缀」注入(coordinator thread 跨轮累积历史,
首轮带上即长期可见),无需单独写 checkpoint。
"""
ensure_roundtable_functional_agents()
thread_ids: dict[str, str] = {}
for a in self._agents:
thread_ids[a.agent_id] = await _create_thread(self._app)
thread_ids[COORDINATOR_AGENT_ID] = await _create_thread(self._app)
self._seat_intro = _build_seat_intro(self._agents)
return thread_ids, COORDINATOR_AGENT_ID
async def _run_turn_with_fallback(
self,
*,
thread_id: str,
agent_name: str,
message: str,
requested_model: str | None,
excluded_tools: list[str],
skill_stop_names: list[str],
thinking_enabled: bool,
reasoning_effort: str,
subagent_enabled: bool = False,
stage: str,
cycle: int | None = None,
validate: Callable[[list[dict[str, Any]]], str | None] | None = None,
) -> tuple[list[dict[str, Any]], list[ModelSwitch], str | None]:
"""跑一轮 agent,并在「模型出问题」时自动切到下一个模型重试。
触发换模型的三种情形(与用户选择一致):
① ``_run_agent_turn`` 抛异常(start_run 失败 / idle 卡死超时);
② 返回了中间件兜底错误文案(``is_llm_error_message`` 命中);
③ 调用方 ``validate`` 判定无效(如席位交付为空 / 总控既无正文又无派活)。
候选模型链按 ``config.yaml models[]`` 顺序(``fallback_model_chain``),首个为
``requested_model``(空则作业级 ``self._model``,再空则配置默认)。返回
``(最终 messages, 发生的换模型记录, 实际所用模型)``。全部模型都失败时返回最后一次的
messages(让既有的「空交付 / 产物缺失」逻辑兜底,绝不抛异常拖垮整个作业)。
"""
base = requested_model if (requested_model not in (None, "")) else self._model
chain = fallback_model_chain(base) or [base]
switches: list[ModelSwitch] = []
messages: list[dict[str, Any]] = []
prev_model: str | None = None
prev_reason: str | None = None
used_model: str | None = chain[0]
for idx, candidate in enumerate(chain):
if idx > 0:
label = reason_text(prev_reason)
switches.append(ModelSwitch(failed_model=prev_model or "", next_model=candidate or "", reason_label=label))
logger.warning(
"roundtable %s '%s' switching model %s -> %s (%s)",
stage, agent_name, prev_model, candidate, prev_reason,
)
await self._diag.record(
stage=stage, level="warning", event="model_switch",
message=f"「{agent_name}」模型「{prev_model}」{label},已自动切换到「{candidate}」重试",
detail={"failed_model": prev_model, "next_model": candidate, "reason": prev_reason, "attempt": idx + 1},
agent_name=agent_name, cycle=cycle,
)
used_model = candidate
try:
messages = await _run_agent_turn(
self._app,
user_id=self._user_id,
thread_id=thread_id,
agent_name=agent_name,
message=message,
model=candidate,
excluded_tools=excluded_tools,
skill_stop_names=skill_stop_names,
thinking_enabled=thinking_enabled,
reasoning_effort=reasoning_effort,
subagent_enabled=subagent_enabled,
diag=self._diag,
stage=stage,
cycle=cycle,
)
except Exception as exc: # noqa: BLE001 — 模型故障:归类后尝试下一个模型
prev_model, prev_reason, messages = candidate, _exception_reason(exc), []
if idx + 1 < len(chain):
continue
await self._diag.record(
stage=stage, level="error", event="model_fallback_exhausted",
message=f"「{agent_name}」所有备选模型均失败(最后 {candidate}:{prev_reason})",
detail={"reason": prev_reason, "tried": chain}, agent_name=agent_name, cycle=cycle,
)
return messages, switches, used_model
# 跑通了:再看是否「兜底错误文案」或 validate 判定无效。
reason = is_llm_error_message(_latest_ai_message(messages))
if reason is None and validate is not None:
reason = validate(messages)
if reason is None:
return messages, switches, used_model
prev_model, prev_reason = candidate, reason
if idx + 1 >= len(chain):
await self._diag.record(
stage=stage, level="error", event="model_fallback_exhausted",
message=f"「{agent_name}」所有备选模型均未产出有效结果(最后 {candidate}:{reason})",
detail={"reason": reason, "tried": chain}, agent_name=agent_name, cycle=cycle,
)
return messages, switches, used_model
return messages, switches, used_model
async def run_leader(self, *, thread_ids, coordinator_name, message, model) -> LeaderTurn:
msg = message
if self._first_leader and self._seat_intro:
msg = self._seat_intro + "\n\n" + message
self._first_leader = False
policy = roundtable_run_policy("leader")
messages, switches, _used = await self._run_turn_with_fallback(
thread_id=thread_ids[coordinator_name],
agent_name=coordinator_name,
message=msg,
requested_model=model,
excluded_tools=policy.excluded_tools,
skill_stop_names=policy.skill_stop_names,
thinking_enabled=policy.thinking_enabled,
reasoning_effort=policy.reasoning_effort,
stage="step2_leader",
validate=_leader_failed_reason,
)
ai = _latest_ai_message(messages)
# tool_calls 取最后一条 AI 消息(派活就在这条);正文取最后一条有文字的 AI 消息。
final_text = _latest_ai_text(messages)
tool_calls = (ai.get("tool_calls") if ai else None) or []
# 本轮工具调用(派活 / 澄清),带进 LeaderTurn 供草稿 run 重建步骤卡。
turn_tool_calls = _normalize_tool_calls(tool_calls)
# 澄清优先
clar = next((tc for tc in tool_calls if (tc.get("name") or "").lower() == "ask_clarification"), None)
if clar is not None:
cargs = clar.get("args") or {}
question = str(cargs.get("question") or "").strip() or "总控需要补充信息,请补充后继续。"
return LeaderTurn(final_text=final_text, clarification=question, tool_calls=turn_tool_calls, switches=switches)
# 派活解析(复用 multi_agent 纯函数)
from app.gateway.routers.multi_agent import _extract_dispatch
dispatched: list[tuple[str, str]] = []
for tc in tool_calls:
try:
name, task = _extract_dispatch(tc)
if name and task:
dispatched.append((name, task))
except Exception: # noqa: BLE001
logger.warning("dispatch parse failed for tool_call", exc_info=True)
return LeaderTurn(final_text=final_text, dispatched=dispatched, tool_calls=turn_tool_calls, switches=switches)
async def run_seat(self, *, thread_ids, agent_id, task, model) -> SeatTurn:
policy = roundtable_run_policy("seat")
# 弱模型补偿:把该席位配置的技能(名称 + SKILL.md 路径 + 必须先 read_file 再按其
# 流程取信息的硬指令)显式拼进本轮任务,避免弱模型无视系统提示里的技能块(见
# roundtable_seat_skills 模块说明)。默认关:是否注入 = 该席位 agent 的 config.yaml
# 开关 OR 该席位在所选业务链条里的开关(按 agent_id 查表)。无配置技能 / 两开关都关时
# 原样返回 task。
seat_task = append_seat_skill_directive(
agent_id, task, chain_enabled=bool(self._seat_skill_directives.get(agent_id, False))
)
# 后台挂起路径与前台保持同一角色边界:席位只能执行具体任务,编排能力永远只属于总控。
seat_task = append_seat_role_boundary(seat_task)
# 每席位「推理深度覆盖」:业务链条里给该席位单独配了 seatMode → 用该档派生的
# thinking/reasoning/subagent **覆盖作业级**(self._seat_*);未配则跟随作业级。
# 作业级仍为 None 时落 policy 默认。与前端 useStep2Orchestration「覆盖优先、否则跟随
# selectedMode」逐席位语义一致。
seat_thinking, seat_reasoning, seat_subagent = (
self._seat_thinking_enabled,
self._seat_reasoning_effort,
self._seat_subagent_enabled,
)
override = _SEAT_MODE_PARAMS.get(self._seat_modes.get(agent_id, ""))
if override is not None:
seat_thinking, seat_reasoning, seat_subagent = override
messages, switches, _used = await self._run_turn_with_fallback(
thread_id=thread_ids[agent_id],
agent_name=agent_id,
message=seat_task,
requested_model=model,
excluded_tools=policy.excluded_tools,
skill_stop_names=policy.skill_stop_names,
thinking_enabled=(
policy.thinking_enabled if seat_thinking is None else seat_thinking
),
reasoning_effort=(
policy.reasoning_effort if seat_reasoning is None else seat_reasoning
),
subagent_enabled=seat_subagent,
stage="step2_seat",
validate=_seat_empty_reason,
)
return SeatTurn(
final_text=_latest_ai_text(messages),
tool_calls=_tool_calls_since_last_human(messages),
switches=switches,
)
async def create_report_thread(self) -> str:
ensure_roundtable_functional_agents()
return await _create_thread(self._app)
async def run_report(self, *, report_thread_id, prompt, model) -> ReportResult:
"""跑报告 agent,**带产物自愈**:首轮后 HTML 缺失/截断(无 </html>)→ 同线程追加
纠偏消息(分段 + append 续写)重试,至多 ``_ARTIFACT_RETRY_LIMIT`` 轮。"""
policy = roundtable_run_policy("report")
summary = ""
html: str | None = None
message = prompt
switches: list[ModelSwitch] = []
current_model = model # 跨产物自愈轮记住实际所用模型:一旦换到可用模型,后续续写轮就从它起。
for attempt in range(1 + _ARTIFACT_RETRY_LIMIT):
messages, sw, used = await self._run_turn_with_fallback(
thread_id=report_thread_id,
agent_name=REPORT_AGENT_ID,
message=message,
requested_model=current_model,
excluded_tools=policy.excluded_tools,
skill_stop_names=policy.skill_stop_names,
thinking_enabled=policy.thinking_enabled,
reasoning_effort=policy.reasoning_effort,
stage="step3_report",
)
switches.extend(sw)
current_model = used or current_model
summary = _latest_ai_text(messages) or summary
html = _read_report_html(self._app, report_thread_id, self._user_id)
if _html_complete(html):
break
if attempt >= _ARTIFACT_RETRY_LIMIT:
await self._diag.record(
stage="step3_report", level="error", event="report_artifact_incomplete",
message=f"结果绘制重试 {_ARTIFACT_RETRY_LIMIT} 次后 HTML 仍缺失/截断,回退为正文展示",
detail={"thread_id": report_thread_id, "has_partial_html": bool(html and html.strip())},
agent_name=REPORT_AGENT_ID,
)
break
kind = "truncated" if (html and html.strip()) else "missing"
message = _REPORT_RETRY_TRUNCATED if (html and html.strip()) else _REPORT_RETRY_MISSING
logger.warning(
"report html %s after attempt %d for thread %s; retrying with corrective prompt",
kind, attempt + 1, report_thread_id,
)
await self._diag.record(
stage="step3_report", level="warning", event=f"report_{kind}",
message=f"结果绘制第 {attempt + 1} 次产物{('被截断' if kind == 'truncated' else '缺失')},已追加纠偏提示重试",
detail={"thread_id": report_thread_id, "attempt": attempt + 1},
agent_name=REPORT_AGENT_ID,
)
# 截断的 HTML 也比纯文字强(浏览器容忍未闭合标签),仅在完全没有文件时才回退正文。
final_html = html if (html and html.strip()) else summary
return ReportResult(html=final_html, summary=summary, model=current_model or self._model, switches=switches)
async def create_summary_thread(self) -> str:
ensure_roundtable_functional_agents()
return await _create_thread(self._app)
async def run_summary(self, *, summary_thread_id, prompt, model) -> SummaryResult:
"""进程内驱动 roundtable-summary 写「方案总结报告.md」+ 对话正文(与前端
useStep3Summary 同一智能体 / 同一 run policy / 同款产物),返回与
DraftStep3SummarySnapshot 对齐的结果。线程是真实持久线程——加载草稿后可续问。
带产物自愈:md 文件缺失时同线程追加纠偏消息重试(同 ``run_report``)。"""
policy = roundtable_run_policy("summary")
closing = ""
md: str | None = None
message = prompt
switches: list[ModelSwitch] = []
current_model = model
for attempt in range(1 + _ARTIFACT_RETRY_LIMIT):
messages, sw, used = await self._run_turn_with_fallback(
thread_id=summary_thread_id,
agent_name=SUMMARY_AGENT_ID,
message=message,
requested_model=current_model,
excluded_tools=policy.excluded_tools,
skill_stop_names=policy.skill_stop_names,
thinking_enabled=policy.thinking_enabled,
reasoning_effort=policy.reasoning_effort,
stage="step3_summary",
)
switches.extend(sw)
current_model = used or current_model
closing = _latest_ai_text(messages) or closing
md = _read_summary_md(self._app, summary_thread_id, self._user_id)
if md and md.strip():
break
if attempt >= _ARTIFACT_RETRY_LIMIT:
await self._diag.record(
stage="step3_summary", level="error", event="summary_artifact_missing",
message=f"方案总结报告重试 {_ARTIFACT_RETRY_LIMIT} 次后 .md 仍缺失,回退为对话正文",
detail={"thread_id": summary_thread_id},
agent_name=SUMMARY_AGENT_ID,
)
break
message = _SUMMARY_RETRY_MISSING
logger.warning(
"summary md missing after attempt %d for thread %s; retrying with corrective prompt",
attempt + 1,
summary_thread_id,
)
await self._diag.record(
stage="step3_summary", level="warning", event="summary_missing",
message=f"方案总结报告第 {attempt + 1} 次未生成 .md,已追加纠偏提示重试",
detail={"thread_id": summary_thread_id, "attempt": attempt + 1},
agent_name=SUMMARY_AGENT_ID,
)
final_md = md if (md and md.strip()) else closing
return SummaryResult(md=final_md, closing=closing, model=current_model or self._model, thread_id=summary_thread_id, switches=switches)
async def run_dashboard_data(self, *, prompt, model) -> str | None:
"""大屏多页的结构化数据:进程内驱动**专职**「大屏多页数据智能体」
(``roundtable-dashboard``) 逐席位分析研讨成果,在对话里输出 ```report-json。
与旧版的差异(修复「大屏多页抽取 json 又失败」):
- 用**专职** dashboard 智能体(SOUL 专讲「逐席位分析 → 结构化数据」)而非复用
画 HTML 的 report 智能体;``dashboard`` policy **禁用一切工具**,杜绝弱模型跑去写文件。
- 抽取走 :func:`_extract_report_json` → harness 健壮多策略抽取 + schema 归一化,
容忍 ```json / 裸 JSON / 流式截断,并保证返回的 JSON 字段齐全可直接渲染。
- 仍保留同线程纠偏重试(缺失/坏 JSON 时追加 ``_DASHBOARD_RETRY``)。"""
ensure_roundtable_functional_agents()
thread_id = await _create_thread(self._app)
policy = roundtable_run_policy("dashboard")
message = prompt
current_model = model
for attempt in range(1 + _ARTIFACT_RETRY_LIMIT):
messages, _sw, used = await self._run_turn_with_fallback(
thread_id=thread_id,
agent_name=DASHBOARD_AGENT_ID,
message=message,
requested_model=current_model,
excluded_tools=policy.excluded_tools,
skill_stop_names=policy.skill_stop_names,
thinking_enabled=policy.thinking_enabled,
reasoning_effort=policy.reasoning_effort,
stage="step3_dashboard",
)
current_model = used or current_model
body = _extract_report_json(_latest_ai_text(messages))
if body:
return body
if attempt >= _ARTIFACT_RETRY_LIMIT:
await self._diag.record(
stage="step3_dashboard", level="error", event="dashboard_json_invalid",
message=f"大屏数据重试 {_ARTIFACT_RETRY_LIMIT} 次后仍抽取不到合法 report-json,放弃大屏产物",
detail={"thread_id": thread_id},
agent_name=DASHBOARD_AGENT_ID,
)
break
message = _DASHBOARD_RETRY
logger.warning(
"dashboard report-json missing/invalid after attempt %d for thread %s; retrying",
attempt + 1,
thread_id,
)
await self._diag.record(
stage="step3_dashboard", level="warning", event="dashboard_json_retry",
message=f"大屏数据第 {attempt + 1} 次抽取失败(缺 report-json / JSON 非法),已追加纠偏提示重试",
detail={"thread_id": thread_id, "attempt": attempt + 1},
agent_name=DASHBOARD_AGENT_ID,
)
return None