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; 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 { // 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; 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, 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 | null { if (!s) return null; try { const v = JSON.parse(s); return v && typeof v === "object" ? (v as Record) : 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) : 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; 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; 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).reasoning_content ?? (additionalKwargs as Record).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 { // 整段包一层兜底:除已在具体分支上报的(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) => { 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 // `` renderer remains the single display path for Steps 2 and 3. const inlineAiDelta = createInlineReasoningDeltaFormatter(); const setPhase = ( next: StreamPhase, name?: string, args?: Record, 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; }; 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; 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) : { 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).status === "model_switch" ) { const obj = parsed as Record; 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; 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).args && typeof (toolCall as Record).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) : Array.isArray(toolCallChunks) && toolCallChunks.length > 0 ? (toolCallChunks[toolCallChunks.length - 1] as Record) : null; const toolName = activeCall && typeof activeCall.name === "string" ? activeCall.name : ""; const args = activeCall && typeof activeCall.args === "object" && activeCall.args !== null ? (activeCall.args as Record) : 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).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).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 { 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 { 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 { 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 { 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 { 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 { 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 }; }