deerflow-code/frontend-web/src/roundtable-planning/api/multi-agent.ts
2026-09-07 18:24:55 +08:00

1216 lines
51 KiB
TypeScript
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.

import { apiFetch } from "@/core/api/fetch-client";
import { getBackendBaseURL } from "@/core/config";
import { OptType, writeYLog } from "@/core/log/y-log";
import { resolveAllowMultiple } from "@/roundtable-planning/lib/clarification";
import { reportRoundtableError, stageForRoundtableAgent } from "@/roundtable-planning/api/diagnostics-report";
import { createInlineReasoningDeltaFormatter } from "@/roundtable-planning/lib/stream-reasoning";
export type ThreadIdMap = Record<string, string>;
export interface SeatReadyEvent {
agentName: string;
threadId: string;
displayName: string;
description: string;
}
export interface CoordinatorReadyEvent {
mainAgentName: string;
threadId: string;
}
export interface InitMultiAgentRequest {
agentNames: string[];
mainAgentName: string;
model?: string;
/** Fires the moment each sub-agent's thread is ready (real completion order). */
onSeatReady?: (event: SeatReadyEvent) => void;
/** Fires when the coordinator's thread + agent have been created. */
onCoordinatorReady?: (event: CoordinatorReadyEvent) => void;
signal?: AbortSignal;
}
export interface InitMultiAgentResponse {
status: "success";
main_agent_name: string;
thread_ids: ThreadIdMap;
}
/**
* Open the /api/multi-agent/init SSE stream. Fires per-seat / coordinator
* callbacks as events arrive so the UI can light up roles individually, and
* resolves with the final {main_agent_name, thread_ids} once init_done lands.
* If the server emits an "error" frame mid-stream, the promise rejects.
*/
export async function initMultiAgent(
payload: InitMultiAgentRequest,
): Promise<InitMultiAgentResponse> {
// Timing telemetry for /api/multi-agent/init — enabled by default so we can
// see where time goes when users report "Step 2 init feels slow". Disable
// with `window.__MULTI_AGENT_INIT_TIMING__ = false` in devtools.
const timingEnabled =
typeof window === "undefined"
? false
: (window as unknown as { __MULTI_AGENT_INIT_TIMING__?: boolean }).__MULTI_AGENT_INIT_TIMING__ !== false;
const t0 = performance.now();
let tFirstByte = 0;
let tCoord = 0;
let seatCount = 0;
const tlog = (label: string) => {
if (!timingEnabled) return;
const elapsed = (performance.now() - t0).toFixed(0);
console.debug(`[init-timing] +${elapsed}ms ${label}`);
};
tlog("fetch.start");
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/init`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
agent_names: payload.agentNames,
main_agent_name: payload.mainAgentName,
model: payload.model,
}),
signal: payload.signal,
});
tlog(`fetch.headers status=${res.status}`);
if (!res.ok || !res.body) {
const text = await res.text().catch(() => "");
reportRoundtableError({
stage: "step2_init", event: "init_http_error",
message: `第二步会商初始化失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status,
detail: { agentNames: payload.agentNames, mainAgentName: payload.mainAgentName },
});
throw new Error(`init multi-agent failed (${res.status}): ${text}`);
}
const reader = res.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
let mainAgentName = "";
let threadIds: ThreadIdMap | null = null;
let aborted = false;
let errorDetail: string | null = null;
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
if (tFirstByte === 0) {
tFirstByte = performance.now();
tlog(`first-byte (TTFB-after-headers=${(tFirstByte - t0).toFixed(0)}ms)`);
}
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split("\n");
buffer = lines.pop() ?? "";
for (const rawLine of lines) {
const line = rawLine.trimEnd();
if (!line.startsWith("data:")) continue;
const data = line.slice(5).trim();
if (!data) continue;
let parsed: unknown;
try {
parsed = JSON.parse(data);
} catch {
continue;
}
if (!parsed || typeof parsed !== "object") continue;
const obj = parsed as Record<string, unknown>;
const event = typeof obj.event === "string" ? obj.event : "";
switch (event) {
case "seat_ready": {
seatCount += 1;
tlog(`seat_ready #${seatCount} ${String(obj.agent_name ?? "?")}`);
payload.onSeatReady?.({
agentName: String(obj.agent_name ?? ""),
threadId: String(obj.thread_id ?? ""),
displayName: String(obj.display_name ?? ""),
description: String(obj.description ?? ""),
});
break;
}
case "coordinator_ready": {
tCoord = performance.now();
tlog(`coordinator_ready (since-first-seat=${(tCoord - tFirstByte).toFixed(0)}ms)`);
payload.onCoordinatorReady?.({
mainAgentName: String(obj.main_agent_name ?? ""),
threadId: String(obj.thread_id ?? ""),
});
break;
}
case "init_done": {
tlog(`init_done seats=${seatCount}`);
mainAgentName = String(obj.main_agent_name ?? "");
threadIds = (obj.thread_ids ?? {}) as ThreadIdMap;
break;
}
case "error": {
errorDetail = String(obj.detail ?? "unknown init error");
break;
}
default:
// Unknown event — ignore for forward-compat.
break;
}
}
}
} catch (err) {
aborted = (err as Error)?.name === "AbortError" || payload.signal?.aborted === true;
if (!aborted) {
reportRoundtableError({
stage: "step2_init", event: "init_stream_exception",
message: `第二步会商初始化中断(网络不可达 / 连接中断 / 超时):${(err as Error)?.message ?? String(err)}`,
detail: { seatCount },
});
throw err;
}
}
tlog(`total ${(performance.now() - t0).toFixed(0)}ms (seats=${seatCount})`);
if (errorDetail) {
reportRoundtableError({
stage: "step2_init", event: "init_error_frame",
message: `第二步会商初始化报错:${errorDetail.slice(0, 800)}`,
});
throw new Error(`init multi-agent failed: ${errorDetail}`);
}
if (aborted) {
throw new DOMException("init aborted", "AbortError");
}
if (!threadIds || !mainAgentName) {
reportRoundtableError({
stage: "step2_init", event: "init_incomplete",
message: "第二步会商初始化流提前结束(未收到 init_done,可能连接被中断)",
detail: { seatCount, hasThreadIds: !!threadIds, mainAgentName },
});
throw new Error("init multi-agent stream ended without init_done");
}
return {
status: "success",
main_agent_name: mainAgentName,
thread_ids: threadIds,
};
}
export type StreamAgentType = "leader" | "special";
export type LeaderDispatch = [string, string];
/**
* High-level phases the SSE stream goes through. Used by the UI to show
* "what the backend is doing" while the user waits — most of the time before
* a first token is upstream/LLM latency, which feels like a stall otherwise.
*
* - `connecting` — fetch sent, no frames yet
* - `connected` — got the `event: metadata` frame (run_id assigned)
* - `streaming` — first AI text chunk arrived; tokens flowing
* - `tool_calling` — model emitted a tool_call chunk (e.g. agent_orchestration,
* present_files); usually means the run is about to interrupt
* - `tool_result` — a tool result message came back; brief, may go back to streaming
*/
export type StreamPhase = "connecting" | "connected" | "streaming" | "tool_calling" | "tool_result";
export interface StreamMultiAgentRequest {
agentType: StreamAgentType;
dAgentThreadId: ThreadIdMap;
agentName: string;
/**
* 人类可读名称(中文/英文)。仅 special 走广播时使用 —— 后端把它拼进
* 「子智能体X完成Y工作,交付内容为:...」里,这样总控 LLM 在下一轮看到自己
* 历史里的交付摘要时能用"语义化的名字"识别,而不是 UUID 字符串(LLM 对
* UUID 不友好,会导致重复派活)。后端缺失时 fallback 到 agentName。
*/
displayName?: string;
newMessage: string;
model?: string;
/**
* 仅 special 用。席位执行模式(快速/思考/专业问答/多智能体问答)派生的运行参数,
* 覆盖后端 `roundtable_run_policy` 的 thinking/reasoning 默认值。省略时后端按 policy
* 默认执行(向后兼容,旧前端不传不破坏)。三者通常一起传(见 `lib/seat-mode.ts`)。
*/
thinkingEnabled?: boolean;
reasoningEffort?: "minimal" | "low" | "medium" | "high";
/** 仅 special 用;ultra 档置真,给席位解锁 `task` 子代理工具。 */
subagentEnabled?: boolean;
/**
* 仅 special 用。会商席位「技能识别强化」开关(沿用所选业务链条配置)。置 true 时后端把该
* 席位配置的技能(含 SKILL.md 路径 + 必须先 read_file 的硬指令)拼进本轮任务,补偿弱模型。
* 与该席位 agent 自身 config.yaml 的同名开关取 OR。省略 = false(关)。
*/
seatSkillDirective?: boolean;
/**
* 深链接业务码(rwfx→6BF 等)。**仅方案总结报告(summary)单例消费**:非空 → 后端按业务链抽取规范
* 强制注入「完整性核对」指令,让总结核对产出相对业务链完不完整。普通会商省略 = 不注入。
*/
businessCode?: string | null;
skillStopNames?: string[];
/**
* 仅 leader 用。置 true 表示这是「派活后让总控判断继续/综合」的轮次:后端会把各已交付
* 席位的**完整交付正文**额外注入总控输入,供其写最终方案。派活决策轮(首轮等)传 false/省略,
* 后端只给摘要,使总控上下文精简、响应快。
*/
synthesisMode?: boolean;
/**
* 仅为兼容旧协议保留。总控 prompt 属于总控私有上下文,前端恒传 true,后端也会无条件隔离。
*/
suppressPreBroadcast?: boolean;
/**
* 仅 special 用。置 true 表示这是 DAG 并行批次里的席位:后端按 `seat_parallel`
* run policy **禁用** write_file / bash / str_replace / present_files 等写文件 /
* 沙箱工具(§4),避免同批多席位并发抢同一沙箱。产物以消息正文交付。后端缺失
* 该字段时按 false 处理(普通 seat policy),旧前端不受影响。
*/
parallelNoFile?: boolean;
/**
* 仅 special 用。前序席位的**完整交付** `[{name, content}]`。后端注入 run context
* (`roundtable_peer_deliveries`),席位用 `read_peer_delivery` 工具**按需**取全文
* (平时只看广播摘要)。不进 prompt → 不调工具就不耗 token。
*/
peerDeliveries?: Array<{ name: string; content: string }>;
/**
* 用户随干预消息上传的参考文件(已经过 `/api/threads/{thread_id}/uploads`)。
* 透传给后端 `additional_kwargs.files`,由 UploadsMiddleware 注入清单 + 解锁
* read_file 等工具。字段 `path` 用上传返回的 `virtual_path`。
*/
files?: Array<{ filename: string; size: number; path: string; status?: string }>;
signal?: AbortSignal;
/** Fires for every non-status data frame the upstream emits. */
onUpstreamFrame?: (frame: unknown) => void;
/**
* Fires for every AI text chunk parsed out of LangGraph messages-tuple frames.
* `text` is the incremental delta (concatenate to rebuild the message);
* `messageId` is stable per AI message so callers can group deltas.
*/
onTextDelta?: (text: string, messageId: string) => void;
/**
* Fires whenever the inferred SSE phase changes. The UI can use this to
* show progressive status messages during the inevitable first-token wait.
* `name` is the tool name when phase is `tool_calling` (or empty otherwise).
* `args` 是 tool_calling 阶段的工具入参(取自上游 LangGraph messages-tuple
* 帧中当前 tool call 的 args),供 UI 渲染路径 chip / 代码块等子节点。
* `toolCallId` 可将同名工具的连续调用与同一次调用的增量参数区分开。
*/
onPhaseChange?: (
phase: StreamPhase,
name?: string,
args?: Record<string, unknown>,
toolCallId?: string,
) => void;
/**
* 模型自动容错:后端检测到当前模型出问题(兜底错误文案 / 交付为空等),已自动切到下一个
* 模型并即将重跑本轮时触发。上层据此把该角色「正在流式的失败气泡」重置(清空已累计的增量
* 文本)并给用户一条「『X』模型异常,已切换『Y』继续」的提示。此回调可能多次触发(连续多个
* 模型都失败时)。触发后,本次 streamMultiAgent 仍会继续用新模型流式,最终照常 resolve。
*/
onModelSwitch?: (info: ModelSwitchInfo) => void;
}
/** 模型自动容错通知(`streamMultiAgent` 的 `onModelSwitch` 回调入参)。 */
export interface ModelSwitchInfo {
/** 角色:leader / seat / seat_parallel / report / summary / dashboard / ingest。 */
role: string;
/** 出问题的智能体 agent_name。 */
agentName: string;
/** 失败的模型名。 */
failedModel: string;
/** 已切换到的下一个模型名。 */
nextModel: string;
/** 失败原因(已翻译成简短中文,如「服务繁忙」)。 */
reason: string;
}
export interface ClarificationPayload {
/**
* Stable id for this clarification — the originating `ask_clarification`
* tool_call_id. Lets the UI key per-question selections when the agent asks
* several clarifications in one turn (Step 1 multi-question support).
* Omitted on the Step 2 leader path (single clarification per turn).
*/
id?: string;
/** Raw question text from the model's ask_clarification args. */
question: string;
/** Optional category hint from the model (e.g. "missing_info"). */
clarificationType: string;
/** Free-text reason the model needs clarification. */
context: string;
/** Suggested answer chips (may be empty). */
options: unknown[];
/** Whether the user is allowed to type a custom answer. */
allowCustom: boolean;
/** When false, option cards are single-select and send on click. Omitted → inferred from clarificationType. */
allowMultiple?: boolean;
}
export interface StreamMultiAgentResult {
/** For leader runs: dispatched [agentName, task] tuples. Empty = consensus. */
dispatched: LeaderDispatch[];
/** For special runs: the completion status string. */
message: string | null;
/** Final assistant text emitted by the underlying agent. */
content: string;
/** Echoed agent_name of the agent that produced this run. */
agentName: string | null;
/** 实际产出本轮结果的模型(发生模型自动容错时即最终切换到的模型);后端未下发时为 null。 */
modelUsed?: string | null;
/**
* Populated when the leader called ask_clarification — the run ended without
* dispatching, and the model is asking the user a question instead. When
* present, callers should render `content` (or `clarification.question` as
* a fallback) as a question and collect a user reply rather than treating
* an empty `dispatched` as consensus.
*/
clarification: ClarificationPayload | null;
/**
* 席位本轮无可见正文(仅思考 / 断流 / 空串)。此时 `content` 为空,调用方不得把
* 流式过程文本当成已交付,也不得写入 priorDeliveries。
*/
invalidDelivery?: { reason: string } | null;
}
/**
* Best-effort partial JSON parser for accumulated `tool_call_chunks.args`.
*
* LangChain emits tool args as a series of string fragments (one delta per
* chunk). To get progressive `args.path` / `args.content` during streaming —
* the way `useStream` does in the main chat — we concatenate the deltas and
* try to parse. When the JSON is still mid-token (e.g. `{"path":"/mnt`) we
* close any open string + brackets and retry, so callers can read partially
* available keys even before the full payload arrives.
*
* Returns null when nothing useful can be extracted yet.
*/
function parsePartialJson(s: string): Record<string, unknown> | null {
if (!s) return null;
try {
const v = JSON.parse(s);
return v && typeof v === "object" ? (v as Record<string, unknown>) : null;
} catch {
// fall through to best-effort close + retry
}
let trimmed = s;
// A dangling `\` would turn our injected closing `"` into an escape; drop it.
if (trimmed.endsWith("\\")) trimmed = trimmed.slice(0, -1);
let inString = false;
let escaped = false;
const stack: string[] = [];
for (const ch of trimmed) {
if (escaped) {
escaped = false;
continue;
}
if (ch === "\\") {
escaped = true;
continue;
}
if (inString) {
if (ch === '"') inString = false;
continue;
}
if (ch === '"') {
inString = true;
continue;
}
if (ch === "{") stack.push("}");
else if (ch === "[") stack.push("]");
else if (ch === "}" || ch === "]") stack.pop();
}
let closed = trimmed;
if (inString) {
closed += '"';
} else {
// Mid-key / dangling separator (`{"path":`) can't be balanced cleanly.
closed = closed.replace(/[,:]\s*$/, "");
}
while (stack.length) closed += stack.pop();
try {
const v = JSON.parse(closed);
return v && typeof v === "object" ? (v as Record<string, unknown>) : null;
} catch {
return null;
}
}
/**
* Identify the final status frame our backend appends after the LangGraph
* stream. Required keys are `status` (array for leader, string for special)
* and `content` (final assistant text); `agent_name` is also expected.
*/
function isFinalStatusFrame(value: unknown): value is {
status: LeaderDispatch[] | string;
content?: string;
agent_name?: string;
question?: string;
clarification_type?: string;
clarification_context?: string;
options?: unknown[];
allow_custom?: boolean;
allow_multiple?: boolean;
} {
if (!value || typeof value !== "object") return false;
if (!("status" in value) || !("content" in value)) return false;
const status = (value as { status: unknown }).status;
if (typeof status === "string") return true;
if (Array.isArray(status)) {
return status.every(
(item) =>
Array.isArray(item) && item.length === 2 && typeof item[0] === "string" && typeof item[1] === "string",
);
}
return false;
}
/**
* Extract an AI text delta from a LangGraph `messages-tuple` data frame.
* Frames look like `[chunk, metadata]`; an AI text chunk carries `type` in
* {"ai","AIMessageChunk"} and a `content` that is either a string or a list of
* structured content blocks. Returns null when the frame is not an AI text
* chunk (e.g., tool calls, tool results, value snapshots).
*/
function extractAiTextDelta(
parsed: unknown,
): { content: string; reasoning: string; id: string } | null {
if (!Array.isArray(parsed) || parsed.length < 1) return null;
const chunk = parsed[0];
if (!chunk || typeof chunk !== "object") return null;
const obj = chunk as Record<string, unknown>;
const type = String(obj.type ?? "").toLowerCase();
// Accept any AI-shaped message chunk LangChain might emit across versions:
// "ai" (older), "AIMessageChunk" (newer streaming), or plain "AIMessage".
if (type !== "ai" && type !== "aimessagechunk" && type !== "aimessage") return null;
const content = obj.content;
let text = "";
if (typeof content === "string") {
text = content;
} else if (Array.isArray(content)) {
for (const item of content) {
if (typeof item === "string") text += item;
else if (item && typeof item === "object") {
const block = item as Record<string, unknown>;
const blockText = block.text;
if (typeof blockText === "string") text += blockText;
}
}
}
const additionalKwargs = obj.additional_kwargs;
const directReasoning = obj.reasoning_content ?? obj.reasoning;
const nestedReasoning =
additionalKwargs && typeof additionalKwargs === "object"
? (additionalKwargs as Record<string, unknown>).reasoning_content ??
(additionalKwargs as Record<string, unknown>).reasoning
: undefined;
const rawReasoning = nestedReasoning ?? directReasoning;
const reasoning = typeof rawReasoning === "string" ? rawReasoning : "";
if (!text && !reasoning) return null;
const id = typeof obj.id === "string" ? obj.id : "";
return { content: text, reasoning, id };
}
/**
* Open an SSE connection to /api/multi-agent/run/stream and resolve once the
* stream completes. The final frame (`{status: ...}`) is interpreted; every
* other `data:` frame is forwarded to `onUpstreamFrame` so the caller can do
* its own rendering (we don't try to mimic LangGraph SDK semantics here).
*/
export async function streamMultiAgent(
req: StreamMultiAgentRequest,
): Promise<StreamMultiAgentResult> {
// 整段包一层兜底:除已在具体分支上报的(HTTP 错误 / error 帧 / worker error 事件)外,
// 再兜住 fetch 直接 reject(网络不可达 / 连接被拒)与流式中途中断(连接重置 / 超时),
// 做到「能收集到的报错越全越好」。`reported` 防止与具体分支重复上报。
let reported = false;
try {
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/run/stream`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
agent_type: req.agentType,
d_agent_thread_id: req.dAgentThreadId,
agent_name: req.agentName,
// 仅 special 用到;后端缺失时 fallback 到 agent_name,旧前端也不破坏。
display_name: req.displayName,
new_message: req.newMessage,
model: req.model,
// 席位模式派生参数(仅 special);undefined 时不带上 → 后端按 policy 默认执行。
thinking_enabled: req.thinkingEnabled,
reasoning_effort: req.reasoningEffort,
subagent_enabled: req.subagentEnabled,
// 仅 special 用;会商技能识别强化(沿用所选业务链条开关)。省略 → 后端按 false(关)处理。
seat_skill_directive: req.seatSkillDirective ?? false,
// 仅 summary 单例用;深链接业务码 → 后端注入业务链完整性核对。省略/空 → 不注入。
business_code: req.businessCode ?? null,
skill_stop_names: req.skillStopNames ?? [],
// 仅 leader 用;综合轮让后端注入各席位完整交付正文。省略 → 后端按 false 处理。
synthesis_mode: req.synthesisMode ?? false,
// 总控 prompt 永不进入席位 thread;后端同样强制执行,避免旧调用方破坏角色隔离。
suppress_pre_broadcast: req.suppressPreBroadcast ?? true,
// 仅 special 用;DAG 并行批次禁写文件(seat_parallel policy)。省略 → 后端按 false 处理。
parallel_no_file: req.parallelNoFile ?? false,
// 仅 special 用;前序席位完整交付,后端注入 context 供 read_peer_delivery 按需取全文。
peer_deliveries: req.peerDeliveries ?? [],
// 用户随干预上传的参考文件;仅在有文件时带上,后端注入 additional_kwargs.files。
files: req.files && req.files.length > 0 ? req.files : undefined,
}),
signal: req.signal,
});
if (!res.ok || !res.body) {
const text = await res.text().catch(() => "");
// 用户实际遇到的报错(含进第三步常见的 502):HTTP 层就失败、SSE 都没建起来。
reportRoundtableError({
stage: stageForRoundtableAgent(req.agentType, req.agentName),
event: "run_stream_http_error",
message: `「${req.displayName || req.agentName}」调用失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status,
agentId: req.agentName,
agentName: req.displayName || req.agentName,
detail: { agentType: req.agentType },
});
reported = true;
throw new Error(`run/stream failed (${res.status}): ${text}`);
}
const reader = res.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
// Track the last `event:` line so logging can show which event a `data:`
// belongs to. The frontend doesn't actually dispatch by event name — every
// data frame goes through the same parser — but the label is invaluable
// for diagnosing why a sub-agent stream "looks empty".
let lastEvent = "";
let frameSeq = 0;
let gotFinalStatus = false;
const debugLabel = `[stream:${req.agentType}:${req.agentName}]`;
// Stream 帧日志(frame# / delta)仅 opt-in 开启。默认关闭,避免 Step 2 流式期间
// 刷屏(含 wujie 嵌入时子应用 devtools);需要排查 SSE 时在控制台执行:
// `window.__MULTI_AGENT_DEBUG__ = true`
const debugEnabled =
typeof window !== "undefined" &&
(window as unknown as { __MULTI_AGENT_DEBUG__?: boolean }).__MULTI_AGENT_DEBUG__ === true;
const result: StreamMultiAgentResult = {
dispatched: [],
message: null,
content: "",
agentName: null,
clarification: null,
modelUsed: null,
invalidDelivery: null,
};
// Inferred SSE phase + a tiny helper so we only fire onPhaseChange when the
// value actually changes (the reader loop sees many frames per phase).
//
// `write_file` 的 `content`、`str_replace` 的 old_str/new_str 会在每个
// tool_call chunk 中持续增长。它们属于实时预览数据,已经由 onUpstreamFrame
// 写入 stub;若用它们判断 phase 变化,会在单次工具调用里反复触发 React 状态
// 更新,严重时会形成 Maximum update depth loop。因此这里只比较工具的身份入参。
const phaseArgsSignature = (args?: Record<string, unknown>) => {
if (!args) return "";
const streamedPayloadKeys = new Set([
"content",
"old_str",
"new_str",
"patch",
"diff",
]);
try {
return JSON.stringify(
Object.fromEntries(
Object.entries(args)
.filter(([key]) => !streamedPayloadKeys.has(key))
.sort(([left], [right]) => left.localeCompare(right)),
),
);
} catch {
return "[unserializable]";
}
};
let currentPhase: StreamPhase = "connecting";
let currentPhaseName = "";
let currentPhaseArgsSignature = "";
let currentPhaseToolCallId = "";
req.onPhaseChange?.("connecting");
// The normal chat renders provider reasoning from a distinct message field.
// Roundtable's bespoke SSE reader keeps that field inline so its existing
// `<think>` renderer remains the single display path for Steps 2 and 3.
const inlineAiDelta = createInlineReasoningDeltaFormatter();
const setPhase = (
next: StreamPhase,
name?: string,
args?: Record<string, unknown>,
toolCallId?: string,
) => {
const nextName = name ?? "";
const nextArgsSignature = phaseArgsSignature(args);
const nextToolCallId = toolCallId ?? "";
// 同 phase、同工具、同身份入参的帧只是流式重复帧,不通知 UI。保留 path/filepaths
// 等身份变化,确保工具参数从不完整变完整时仍可自动打开沙箱。
if (
next === currentPhase &&
nextName === currentPhaseName &&
nextArgsSignature === currentPhaseArgsSignature &&
nextToolCallId === currentPhaseToolCallId
) {
return;
}
currentPhase = next;
currentPhaseName = nextName;
currentPhaseArgsSignature = nextArgsSignature;
currentPhaseToolCallId = nextToolCallId;
req.onPhaseChange?.(next, name, args, toolCallId);
};
/**
* 累积 tool_call_chunks(delta args 字符串) —— 与 LangChain `AIMessageChunk`
* 合并的内部行为对齐,补主聊天靠 `useStream` 暗中做的那部分工作。Key 是 chunk.index
* (并发 tool 调用通常 0/1/2 区分);chunk.id 变化时视作"同一 index 上的新调用"并重置。
* 每帧累积后用 `parsePartialJson` 尽力 parse,成功就把合成的 `tool_calls[]` 注入
* 到 messages-tuple 帧的 chunkObj 上,让上层 `onUpstreamFrame` 和 `onPhaseChange`
* 都能拿到完整的 `args.path` / `args.content` 来做"流式期间打开沙箱 + 实时预览"。
*/
const toolCallAcc = new Map<
number,
{ id?: string; name: string; argsStr: string }
>();
type SyntheticToolCall = {
id?: string;
name: string;
args: Record<string, unknown>;
};
const ingestToolCallChunks = (chunks: unknown): void => {
if (!Array.isArray(chunks)) return;
for (const tcc of chunks) {
if (!tcc || typeof tcc !== "object") continue;
const tccObj = tcc as Record<string, unknown>;
const idx = typeof tccObj.index === "number" ? tccObj.index : 0;
let entry = toolCallAcc.get(idx);
if (!entry) {
entry = { name: "", argsStr: "" };
toolCallAcc.set(idx, entry);
}
if (typeof tccObj.id === "string" && tccObj.id) {
if (entry.id && entry.id !== tccObj.id) {
// 同一 index 上换了新 tool_call,重置累积(不太常见,defensive)。
entry.name = "";
entry.argsStr = "";
}
entry.id = tccObj.id;
}
if (typeof tccObj.name === "string" && tccObj.name) entry.name = tccObj.name;
if (typeof tccObj.args === "string") entry.argsStr += tccObj.args;
}
};
const buildSyntheticToolCalls = (): SyntheticToolCall[] => {
const out: SyntheticToolCall[] = [];
for (const entry of toolCallAcc.values()) {
if (!entry.name) continue;
const parsed = parsePartialJson(entry.argsStr);
if (parsed) {
out.push({ id: entry.id, name: entry.name, args: parsed });
}
}
return out;
};
while (true) {
const { done, value } = await reader.read();
if (done) break;
const chunk = decoder.decode(value, { stream: true });
buffer += chunk;
const lines = buffer.split("\n");
buffer = lines.pop() ?? "";
for (const rawLine of lines) {
const line = rawLine.trimEnd();
if (!line) continue;
if (line.startsWith("event:")) {
lastEvent = line.slice(6).trim();
// First metadata frame = run accepted by upstream. Move phase out
// of "connecting" so the UI can switch from "调度中" to "已连接".
if (lastEvent === "metadata" && currentPhase === "connecting") {
setPhase("connected");
}
continue;
}
if (!line.startsWith("data:")) continue;
const payload = line.slice(5).trim();
if (!payload) continue;
let parsed: unknown;
try {
parsed = JSON.parse(payload);
} catch {
if (debugEnabled) console.debug(debugLabel, `bad JSON in event=${lastEvent}:`, payload.slice(0, 200));
continue;
}
frameSeq += 1;
if (debugEnabled) console.debug(debugLabel, `frame#${frameSeq} event=${lastEvent}`, parsed);
// worker 把 run 内异常(常见:第二步席位/Step3 模型调用失败)发成 `event: error`
// SSE 帧 `{message,name}`——它既不是 HTTP 错误也不是 `{status:"error"}` 终帧,
// 单独在这里捕获上报,否则这类「模型没调用成功」的报错会漏收集。
if (lastEvent === "error") {
const obj = parsed && typeof parsed === "object" ? (parsed as Record<string, unknown>) : { raw: parsed };
reportRoundtableError({
stage: stageForRoundtableAgent(req.agentType, req.agentName),
event: "run_worker_error_event",
message: `「${req.displayName || req.agentName}」运行报错:${String(obj.message ?? obj.name ?? "").slice(0, 800)}`,
agentId: req.agentName,
agentName: req.displayName || req.agentName,
detail: { ...obj, agentType: req.agentType, framesBeforeError: frameSeq },
});
}
// 模型自动容错帧:后端检测到当前模型出问题,已切到下一个模型并即将重跑本轮。
// 不是终态帧(无 content)→ 不结束流。清空已累计的(失败尝试的)tool_call 增量,
// 把 phase 退回 connected,并通知上层重置该角色气泡 + 给用户提示。随后新模型的
// 流式增量会照常通过 onTextDelta 继续累计到(已被上层清空的)气泡。
if (
parsed &&
typeof parsed === "object" &&
(parsed as Record<string, unknown>).status === "model_switch"
) {
const obj = parsed as Record<string, unknown>;
toolCallAcc.clear();
currentPhase = "connected";
req.onModelSwitch?.({
role: typeof obj.role === "string" ? obj.role : "",
agentName: typeof obj.agent_name === "string" ? obj.agent_name : req.agentName,
failedModel: typeof obj.failed_model === "string" ? obj.failed_model : "",
nextModel: typeof obj.next_model === "string" ? obj.next_model : "",
reason: typeof obj.reason === "string" ? obj.reason : "",
});
continue;
}
// Phase inference from LangGraph messages-tuple frames. We only flip
// to tool_calling/tool_result on chunks that clearly indicate them;
// everything else with non-empty AI text means we're streaming.
if (lastEvent === "messages" && Array.isArray(parsed) && parsed.length > 0) {
const chunk = parsed[0];
if (chunk && typeof chunk === "object") {
const chunkObj = chunk as Record<string, unknown>;
const chunkType = String(chunkObj.type ?? "").toLowerCase();
const toolCalls = chunkObj.tool_calls;
const toolCallChunks = chunkObj.tool_call_chunks;
const content = chunkObj.content;
const streamDelta = extractAiTextDelta(parsed);
// 1) 累积本帧的 tool_call_chunks delta(若有),再合成可能可读的 tool_calls。
// 主聊天 useStream 在 SDK 层做了这件事,我们在自定义 SSE 里手动补上。
ingestToolCallChunks(toolCallChunks);
// Some providers serialize `tool_calls` before their `args` have
// been decoded, leaving args as the raw JSON string. Treating that
// as a complete call suppresses the accumulated tool_call_chunks,
// so write_file reaches the sandbox with a path but no content.
const hasCompleteToolCalls =
Array.isArray(toolCalls) &&
toolCalls.length > 0 &&
toolCalls.every(
(toolCall) =>
!!toolCall &&
typeof toolCall === "object" &&
!!(toolCall as Record<string, unknown>).args &&
typeof (toolCall as Record<string, unknown>).args === "object",
);
let synthetic: SyntheticToolCall[] = [];
if (!hasCompleteToolCalls && toolCallAcc.size > 0) {
synthetic = buildSyntheticToolCalls();
if (synthetic.length > 0) {
// 把累积出的 tool_calls 注入回帧的 chunkObj —— 这样下面 onUpstreamFrame
// 那边的 handleUpstreamFrame 只读 tool_calls 字段就能拿到完整 args,
// 不必再单独做一份累积。
chunkObj.tool_calls = synthetic;
}
}
const effectiveToolCalls: unknown[] | null = hasCompleteToolCalls
? (toolCalls as unknown[])
: synthetic.length > 0
? synthetic
: null;
const hasToolCall =
(effectiveToolCalls !== null && effectiveToolCalls.length > 0) ||
(Array.isArray(toolCallChunks) && toolCallChunks.length > 0);
if (chunkType === "tool") {
setPhase("tool_result");
} else if (hasToolCall) {
// Pick the active (latest) call's name + args. A leader can dispatch
// several seats with the same `agent_orchestration` tool in one run;
// choosing the first accumulated call would leave the UI stuck on an
// earlier assignment until the final response arrives.
// 优先用累积出的 / 完整的 tool_calls(含可读 args 对象);若都还不到火候,
// 回退到 tool_call_chunks 拿名字(args 还是 string,不透传)。
const activeCall =
effectiveToolCalls && effectiveToolCalls.length > 0
? (effectiveToolCalls[effectiveToolCalls.length - 1] as Record<string, unknown>)
: Array.isArray(toolCallChunks) && toolCallChunks.length > 0
? (toolCallChunks[toolCallChunks.length - 1] as Record<string, unknown>)
: null;
const toolName =
activeCall && typeof activeCall.name === "string" ? activeCall.name : "";
const args =
activeCall &&
typeof activeCall.args === "object" &&
activeCall.args !== null
? (activeCall.args as Record<string, unknown>)
: undefined;
const toolCallId =
activeCall && typeof activeCall.id === "string" ? activeCall.id : undefined;
setPhase("tool_calling", toolName, args, toolCallId);
} else if (
(chunkType === "ai" || chunkType === "aimessagechunk" || chunkType === "aimessage") &&
(Boolean(streamDelta) ||
(typeof content === "string" && content.length > 0) ||
(Array.isArray(content) && content.length > 0))
) {
setPhase("streaming");
}
}
}
if (isFinalStatusFrame(parsed)) {
gotFinalStatus = true;
if (Array.isArray(parsed.status)) {
result.dispatched = parsed.status;
} else {
result.message = parsed.status;
// "clarification" status carries the leader's ask_clarification
// payload. Extract the structured fields so the page can render
// a question (or chip choices) instead of mistaking the run for
// a consensus.
if (parsed.status === "clarification") {
const clarificationType =
typeof parsed.clarification_type === "string" ? parsed.clarification_type : "";
result.clarification = {
question: typeof parsed.question === "string" ? parsed.question : "",
clarificationType,
context:
typeof parsed.clarification_context === "string" ? parsed.clarification_context : "",
options: Array.isArray(parsed.options) ? parsed.options : [],
allowCustom: parsed.allow_custom !== false,
allowMultiple: resolveAllowMultiple(
typeof parsed.allow_multiple === "boolean" ? parsed.allow_multiple : undefined,
clarificationType,
),
};
}
if (parsed.status === "invalid_delivery") {
const reasonRaw = (parsed as Record<string, unknown>).reason;
result.invalidDelivery = {
reason: typeof reasonRaw === "string" && reasonRaw ? reasonRaw : "empty",
};
result.content = "";
}
}
if (typeof parsed.content === "string" && !result.invalidDelivery) {
result.content = parsed.content;
}
if (typeof parsed.agent_name === "string") {
result.agentName = parsed.agent_name;
}
{
const mu = (parsed as Record<string, unknown>).model_used;
if (typeof mu === "string" && mu) result.modelUsed = mu;
}
// SSE error 帧(后端 stream_generator 在席位/Step3 run 失败时下发
// `{status:"error", content, agent_name}`):上报,便于收集用户侧报错。
if (result.message === "error") {
reportRoundtableError({
stage: stageForRoundtableAgent(req.agentType, result.agentName || req.agentName),
event: "run_stream_error_frame",
message: `「${req.displayName || req.agentName}」运行报错:${(result.content || "").slice(0, 800)}`,
agentId: result.agentName || req.agentName,
agentName: req.displayName || req.agentName,
detail: { agentType: req.agentType, framesBeforeError: frameSeq },
});
}
continue;
}
const delta = extractAiTextDelta(parsed);
if (delta) {
const text = inlineAiDelta(delta);
if (debugEnabled) console.debug(debugLabel, `→ delta (${text.length} chars)`);
if (text) req.onTextDelta?.(text, delta.id);
}
req.onUpstreamFrame?.(parsed);
}
}
if (debugEnabled) {
console.debug(debugLabel, `stream ended after ${frameSeq} frames; final content length=${result.content.length}`);
}
if (req.agentType === "special" && !gotFinalStatus) {
result.invalidDelivery = { reason: "truncated" };
result.content = "";
}
return result;
} catch (err) {
if (!reported && req.signal?.aborted !== true && (err as Error)?.name !== "AbortError") {
reportRoundtableError({
stage: stageForRoundtableAgent(req.agentType, req.agentName),
event: "run_stream_exception",
message: `「${req.displayName || req.agentName}」调用异常(网络不可达 / 连接中断 / 超时):${(err as Error)?.message ?? String(err)}`,
agentId: req.agentName,
agentName: req.displayName || req.agentName,
detail: { agentType: req.agentType },
});
}
throw err;
}
}
export interface SessionArtifactFile {
/** Virtual path under the sandbox, e.g. `mnt/user-data/outputs/foo.md`. */
path: string;
/** Bare filename (basename of `path`). */
name: string;
sizeBytes: number;
/** MIME type guessed from the extension, or null when unknown. */
mimeType: string | null;
/** ISO 8601 UTC timestamp of last modification. */
modifiedAt: string;
}
export interface SessionArtifactsResponse {
threadId: string;
files: SessionArtifactFile[];
}
/**
* List every file written to `/mnt/user-data/outputs/` for a thread. Used by
* the Step 3 final-report view to surface a thread's produced artifacts so the
* user can preview/download them via {@link buildArtifactDownloadUrl}.
*
* Returns an empty `files` array when the thread produced no artifacts (the
* backend treats a missing outputs directory as empty, not 404).
*/
export async function listSessionArtifacts(
threadId: string,
options: { signal?: AbortSignal } = {},
): Promise<SessionArtifactsResponse> {
const res = await apiFetch(
`${getBackendBaseURL()}/api/threads/${encodeURIComponent(threadId)}/artifacts`,
{ method: "GET", signal: options.signal },
);
if (!res.ok) {
const text = await res.text().catch(() => "");
throw new Error(`list artifacts failed (${res.status}): ${text}`);
}
const json = (await res.json()) as {
thread_id: string;
files: Array<{
path: string;
name: string;
size_bytes: number;
mime_type: string | null;
modified_at: string;
}>;
};
const result = {
threadId: json.thread_id,
files: (json.files ?? []).map((f) => ({
path: f.path,
name: f.name,
sizeBytes: f.size_bytes,
mimeType: f.mime_type,
modifiedAt: f.modified_at,
})),
};
writeYLog({
optType: OptType.search,
optDetail: "CM助手-全量-会商产物-查询",
url: `/api/threads/${threadId}/artifacts`,
reqParam: { threadId },
returnParam: { count: result.files.length },
});
return result;
}
/**
* Build a URL that previews (inline) or downloads (attachment) an artifact
* file. The `path` must be a virtual path returned by
* {@link listSessionArtifacts} (e.g. `mnt/user-data/outputs/foo.md`).
*/
export function buildArtifactDownloadUrl(
threadId: string,
virtualPath: string,
options: { download?: boolean } = {},
): string {
const segments = virtualPath.split("/").map((s) => encodeURIComponent(s)).join("/");
const base = `${getBackendBaseURL()}/api/threads/${encodeURIComponent(threadId)}/artifacts/${segments}`;
return options.download ? `${base}?download=true` : base;
}
export interface ReportThreadResponse {
/** Always the built-in singleton id `roundtable-report`. */
agentName: string;
/** Freshly created lightweight thread to run the report agent on. */
threadId: string;
}
/**
* Create a thread for the Step 3「结果绘制」visualization agent
* (`roundtable-report`, a built-in singleton — this only spins up a cheap
* thread, it never re-creates the agent). Run it afterwards with
* {@link streamMultiAgent} (`agentType: "special"`,
* `agentName: "roundtable-report"`, `dAgentThreadId: { [agentName]: threadId }`).
*/
export async function initReportThread(
options: { signal?: AbortSignal } = {},
): Promise<ReportThreadResponse> {
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/report/init`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: "{}",
signal: options.signal,
});
if (!res.ok) {
const text = await res.text().catch(() => "");
reportRoundtableError({
stage: "step3_report", event: "init_thread_http_error",
message: `结果绘制线程创建失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status, agentId: "roundtable-report", agentName: "结果绘制",
});
throw new Error(`init report thread failed (${res.status}): ${text}`);
}
const json = (await res.json()) as { agent_name: string; thread_id: string };
return { agentName: json.agent_name, threadId: json.thread_id };
}
/**
* Create a thread for the Step 3「大屏多页」dashboard-data agent
* (`roundtable-dashboard`, a built-in singleton — distinct from `roundtable-report`;
* its run policy disables ALL tools so it only emits a ```report-json block). Run it
* with {@link streamMultiAgent} (`agentType: "special"`, `agentName: "roundtable-dashboard"`,
* `dAgentThreadId: { [agentName]: threadId }`); the frontend then `extractReportData`s
* the streamed text. Follow-up Q&A / refinements reuse the SAME thread (do not re-init)
* so the agent keeps the prior dashboard context.
*/
export async function initDashboardThread(
options: { signal?: AbortSignal } = {},
): Promise<ReportThreadResponse> {
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/dashboard/init`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: "{}",
signal: options.signal,
});
if (!res.ok) {
const text = await res.text().catch(() => "");
reportRoundtableError({
stage: "step3_dashboard", event: "init_thread_http_error",
message: `大屏页面线程创建失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status, agentId: "roundtable-dashboard", agentName: "大屏页面",
});
throw new Error(`init dashboard thread failed (${res.status}): ${text}`);
}
const json = (await res.json()) as { agent_name: string; thread_id: string };
return { agentName: json.agent_name, threadId: json.thread_id };
}
/**
* Create a thread for the Step 3「总结报告」summary agent (`roundtable-summary`,
* a built-in singleton — distinct from `roundtable-report`; this only spins up a
* cheap thread, it never re-creates the agent). Run it afterwards with
* {@link streamMultiAgent} (`agentType: "special"`, `agentName: "roundtable-summary"`,
* `dAgentThreadId: { [agentName]: threadId }`). Follow-up Q&A / refinements reuse
* the SAME thread (do not re-init) so the agent keeps the report context.
*/
export async function initSummaryThread(
options: { signal?: AbortSignal } = {},
): Promise<ReportThreadResponse> {
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/summary/init`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: "{}",
signal: options.signal,
});
if (!res.ok) {
const text = await res.text().catch(() => "");
reportRoundtableError({
stage: "step3_summary", event: "init_thread_http_error",
message: `方案报告线程创建失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status, agentId: "roundtable-summary", agentName: "方案报告",
});
throw new Error(`init summary thread failed (${res.status}): ${text}`);
}
const json = (await res.json()) as { agent_name: string; thread_id: string };
return { agentName: json.agent_name, threadId: json.thread_id };
}
/**
* Create a per-session thread for the position-roundtable action-planning
* built-in agent. The agent id is fixed; callers persist/reuse the returned
* thread id for subsequent plan refinements.
*/
export async function initPositionActionPlanThread(
options: { signal?: AbortSignal } = {},
): Promise<ReportThreadResponse> {
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/position-action-plan/init`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: "{}",
signal: options.signal,
});
if (!res.ok) {
const text = await res.text().catch(() => "");
reportRoundtableError({
stage: "position_action_plan",
event: "init_thread_http_error",
message: `行动规划线程创建失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status,
agentId: "position-action-planner",
agentName: "行动规划智能体",
});
throw new Error(`init position action-plan thread failed (${res.status}): ${text}`);
}
const json = (await res.json()) as { agent_name: string; thread_id: string };
return { agentName: json.agent_name, threadId: json.thread_id };
}
/**
* Create a thread for the Step 3「入库」structure agent (`roundtable-structure`,
* a built-in singleton — the same agent that extracts the flowchart, but this run
* uses the `ingest` run policy so its attached ingest skill actually executes).
* Run it afterwards with {@link streamMultiAgent} (`agentType: "special"`,
* `agentName: "roundtable-structure"`, `dAgentThreadId: { [agentName]: threadId }`):
* the agent picks the suitable ingest skill (its name may contain「六步法」), reads
* its SKILL.md and calls it — passing the taskId + summary report + structured flow.
*/
export async function initStructureThread(
options: { signal?: AbortSignal } = {},
): Promise<ReportThreadResponse> {
const res = await apiFetch(`${getBackendBaseURL()}/api/multi-agent/structure/init`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: "{}",
signal: options.signal,
});
if (!res.ok) {
const text = await res.text().catch(() => "");
reportRoundtableError({
stage: "step3_structure", event: "init_thread_http_error",
message: `结构化输出线程创建失败 (HTTP ${res.status}): ${text.slice(0, 500)}`,
httpStatus: res.status, agentId: "roundtable-structure", agentName: "结构化输出",
});
throw new Error(`init structure thread failed (${res.status}): ${text}`);
}
const json = (await res.json()) as { agent_name: string; thread_id: string };
return { agentName: json.agent_name, threadId: json.thread_id };
}