1010 lines
51 KiB
Python
1010 lines
51 KiB
Python
"""圆桌真网关(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
|