Files
metona-ai-desktop/electron/ipc/agent.ts
T
thzxx 5b9d4d19b3
CI / 类型检查 + Lint + 单元测试 (push) Failing after 9m16s
CI / 全量测试 (Electron ABI) (push) Failing after 6m4s
CI / 产物编译验证 (push) Successful in 11m1s
feat: v0.8.0 流语义补全 · 会话可靠 · 恢复力 — finish_reason 全链路贯通根治"思考中停止" · 2445 用例全量回归
P0 会话可靠性收口(根治"模型思考着会话就停止"):
- P0-1 finish_reason 全链路贯通:DONE 事件与 IterationStep 新增 finishReason,OpenAI 共享 SSE / Anthropic message_delta.stop_reason / Ollama done_reason 三路采集,TRACE 层弃用硬编码 'stop' 记录真值
- P0-2 空响应守卫 + 降级重试:零产出流→可重试错误走退避;思考耗尽输出预算(reasoning-only + length)→自动关闭思考降级重试一次;仍失败→OUTPUT_LENGTH_EXCEEDED 结构化错误 + 故障转移;附带根治 abort 恰逢零工具调用轮被 COMPLETED 抢占的真实缺陷
- P0-3 思考×能力×预算三对齐:DeepSeek/MiMo/Agnes/Ollama 四家 supportsThinking=false 强制不发思考参数;小输出预算告警;设置页联动提示
- P0-4 渲染层可见性:截断/空完成/友好错误三类提示,i18n 全部出层
- P0-5 回归四件套:reasoning-only 终止判定、集成级空闲超时、504 引擎重试归类、思考中 abort→USER_INTERRUPT、P4-2 强制收尾路径

FEAT-1:LLM 设置新增「最大输出上限」——Provider 支持矩阵显隐 + 模型上限钳制提示 + 超限保存警告 + llm.maxTokens 热生效

P1 修复面收口:
- 渲染层三缺陷根治:后台会话回放缓冲(2000 条/4MB 有界 + agent:getReplayState + 事件总线)+ abort 双层自愈 + sendMessage 收尾兜底 + 中断卡片清扫
- 工具 abort 信号全覆盖:web_search/web_fetch/http_request/code_search/git 系列/delegate_task 全部接入引擎中断;web_search 时间预算收敛(720s→≤240s);移除伪造 ToolExecutionContext 与死代码
- 安全:本地 Pinned CONNECT 代理根治浏览器通道 DNS rebinding(校验期 IP pinning,可注入 resolver 表测);配置 URL 域名解析深校验(DeepCheckSoftFailure 软失败);SSE 空 error 帧防御修复;Ollama generate/embed AbortSignal.any 合并
- 缺陷清单:UTF-16 BOM 读取、tmp 同毫秒碰撞(nanoid 后缀)、code_search JS 回退参数对称(case_sensitive/前后文独立)、list_directory include_node_modules、崩溃自愈退避(60s 窗 ≥3 次停 reload)、MemoryViewer/Sidebar i18n 收口

P2 能力演进:
- 会话回收站:SCHEMA_VERSION 3 + 迁移 10(deleted_at,存在性守卫),软删除/恢复/彻底删除/30 天自动清理(启动+24h),searchMessages 聚合剔除,Sidebar 回收站面板
- 会话回放播放器:sessions:listRecordings/readRecording(白名单+目录边界+20MB 上限),SessionReplayPlayer 时间轴/步进/变速,Trace 面板入口
- electron-updater 自动更新:双轨(手动 feed 比对保留),生产环境启动静默检查 + update:status 广播 + app:updateInstall + LogsSettings UpdatePanel + builder publish 配置
- @ 文件提及:workspace.listFiles/readFileClip(边界/512KB/NUL 拒绝/MEMORY.md 保护),ChatInput Fuse 联想+键盘导航+附件管线注入
- MCP Resources/Prompts 发现:可选能力 try/catch 降级,mcp:listServerContents,MCPSettings 展开视图
- 文档对齐:内部 API 标准 HTML(Adapter 清单补 MiMo/已实现注记/STREAM_RESET/DONE.finishReason/ repetition_truncation 映射);README v0.8.0 亮点表

P3 测试基建:
- 新增 4 个测试文件:engine-stream-contract(6)、engine-stream-reliability(4:集成空闲超时/504 重试/思考中 abort/P4-2 强制收尾)、thinking-capability-gate(7)、pinned-proxy(9,含深校验 5)、session-trash(5,DB 域)、use-agent-stream hook 级(5)、agent.test 回放缓冲(2)
- 契约更新:orchestrator 被中断 SubAgent success=false(abort 优先级修复语义)、SSE 空 error 帧、UTF-16 正常读取、DeepSeek 未配置思考显式 disabled、迁移矩阵 v2→3
- 弱断言根治:registry WEBP 单向断言、hooks-contracts 自比恒真、memory 空 token 补强

全量验证:typecheck 0 错误 / lint 0 问题 / 系统 Node 2144 通过(301 DB 用例按 ABI 跳过)/ Electron ABI 2445/2445 全量通过 0 跳过
2026-09-05 20:06:26 +08:00

1353 lines
54 KiB
TypeScript
Raw 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.
/**
* IPC Agent Handlers — Agent 交互域(P2-9 从 handlers.ts 拆分)
*
* 职责:
* 1. agent:sendMessage — 消息发送编排(历史加载/记忆注入/注入检测/runStream/持久化)
* 2. agent:abortSession — 中断会话(联动 SubAgent
* 3. 常驻引擎事件管道 — 流式转发(按会话节流)、状态广播、TRACE 录制(P1-6 补全 9 种事件)、
* 压缩/死循环/Provider 切换通知
*
* P2-10: 事件监听从"每消息 attach/detach"改为常驻管道(按 sessionId 隔离状态),
* 支持多会话并发流式与多窗口广播。
*/
import { ipcMain } from 'electron';
import type { IPCContext } from './context';
import { broadcast } from './context';
import type {
MetonaMessage,
MetonaStreamEvent,
MetonaError,
MetonaSystemPrompt,
} from '../harness/types';
import { MetonaErrorCode, MetonaStreamEventType } from '../harness/types';
import { estimateMessagesTokens } from '../harness/utils/token-estimator';
import { DeepSeekAdapter } from '../harness/adapters/deepseek.adapter';
import { OllamaAdapter } from '../harness/adapters/ollama.adapter';
// v0.7.3 P1-1: 用户上下文前置块(动态内容出 system,保 prompt cache 前缀稳定)
import { buildUserContextPrefix, withUserContextPrefix } from '../harness/prompts/user-context';
// v0.7.3 P1-5: 记忆固化触发决策(纯函数)
import { shouldConsolidate } from '../harness/memory/consolidation-policy';
import log from 'electron-log';
/** 构建 "时区名 (UTC±N)" 标签(注入用户上下文前置块;失败回退 UTC) */
function buildTimezoneLabel(): string {
try {
const tz = Intl.DateTimeFormat().resolvedOptions().timeZone ?? 'UTC';
const offset = -new Date().getTimezoneOffset() / 60;
const offsetStr = offset >= 0 ? `UTC+${offset}` : `UTC${offset}`;
return `${tz} (${offsetStr})`;
} catch {
return 'UTC';
}
}
/** 单会话的 text_delta 节流状态 */
interface ThrottleState {
buffer: string;
lastEventMeta: Pick<
MetonaStreamEvent,
'requestId' | 'sessionId' | 'iteration' | 'seq' | 'timestamp' | 'runId'
> | null;
flushTimer: ReturnType<typeof setTimeout> | null;
}
/** 单会话的迭代录制状态(TRACE 层) */
interface IterationTrace {
iteration: number;
startedAt: number;
text: string;
usage?: { input: number; output: number; total: number };
responded: boolean; // 本轮 llm_response 是否已记录(PARSING 与下一轮 THINKING 去重)
/**
* v0.8.0 P0-1: 本轮 Provider 原生停止原因(DONE 事件携带的 finishReason)。
* v0.7.4 及之前 recordLLMResponse 硬编码 'stop'length 截断在 TRACE 文件中
* 不可归因;现记录真值,缺省回退 'stop'。
*/
finishReason?: string;
}
/**
* v0.8.0 P1-1a: 单会话的流式回放缓冲 —— 渲染层切走会话期间,事件管道仍照常
* 投递(按会话隔离被渲染层准入拒绝),此缓冲按序保存 streamEvent + stateChange
* 双通道事件;用户切回会话时经 agent:getReplayState 拉取并灌入渲染层事件总线,
* 实现"后台会话运行内容不丢失"。
*
* 有界设计:单会话上限 MAX_REPLAY_EVENTS 条 / MAX_REPLAY_BYTES 字节,溢出丢弃
* 最旧并置 truncated 标记;TERMINATED 后保留(供"完成后切回"看到最终内容),
* 下一次 run 启动(INIT stateChange)时清空重建。
*/
interface ReplayEntry {
channel: 'streamEvent' | 'stateChange';
payload: unknown;
ts: number;
}
interface ReplayBuffer {
events: ReplayEntry[];
bytes: number;
truncated: boolean;
/** 最近一次观察到的 runId(供渲染层回放后对齐 run 守卫) */
runId?: string;
terminated: boolean;
}
const MAX_REPLAY_EVENTS = 2000;
const MAX_REPLAY_BYTES = 4 * 1024 * 1024;
const MAX_REPLAY_SESSIONS = 50;
export function registerAgentHandlers(ctx: IPCContext): void {
const {
agentEngineManager,
sessionRecorder,
configService,
sessionService,
workspaceService,
contextBuilder,
auditService,
memoryManager,
promptInjectionDefender,
outputValidator,
memoryConsolidator,
sessionSummaryService,
orchestrator,
confirmationHook,
titleGenerator,
} = ctx;
// ===== 常驻事件管道:text_delta 按会话节流(F8 =====
const throttleStates = new Map<string, ThrottleState>();
const iterationTraces = new Map<string, IterationTrace>();
// v0.8.0 P1-1a: 每会话流式回放缓冲(后台会话内容恢复)
const replayBuffers = new Map<string, ReplayBuffer>();
/** 追加一条事件到会话回放缓冲(有界:条数/字节双上限,溢出丢最旧) */
const appendReplay = (sessionId: string, entry: ReplayEntry): void => {
let buf = replayBuffers.get(sessionId);
if (!buf) {
// 会话数上限保护:超出后淘汰最早的缓冲(Map 迭代序 = 插入序)
if (replayBuffers.size >= MAX_REPLAY_SESSIONS) {
const oldest = replayBuffers.keys().next().value as string | undefined;
if (oldest !== undefined) replayBuffers.delete(oldest);
}
buf = { events: [], bytes: 0, truncated: false, terminated: false };
replayBuffers.set(sessionId, buf);
}
// run 启动(INIT)后旧 run 缓冲作废 —— 由 stateChange 监听器在 INIT 时清空,
// 此处仅在 terminated 缓冲上遇到新内容时兜底重置(正常路径 INIT 先行)
if (buf.terminated) {
replayBuffers.delete(sessionId);
buf = { events: [], bytes: 0, truncated: false, terminated: false };
replayBuffers.set(sessionId, buf);
}
let size = 0;
try {
size = JSON.stringify(entry.payload).length;
} catch {
size = 256; // 序列化失败按保守值计入
}
buf.events.push(entry);
buf.bytes += size;
while (
buf.events.length > MAX_REPLAY_EVENTS ||
(buf.bytes > MAX_REPLAY_BYTES && buf.events.length > 1)
) {
const dropped = buf.events.shift();
buf.truncated = true;
if (dropped) {
try {
buf.bytes -= JSON.stringify(dropped.payload).length;
} catch {
buf.bytes -= 256;
}
}
}
// 记录 runId(首个携带者)
const payload = entry.payload as { runId?: string } | null;
if (!buf.runId && payload?.runId) buf.runId = payload.runId;
};
// v0.7.3 P1-5: 会话级记忆固化时间戳(consolidation-policy 频率门控的状态持有方)
// v0.7.4 P3-5: LRU 化 —— 旧实现每会话一条永不清除,长期运行后无界增长。
// 容量 200:超过后淘汰最旧的(固化频率窗口 10 分钟,200 个会话远超实际并发上限)。
const lastConsolidationBySession = new Map<string, number>();
const MAX_CONSOLIDATION_TRACKED = 200;
const touchConsolidation = (sessionId: string): number => {
const now = Date.now();
lastConsolidationBySession.set(sessionId, now);
if (lastConsolidationBySession.size > MAX_CONSOLIDATION_TRACKED) {
// 淘汰最旧(Map 迭代顺序 = 插入顺序)
const oldestKey = lastConsolidationBySession.keys().next().value as string | undefined;
if (oldestKey !== undefined) lastConsolidationBySession.delete(oldestKey);
}
return now;
};
const flushThrottle = (sessionId: string): void => {
const st = throttleStates.get(sessionId);
if (!st) return;
st.flushTimer = null;
if (!st.buffer || !st.lastEventMeta) {
st.buffer = '';
st.lastEventMeta = null;
return;
}
// 构建合并的 text_delta 事件,保留最后一个 delta 的元数据
const mergedEvent: MetonaStreamEvent = {
...st.lastEventMeta,
type: MetonaStreamEventType.TEXT_DELTA,
delta: st.buffer,
seq: st.lastEventMeta.seq,
timestamp: Date.now(),
};
broadcast('agent:streamEvent', mergedEvent);
st.buffer = '';
st.lastEventMeta = null;
};
const cleanupSessionState = (sessionId: string): void => {
const st = throttleStates.get(sessionId);
if (st?.flushTimer) {
clearTimeout(st.flushTimer);
flushThrottle(sessionId);
}
throttleStates.delete(sessionId);
iterationTraces.delete(sessionId);
};
// ===== 常驻监听:流式事件(节流转发 + TRACE 录制) =====
agentEngineManager.on('streamEvent', (event: MetonaStreamEvent) => {
if (!event.sessionId) return;
const sessionId = event.sessionId;
// v0.8.0 P1-1a: 全量事件进入回放缓冲(含 text_delta —— 必须在节流分支
// return 之前入队,否则切回会话时已缓冲的正文增量缺失)
appendReplay(sessionId, { channel: 'streamEvent', payload: event, ts: Date.now() });
const trace = iterationTraces.get(sessionId);
// F8: text_delta 聚合,其他事件立即转发(先 flush 保证顺序)
if (event.type === MetonaStreamEventType.TEXT_DELTA && event.delta) {
const st = throttleStates.get(sessionId) ?? {
buffer: '',
lastEventMeta: null,
flushTimer: null,
};
throttleStates.set(sessionId, st);
if (st.buffer === '') {
st.lastEventMeta = {
requestId: event.requestId,
sessionId: event.sessionId,
iteration: event.iteration,
seq: event.seq,
timestamp: event.timestamp,
runId: event.runId,
};
} else if (st.lastEventMeta) {
st.lastEventMeta = {
...st.lastEventMeta,
seq: event.seq,
timestamp: event.timestamp,
runId: event.runId,
};
}
st.buffer += event.delta;
if (st.flushTimer === null) {
st.flushTimer = setTimeout(() => flushThrottle(sessionId), 32);
}
// TRACE: 累积本轮 LLM 文本(供 llm_response 记录)
if (trace) trace.text += event.delta;
return;
}
// 非 text_delta 事件:先 flush 缓冲区,再立即转发(保证事件顺序)
const st = throttleStates.get(sessionId);
if (st?.flushTimer) {
clearTimeout(st.flushTimer);
flushThrottle(sessionId);
}
// TRACE: 工具调用/结果/usage 录制(P1-6 补全)
switch (event.type) {
case MetonaStreamEventType.TOOL_CALL_COMPLETE:
if (event.toolCall) {
sessionRecorder.recordToolCall({
sessionId,
iteration: event.iteration,
toolName: event.toolCall.name,
args: event.toolCall.args,
});
}
break;
case MetonaStreamEventType.TOOL_RESULT:
if (event.toolResult) {
const resultPreview =
typeof event.toolResult.result === 'string'
? event.toolResult.result
: JSON.stringify(event.toolResult.result);
sessionRecorder.recordToolResult({
sessionId,
iteration: event.iteration,
toolName: event.toolResult.toolName,
success: event.toolResult.success,
durationMs: event.toolResult.durationMs ?? 0,
resultPreview,
error: event.toolResult.error,
});
}
break;
case MetonaStreamEventType.USAGE:
if (trace && event.usage) {
trace.usage = {
input: event.usage.inputTokens ?? 0,
output: event.usage.outputTokens ?? 0,
total: event.usage.totalTokens ?? 0,
};
}
break;
case MetonaStreamEventType.DONE:
// v0.7.4 P1-6: DONE 不再执行 cleanupSessionState —— 引擎的 finish() 顺序是
// 先发 DONE 再发 TERMINATED stateChange,而 TERMINATED 分支(stateChange 监听)
// 需要读取 iterationTraces 补记最终迭代的 iteration_end/llm_response。
// 旧实现在 DONE 时删除 iterationTraces → 最终迭代的 TRACE 完整性缺失。
// 现仅在 TERMINATED 分支统一收尾(flush throttle buffer + 删 iterationTraces)。
// 注意:此处仍需转发 DONE 事件到前端(下方 broadcast 统一执行)。
// v0.8.0 P0-1: 捕获 Provider 原生停止原因(recordLLMResponse 记录真值)
if (trace && event.finishReason) trace.finishReason = event.finishReason;
break;
default:
break;
}
broadcast('agent:streamEvent', event);
});
// ===== 常驻监听:状态变化(广播 + TRACE 迭代录制) =====
agentEngineManager.on(
'stateChange',
(data: {
previous?: string;
current?: string;
state?: string;
sessionId?: string;
iteration?: number;
runId?: string;
}) => {
if (data.previous) log.info(`[AGENT] State: ${data.previous}${data.current}`);
broadcast('agent:stateChange', data);
const sessionId = data.sessionId;
// v0.8.0 P1-1a: stateChange 事件进入回放缓冲。INIT(run 启动)清空上一 run
// 的缓冲;TERMINATED 置终态标记(保留缓冲供"完成后切回"回放最终内容)。
if (sessionId) {
const stateRaw = data.state ?? data.current ?? '';
if (stateRaw === 'INIT') {
replayBuffers.delete(sessionId);
} else {
appendReplay(sessionId, { channel: 'stateChange', payload: data, ts: Date.now() });
if (stateRaw === 'TERMINATED') {
const buf = replayBuffers.get(sessionId);
if (buf) buf.terminated = true;
}
}
}
if (!sessionId || data.iteration == null) return;
const stateValue = data.state ?? data.current ?? '';
const trace = iterationTraces.get(sessionId);
// THINKING 且迭代号变化 → 新迭代开始(关闭上一迭代)
if (stateValue === 'THINKING' && (!trace || trace.iteration !== data.iteration)) {
if (trace && !trace.responded) {
sessionRecorder.recordLLMResponse({
sessionId,
iteration: trace.iteration,
content: trace.text,
// v0.8.0 P0-1: 记录真实停止原因(旧实现硬编码 'stop'length 截断不可归因)
finishReason: trace.finishReason ?? 'stop',
tokenUsage: trace.usage ?? { input: 0, output: 0, total: 0 },
});
}
if (trace) {
sessionRecorder.recordIterationEnd(sessionId, {
iteration: trace.iteration,
durationMs: Date.now() - trace.startedAt,
});
}
iterationTraces.set(sessionId, {
iteration: data.iteration,
startedAt: Date.now(),
text: '',
responded: false,
});
sessionRecorder.recordIterationStart(sessionId, data.iteration);
const provider = configService.get<string>('llm.provider') ?? '';
const model = configService.get<string>('llm.model') ?? '';
sessionRecorder.recordLLMRequest({
sessionId,
iteration: data.iteration,
provider,
model,
messageCount: data.iteration + 1,
});
}
// PARSING → 本轮流式结束,记录 llm_response
if (
stateValue === 'PARSING' &&
trace &&
trace.iteration === data.iteration &&
!trace.responded
) {
trace.responded = true;
sessionRecorder.recordLLMResponse({
sessionId,
iteration: trace.iteration,
content: trace.text,
// v0.8.0 P0-1: 记录真实停止原因(旧实现硬编码 'stop'length 截断不可归因)
finishReason: trace.finishReason ?? 'stop',
tokenUsage: trace.usage ?? { input: 0, output: 0, total: 0 },
});
}
// TERMINATED → 补记最终迭代的 iteration_end(正常流程只在下一轮 THINKING 补记,
// 最终轮无后续迭代,需在此补齐 TRACE 完整性)+ 兜底清理会话管道状态
if (stateValue === 'TERMINATED') {
if (trace) {
sessionRecorder.recordIterationEnd(sessionId, {
iteration: trace.iteration,
durationMs: Date.now() - trace.startedAt,
});
}
cleanupSessionState(sessionId);
}
},
);
// ===== 常驻监听:上下文压缩(toast + streamEvent 通知) =====
agentEngineManager.on(
'compressed',
(data: {
sessionId?: string;
iteration?: number;
originalTokens?: number;
compressedTokens?: number;
}) => {
const savedTokens = Math.max(0, (data.originalTokens ?? 0) - (data.compressedTokens ?? 0));
// toast 通知用户压缩已发生
broadcast('toast:show', {
type: 'info',
message: `上下文压缩: ${data.originalTokens ?? '?'}${data.compressedTokens ?? '?'} tokens(节省 ${savedTokens}`,
});
// 通过 streamEvent 转发,前端 useAgentStream 监听 'compressed' 类型后更新 store
broadcast('agent:streamEvent', {
type: 'compressed',
sessionId: data.sessionId ?? '',
iteration: data.iteration ?? 0,
originalTokens: data.originalTokens,
compressedTokens: data.compressedTokens,
savedTokens,
timestamp: Date.now(),
});
},
);
// ===== 常驻监听:死循环检测(toast 警告) =====
agentEngineManager.on('deadLoop', (data: { iteration?: number; sessionId?: string }) => {
log.warn(`[AGENT] Dead loop detected at iteration ${data.iteration ?? '?'}`);
broadcast('toast:show', {
type: 'warning',
message: `检测到死循环(第 ${data.iteration ?? '?'} 轮):连续3轮重复相同工具调用,已自动终止`,
});
});
// ===== 常驻监听:Provider 故障转移(P1,通知前端 + toast =====
agentEngineManager.on(
'providerSwitched',
(data: { from?: string; to?: string; reason?: string; sessionId?: string }) => {
broadcast('agent:providerSwitched', {
from: data.from,
to: data.to,
reason: data.reason ?? 'failover',
sessionId: data.sessionId ?? '',
});
broadcast('toast:show', {
type: 'warning',
message: `Provider 故障转移: ${data.from ?? '?'}${data.to ?? '?'}(主 Provider 请求失败)`,
});
},
);
// ===== Agent 消息发送 =====
ipcMain.handle(
'agent:sendMessage',
async (_event, userMessage: MetonaMessage, sessionId: string) => {
// M-33 修复: 参数校验,防止 undefined/非字符串导致下游异常
// P1-5 修复: 校验失败时也发 ERROR+DONE 流事件,防止 isStreaming 永久卡死
const sendErrorEvent = (message: string, sid: string): void => {
const errorEvent: MetonaStreamEvent = {
type: MetonaStreamEventType.ERROR,
requestId: '',
sessionId: sid,
iteration: 0,
seq: 0,
timestamp: Date.now(),
error: { code: MetonaErrorCode.UNKNOWN, message, retryable: false },
};
broadcast('agent:streamEvent', errorEvent);
broadcast('agent:streamEvent', { ...errorEvent, type: MetonaStreamEventType.DONE });
};
if (!sessionId || typeof sessionId !== 'string') {
log.warn('[AGENT] sendMessage rejected: invalid sessionId');
sendErrorEvent('无效的会话 ID', sessionId ?? '');
return { success: false, error: 'Invalid sessionId' };
}
if (
!userMessage ||
typeof userMessage !== 'object' ||
typeof userMessage.content !== 'string'
) {
log.warn('[AGENT] sendMessage rejected: invalid userMessage');
sendErrorEvent('无效的消息格式', sessionId);
return { success: false, error: 'Invalid message format' };
}
log.info('[AGENT] sendMessage:', sessionId, (userMessage.content ?? '').slice(0, 80));
// v0.7.4 P1-3: 同会话并发防重入 —— 同一 sessionId 的第二个 invoke 在第一个 run
// 未结束时到达,会让同一引擎被并行 runStream(工具副作用并发执行 / 内部 run-lock
// 排队 30s 后强制 abort 旧 run,用户看到"上一次操作未完成")。此处直接拒绝并
// 广播明确错误事件,前端 isStreaming 正常收尾。
if (agentEngineManager.isRunning(sessionId)) {
const busyMsg = '该会话正在执行任务,请等待完成或先中断后再发送';
log.warn(`[AGENT] sendMessage rejected: session ${sessionId} is already running`);
sendErrorEvent(busyMsg, sessionId);
return { success: false, error: busyMsg };
}
// v0.7.4 P1-4: 会话存在性预检 —— sessionId 指向已删除会话时,后续 saveMessage
// 会因 FK 约束(foreign_keys=ON)抛错,而旧实现该调用在 try 块之外,invoke 直接
// reject:不发 ERROR/DONE 事件、不 stopRecording,前端 isStreaming 永久卡死。
// 此处提前校验并走统一的 sendErrorEvent + stopRecording 收尾路径。
const sessionExists = sessionService.getSession(sessionId) != null;
if (!sessionExists) {
const missingMsg = '会话不存在或已被删除,请刷新后重试';
log.warn(`[AGENT] sendMessage rejected: session ${sessionId} not found`);
sendErrorEvent(missingMsg, sessionId);
await sessionRecorder.stopRecording(sessionId, {
totalIterations: 0,
totalTokens: 0,
durationMs: 0,
terminationReason: 'error',
});
return { success: false, error: missingMsg };
}
// 发送消息前确保 Adapter 使用最新配置(失败则中止,防止用旧 Provider 的 adapter 发送)
if (!ctx.reloadAdapter()) {
const errorMsg =
'Adapter 加载失败,请检查 LLM 配置(Provider、API Key、Base URL、Model 是否完整)';
log.error('[AGENT]', errorMsg);
sendErrorEvent(errorMsg, sessionId);
await sessionRecorder.stopRecording(sessionId, {
totalIterations: 0,
totalTokens: 0,
durationMs: 0,
terminationReason: 'error',
});
return { success: false, error: errorMsg };
}
// TRACE 层:开始录制 / TOOL 层:记录会话开始
sessionRecorder.startRecording(sessionId);
auditService.logSessionStart(sessionId);
// v0.7.4 P1-4: 以下数据准备(保存用户消息/加载历史/构建 System Prompt/记忆检索)
// 全部移入 try 块 —— 旧实现这些调用在 try 之外,DB 或文件系统抛错时 invoke 直接
// reject(不发 ERROR/DONE 事件、不 stopRecording,前端 isStreaming 永久卡死)。
// 统一由下方 catch 收尾:ERROR + DONE 双事件 + 审计 + 录制终止。
let history: MetonaMessage[] = [];
let systemPrompt: MetonaSystemPrompt | null = null;
let userContextPrefix = '';
let engineUserMessage: MetonaMessage;
try {
// 保存用户消息到数据库
// v0.7.4 回归修复: 透传前端消息 idChatMessage.id 由 genMsgId 生成)——
// 否则 DB 用 msg_<nanoid> 生成不同 id,用户对刚发送消息"仅保存"
// updateMessageContent 按 id 匹配)会 0 行更新失败。
sessionService.saveMessage({
sessionId,
role: 'user',
content: userMessage.content,
attachments: (userMessage as MetonaMessage & { attachments?: unknown[] }).attachments,
id: (userMessage as MetonaMessage & { id?: string }).id,
});
// P2-11: 分层加载历史——存在滚动摘要时只加载 [摘要 + 近期原文]
history = sessionSummaryService.buildHistoryMessages(sessionId).slice(0, -1);
// 从工作空间文件构建 System Prompt
const workspaceFiles = workspaceService.getFiles();
systemPrompt = contextBuilder.buildSystemPrompt(workspaceFiles, workspaceService.getPath());
// v0.3.18 修复: SOUL.md 为空或不存在时降级到默认身份,向前端发 toast 提示用户
if (contextBuilder.isUsingFallbackRole()) {
broadcast('toast:show', {
type: 'info',
message:
'未找到 SOUL.md 或内容为空,已使用默认 Metona 身份。可在工作空间根目录创建 SOUL.md 自定义 Agent 人格',
});
}
// 检索与用户消息相关的记忆 + 附件元信息 → 构建用户消息上下文前置块。
// v0.7.3 P1-1 根治: 记忆注入与附件提示此前追加进 systemPrompt.dynamicReminders
// 每条消息都改变 system 字节 → 跨 run 缓存全 miss。现随首条 user 消息注入
// LLM 语义等价),system prompt 保持跨 run 字节级稳定。
try {
const memories = memoryManager.search(userMessage.content, {
topK: 5,
minImportance: 0.3,
});
const attachments = (
userMessage as MetonaMessage & {
attachments?: Array<{ name: string; type: string; truncated?: boolean }>;
}
).attachments;
userContextPrefix = buildUserContextPrefix({
now: Date.now(),
memories,
attachments: Array.isArray(attachments) ? attachments : [],
timezoneLabel: buildTimezoneLabel(),
});
if (memories.length > 0) {
log.debug(`[AGENT] Injected ${memories.length} memories into user context prefix`);
}
} catch (err) {
log.warn('[AGENT] Memory retrieval failed, proceeding without memories:', err);
// 记忆检索失败时前置块退化为仅含日期时间(附件提示随之丢失可接受——
// 主进程附件提示是辅助语义,附件内容本体仍在消息中)
userContextPrefix = buildUserContextPrefix({
now: Date.now(),
memories: [],
attachments: [],
timezoneLabel: buildTimezoneLabel(),
});
}
// 构建发送给引擎的用户消息副本 —— 前置块只存在于该副本:
// DB 持久化(上方 saveMessage 已用原始内容)、前端展示、记忆固化、
// 注入检测均使用原始干净内容,互不污染。
engineUserMessage = {
...userMessage,
content: withUserContextPrefix(userContextPrefix, userMessage.content),
};
} catch (err) {
log.error('[AGENT] Failed to prepare message context:', err);
const prepErr = (err as Error).message;
broadcast('agent:streamEvent', {
type: MetonaStreamEventType.ERROR,
requestId: '',
sessionId,
iteration: 0,
seq: 0,
timestamp: Date.now(),
error: { code: MetonaErrorCode.UNKNOWN, message: prepErr, retryable: false },
} satisfies MetonaStreamEvent);
await sessionRecorder.stopRecording(sessionId, {
totalIterations: 0,
totalTokens: 0,
durationMs: 0,
terminationReason: 'error',
});
return { success: false, error: prepErr };
}
// 提示注入检测在 try 内执行(需 systemPrompt 已构建,与引擎运行同域)
try {
// 提示注入检测(安全模块)
// F-8 接通: security.promptInjectionDefense=false 时跳过用户消息检测
// (工具结果侧的 SecurityScanHook 由 main.ts 按同一配置决定是否挂载)
// fail-secure: 仅显式 false 才关闭 —— 配置值异常(空串/null/类型错误)时保持防护开启
const injectionEnabled =
configService.get<boolean>('security.promptInjectionDefense') !== false;
if (injectionEnabled) {
const injectionResult = promptInjectionDefender.detect(userMessage.content);
if (injectionResult.riskScore >= 7) {
log.warn('[PromptInjectionDefender] Blocked message:', injectionResult.findings);
sendErrorEvent(
`Message blocked by prompt injection defense: ${injectionResult.recommendation}`,
sessionId,
);
await sessionRecorder.stopRecording(sessionId, {
totalIterations: 0,
totalTokens: 0,
durationMs: 0,
terminationReason: 'error',
});
return { success: false, error: 'Message blocked by prompt injection defense' };
}
if (injectionResult.riskScore >= 4) {
log.warn(
'[PromptInjectionDefender] Suspicious patterns detected:',
injectionResult.findings,
);
}
}
// TRACE 层:记录上下文构建
sessionRecorder.recordContextBuilt(sessionId, {
tokenCount: estimateMessagesTokens(history),
usageRatio: 0,
});
// 启动 Agent LoopP2-10: 每会话独立引擎)
// v0.7.3 P1-1: 传入带上下文前置块的消息副本 —— 原始 userMessage 保持干净
const engine = agentEngineManager.getEngine(sessionId);
const output = await engine.runStream(engineUserMessage, sessionId, history, systemPrompt);
// 输出验证(不阻塞响应,仅记录警告)
// v0.3.0 修复: 传入 toolResults 和 context,启用事实一致性检查和幻觉检测
// v0.4.1: warning 及以上级别的 issue 通过 VALIDATION 流事件推送前端展示(此前仅写日志,用户不可感知)
try {
const toolResults = output.iterations
.flatMap((step) => step.toolResults ?? [])
.map((r) => (typeof r.result === 'string' ? r.result : JSON.stringify(r.result)));
const context = [...history, { role: 'user', content: userMessage.content }]
.map((m) => `${m.role}: ${m.content}`)
.join('\n');
const validation = await outputValidator.validate(output.finalAnswer, {
toolResults: toolResults.length > 0 ? toolResults : undefined,
context,
});
if (!validation.valid || validation.issues.length > 0) {
log.warn('[OutputValidator] Validation issues:', validation.issues);
}
log.debug(`[OutputValidator] Score: ${validation.score}, Valid: ${validation.valid}`);
// v0.4.1: 推送验证结果到前端 — 只推送 warning/error 级(info 级为噪声)
// 类型谓词收窄 severity,确保与 MetonaValidationPayload.issues 的类型一致
const visibleIssues = validation.issues
.filter(
(i): i is typeof i & { severity: 'warning' | 'error' } =>
i.severity === 'warning' || i.severity === 'error',
)
.slice(0, 5);
if (visibleIssues.length > 0) {
broadcast('agent:streamEvent', {
type: MetonaStreamEventType.VALIDATION,
requestId: '',
sessionId,
iteration: output.iterations.length,
seq: 0,
timestamp: Date.now(),
validation: {
score: validation.score,
issues: visibleIssues.map((i) => ({
severity: i.severity,
type: i.type,
message: i.message,
})),
},
} satisfies MetonaStreamEvent);
}
} catch (err) {
log.error('[OutputValidator] Validation failed:', err);
}
// 保存每轮迭代的 assistant 消息到数据库(含思考内容和工具调用)
// 崩溃修复(会话停止根因的持久化侧): 原 `if (!step.thought) continue;` 把
// 纯工具调用轮(零文本零思考)整个跳过 — assistant 与 tool 结果都不落库。
// 后果:① 重启后历史缺失工具上下文(模型"忘记"做过什么,重复调用 →
// 表现为工具调用不稳定);② 与 engine 运行时缺陷同源。现改为:
// 有 thought 或有 toolCalls 的步骤都保存。
for (const step of output.iterations) {
if (!step.thought && !(step.toolCalls && step.toolCalls.length > 0)) continue;
const toolCallsWithResults = step.toolCalls?.map((tc) => {
const result = step.toolResults?.find((r) => r.toolCallId === tc.id);
return {
id: tc.id,
name: tc.name,
args: tc.args,
status: result?.success ? ('success' as const) : ('error' as const),
result: result?.result,
durationMs: result?.durationMs,
error: result?.error,
};
});
// 只有当有内容、思考内容或工具调用时才保存
if (
step.thought?.content ||
step.thought?.reasoningContent ||
toolCallsWithResults?.length
) {
// C-6 修复: assistant 消息仅有 tool_calls 时 content 必须为 null(而非空字符串)
const assistantContent =
toolCallsWithResults?.length && !step.thought?.content
? null
: (step.thought?.content ?? null);
sessionService.saveMessage({
sessionId,
role: 'assistant',
content: assistantContent,
reasoningContent: step.thought?.reasoningContent || undefined,
toolCalls: toolCallsWithResults,
iteration: step.iteration,
});
}
// v0.3.0 修复: 保存 tool 结果消息到数据库
// OpenAI 兼容 API 要求 assistant 消息有 tool_calls 时,后续必须有对应的 tool 结果消息
if (step.toolResults) {
for (const result of step.toolResults) {
const resultContent =
typeof result.result === 'string' ? result.result : JSON.stringify(result.result);
sessionService.saveMessage({
sessionId,
role: 'tool',
content: result.error ?? resultContent,
toolResult: result,
iteration: step.iteration,
});
}
}
}
// 更新 Token 统计
if (output.totalTokenUsage.totalTokens > 0) {
sessionService.updateTokenUsage(sessionId, output.totalTokenUsage.totalTokens);
}
// 更新 MEMORY.md 时间戳
workspaceService.updateMemoryTimestamp();
// 会话结束:AI 判断本次对话有哪些重要内容需要持久化到 MEMORY.md
// v0.7.3 P1-5 节流: 此前每次 run 无条件发起固化 LLM 请求,短寒暄同样触发。
// 现按 consolidation-policy 决策(总开关 + 内容门控 + 会话级频率窗口)触发;
// 异步执行不阻塞主流程返回;失败仅记录日志。
const lastConsolidationAt = lastConsolidationBySession.get(sessionId) ?? 0;
const hadSuccessfulToolCall = output.iterations.some((step) =>
(step.toolResults ?? []).some((r) => r.success),
);
const decision = shouldConsolidate({
enabled: configService.get<boolean>('memory.consolidationEnabled') !== false,
answerChars: output.finalAnswer?.length ?? 0,
minChars: configService.get<number>('memory.consolidationMinChars') ?? 200,
hadSuccessfulToolCall,
lastConsolidationAt,
now: Date.now(),
intervalMs: configService.get<number>('memory.consolidationIntervalMs') ?? 600_000,
});
if (decision.consolidate) {
memoryConsolidator
.consolidate(userMessage.content, output.finalAnswer, output.iterations)
.then((result) => {
if (result.appended > 0) {
touchConsolidation(sessionId);
log.info(
`[AGENT] Memory consolidated: ${result.appended} entries appended to MEMORY.md`,
);
broadcast('toast:show', {
type: 'info',
message: `AI 已将 ${result.appended} 条重要记忆写入 MEMORY.md`,
});
}
})
.catch((err) => {
log.warn('[AGENT] Memory consolidation failed:', err);
});
} else {
log.debug(`[AGENT] Memory consolidation skipped (${decision.reason})`);
}
// v0.7.3 P4-1: 首个完成的 run 之后生成精炼会话标题(每会话幂等,失败静默)
if (output.terminationReason === 'completed') {
titleGenerator
.maybeGenerateTitle(sessionId, userMessage.content, output.finalAnswer)
.then((title) => {
if (title) {
// 广播重命名结果,前端 Sidebar 实时刷新标题
broadcast('config:changed', { key: `session.title.${sessionId}`, value: title });
}
})
.catch(() => {
/* 静默 —— 标题失败已有 debug 日志 */
});
}
// TOOL 层:记录会话结束 / TRACE 层:停止录制
auditService.logSessionEnd({
sessionId,
totalIterations: output.iterations.length,
totalTokens: output.totalTokenUsage.totalTokens,
durationMs: output.durationMs,
terminationReason: output.terminationReason,
});
await sessionRecorder.stopRecording(sessionId, {
totalIterations: output.iterations.length,
totalTokens: output.totalTokenUsage.totalTokens,
durationMs: output.durationMs,
terminationReason: output.terminationReason,
});
// P2-11: 会话结束后评估滚动摘要(fire-and-forget,失败仅记录)
sessionSummaryService.maybeSummarize(sessionId).catch((err) => {
log.warn('[AGENT] Session summary generation failed:', err);
});
log.info(
`[AGENT] Completed: ${output.terminationReason}, ${output.iterations.length} iterations, ${output.durationMs}ms`,
);
return { success: true };
} catch (error) {
log.error('[AGENT] Error:', error);
// TOOL 层:记录错误
auditService.log({
sessionId,
eventType: 'error',
actor: 'agent',
target: 'agent_loop',
details: { error: (error as Error).message },
outcome: 'error',
});
// TRACE 层:停止录制
await sessionRecorder.stopRecording(sessionId, {
totalIterations: 0,
totalTokens: 0,
durationMs: 0,
terminationReason: 'error',
});
// 发送错误事件到 UI
const metonaError: MetonaError = {
code: MetonaErrorCode.UNKNOWN,
message: (error as Error).message,
retryable: false,
};
broadcast('agent:streamEvent', {
type: MetonaStreamEventType.ERROR,
requestId: '',
sessionId,
iteration: 0,
seq: 0,
timestamp: Date.now(),
error: metonaError,
} satisfies MetonaStreamEvent);
return { success: false, error: (error as Error).message };
}
},
);
// ===== v0.5.0: DeepSeek 余额查询(复用 DeepSeekAdapter.getBalance,原为死代码) =====
ipcMain.handle('llm:getBalance', async () => {
try {
const adapter = agentEngineManager.getAdapter();
if (!(adapter instanceof DeepSeekAdapter)) {
return {
success: false,
error: 'Balance query is only supported for the DeepSeek provider',
};
}
const balance = await adapter.getBalance();
if (!balance) {
return { success: false, error: '余额查询失败(API Key 无效或网络错误)' };
}
return { success: true, data: balance };
} catch (error) {
log.warn('[AGENT] getBalance failed:', (error as Error).message);
return { success: false, error: (error as Error).message };
}
});
// ===== v0.7.2 P3-9: 动态模型列表(六家 adapter 的 listModels 首次获得 IPC 消费者)=====
// 此前 DeepSeek/OpenAI/Ollama 均实现了实时模型发现(/models、/api/tags),
// 但 preload/IPC 层无任何通道 —— 设置页模型输入只能靠用户手填。
// 契约:列表基于"已保存"的 LLM 配置(与引擎实际使用的 adapter 同源);
// 配置不完整时显式失败,绝不返回兜底 adapter 的误导性静态列表。
ipcMain.handle('llm:listModels', async () => {
try {
const provider = configService.get<string>('llm.provider') ?? '';
const model = configService.get<string>('llm.model') ?? '';
const apiKey = configService.get<string>('llm.apiKey') || '';
if (!provider || !model) {
return { success: false, error: 'LLM 未配置(Provider/Model 为空),无法获取模型列表' };
}
if (!apiKey && provider !== 'ollama') {
return { success: false, error: 'API Key 未配置,无法获取模型列表' };
}
// 列表基于当前已保存配置 —— 先幂等重载 adapter(配置签名未变时为 no-op
if (!ctx.reloadAdapter()) {
return { success: false, error: 'LLM 配置校验失败,请先在设置中修正配置' };
}
const adapter = agentEngineManager.getAdapter();
if (!adapter.listModels) {
return { success: false, error: '当前 Provider 不支持模型列表查询' };
}
const models = await adapter.listModels();
return { success: true, data: models };
} catch (error) {
log.warn('[AGENT] listModels failed:', (error as Error).message);
return { success: false, error: (error as Error).message };
}
});
// ===== v0.7.2 P3-10: Ollama 模型下载(adapter.pullModel 首次接线 IPC/UI=====
// v0.6.4 P4-1 已实现 pull 的进度回调与外部取消信号,但主进程侧零消费者。
// 契约:同一时刻仅允许一个下载任务(全局 AbortController);
// 进度经 llm:ollamaPullProgress 广播({ model, status, completed?, total? }),
// 结束(成功/失败/取消)统一广播 llm:ollamaPullEnded 供前端复位 UI。
let ollamaPullController: AbortController | null = null;
ipcMain.handle('llm:ollamaPull', async (_event, modelName: unknown) => {
if (typeof modelName !== 'string' || !modelName.trim()) {
return { success: false, error: 'Invalid model name' };
}
// 模型名仅进入 JSON body(不经 shell),仍做字符白名单防转义边界
// 合法形态如 qwen3:8b / llama3.1:70b-instruct-q4_K_M / user/model
if (!/^[A-Za-z0-9._:/-]+$/.test(modelName.trim())) {
return { success: false, error: `Invalid model name: ${modelName.slice(0, 60)}` };
}
const adapter = agentEngineManager.getAdapter();
if (!(adapter instanceof OllamaAdapter)) {
return { success: false, error: '仅 Ollama Provider 支持模型下载' };
}
if (ollamaPullController) {
return { success: false, error: '已有模型下载任务进行中,请先取消' };
}
const controller = new AbortController();
ollamaPullController = controller;
const trimmed = modelName.trim();
try {
await adapter.pullModel(
trimmed,
(progress) => {
broadcast('llm:ollamaPullProgress', { model: trimmed, ...progress });
},
controller.signal,
);
return { success: true };
} catch (error) {
return {
success: false,
error: (error as Error).message,
aborted: controller.signal.aborted,
};
} finally {
ollamaPullController = null;
broadcast('llm:ollamaPullEnded', { model: trimmed });
}
});
ipcMain.handle('llm:ollamaPullCancel', async () => {
if (!ollamaPullController) {
return { success: false, error: '没有进行中的下载任务' };
}
ollamaPullController.abort();
return { success: true };
});
// ===== 中断会话 =====
ipcMain.handle('agent:abortSession', async (_event, sessionId) => {
log.info('[AGENT] Abort:', sessionId);
// P2-10: 联动中断该会话派生的所有 SubAgent(消除"会话停了子任务还在跑")
// v0.5.1: 记录被中止的 taskId — SubAgent 的 pending 确认以 taskId 为 sessionId
// 需一并清理,否则中止后孤儿工具在用户补批残留弹框时会真实执行副作用
const abortedTaskIds = orchestrator.abortByParent(sessionId);
agentEngineManager.abort(sessionId);
// MT-1 修复: 等待当前 run 完全结束再返回,防止用户立即重发时新消息卡在等待中
await agentEngineManager.waitForAbort(sessionId);
// v0.3.0 修复: 清理所有等待中的工具确认,避免定时器泄漏和超时 toast 在新会话中弹出
// v0.5.0: 按会话清理 — 只拒绝被中断会话的 pending,不影响其他并发会话等待中的确认
confirmationHook.clearPending(sessionId);
// v0.5.1: 被中止 SubAgent 的 pending 确认一并拒绝(含 SubAgent 递归派生的孙任务)
// v0.7.3 P2-3: 被中止的 SubAgent 是会话终态 —— 用 forgetSession 连决策记忆一并清理
// v0.7.4 P3-5: 被中止 SubAgent 的 TRACE 状态(subTraces/subMeta)一并清理 ——
// 旧实现只覆盖 taskCompleted/taskError 路径,被 abort 的 SubAgent 残留 Map 条目
for (const taskId of abortedTaskIds) {
confirmationHook.forgetSession(taskId);
subTraces.delete(taskId);
subMeta.delete(taskId);
// 中止的 SubAgent 录制文件也收尾(TRACE 完整性:标记为中断终止)
await sessionRecorder.stopRecording(taskId, {
totalIterations: 0,
totalTokens: 0,
durationMs: 0,
terminationReason: 'user_interrupt',
});
}
// TOOL 层:记录中断
auditService.log({
sessionId,
eventType: 'session_end',
actor: 'user',
target: 'session',
details: { reason: 'user_abort' },
outcome: 'denied',
});
return { success: true };
});
// ===== v0.8.0 P1-1a: 后台会话流回放(切回会话时恢复运行内容) =====
// 渲染层 setCurrentSession 检测到目标会话在 sessionRunStates 中运行(或刚结束)
// 时调用本通道:返回回放缓冲中按序保存的 streamEvent + stateChange 事件,
// 渲染层经事件总线灌入 useAgentStream 管线重建消息/Trace 状态,实时事件随后
// 无缝衔接(runId 对齐由缓冲记录的 runId 提供)。
ipcMain.handle('agent:getReplayState', (_event, sessionId: unknown) => {
if (typeof sessionId !== 'string' || !sessionId) {
return { success: false, error: 'Invalid sessionId' };
}
const buf = replayBuffers.get(sessionId);
return {
success: true,
data: {
isRunning: agentEngineManager.isRunning(sessionId),
runId: buf?.runId ?? null,
truncated: buf?.truncated ?? false,
events: buf?.events ?? [],
},
};
});
// ===== v0.5.0: SubAgent 可观测性 =====
// 1) 生命周期事件广播给前端(AgentMonitor 的 SubAgent 状态区)
// 2) SubEngine 流事件录制到独立 TRACE 文件(sessionId = taskId),不污染父会话的流
// 此前 orchestrator 的 6 个事件全项目零消费者,SubAgent 执行过程对 UI 与 TRACE 完全不可见
/** SubAgent 元信息(description/depth,供完成/失败事件广播时补全载荷) */
const subMeta = new Map<string, { description: string; depth: number }>();
/** SubAgent 迭代录制状态(THINKING 驱动新迭代,text 累积供 llm_response */
interface SubTraceState {
iteration: number;
startedAt: number;
text: string;
usage: { input: number; output: number; total: number };
responded: boolean;
/** v0.8.0 P0-1: 本轮 Provider 原生停止原因(DONE 事件携带) */
finishReason?: string;
}
const subTraces = new Map<string, SubTraceState>();
const finishSubTrace = async (
taskId: string,
reason: string,
durationMs: number,
iterations: number,
): Promise<void> => {
await sessionRecorder.stopRecording(taskId, {
totalIterations: iterations,
totalTokens: 0,
durationMs,
terminationReason: reason,
});
subTraces.delete(taskId);
subMeta.delete(taskId);
// v0.7.3 P2-3: SubAgent 终态 —— 决策记忆随任务终结清理(防长期运行泄漏)
confirmationHook.forgetSession(taskId);
};
orchestrator.on(
'taskDelegated',
(d: { taskId: string; description: string; parentSessionId: string; depth: number }) => {
broadcast('subagent:event', {
taskId: d.taskId,
parentSessionId: d.parentSessionId,
description: d.description,
status: 'delegated',
depth: d.depth,
});
subMeta.set(d.taskId, { description: d.description, depth: d.depth });
sessionRecorder.startRecording(d.taskId);
},
);
orchestrator.on(
'taskStarted',
(d: { taskId: string; description: string; parentSessionId: string; depth: number }) => {
broadcast('subagent:event', {
taskId: d.taskId,
parentSessionId: d.parentSessionId,
description: d.description,
status: 'running',
depth: d.depth,
});
},
);
orchestrator.on(
'taskCompleted',
(r: {
taskId: string;
parentSessionId: string;
success: boolean;
durationMs: number;
iterations: number;
}) => {
const meta = subMeta.get(r.taskId);
broadcast('subagent:event', {
taskId: r.taskId,
parentSessionId: r.parentSessionId,
description: meta?.description ?? '',
status: r.success ? 'completed' : 'error',
depth: meta?.depth ?? 1,
durationMs: r.durationMs,
iterations: r.iterations,
});
void finishSubTrace(r.taskId, r.success ? 'completed' : 'error', r.durationMs, r.iterations);
},
);
orchestrator.on(
'taskError',
(d: { taskId: string; parentSessionId: string; description: string; error: string }) => {
const meta = subMeta.get(d.taskId);
broadcast('subagent:event', {
taskId: d.taskId,
parentSessionId: d.parentSessionId,
description: d.description ?? meta?.description ?? '',
status: 'error',
depth: meta?.depth ?? 1,
error: d.error,
});
void finishSubTrace(d.taskId, 'error', 0, 0);
},
);
// SubEngine 状态事件 → TRACE 迭代录制(与主管道相同的事件形状,sessionId = taskId
orchestrator.on(
'subStateChange',
({
taskId,
data,
}: {
taskId: string;
data: { state?: string; current?: string; iteration?: number };
}) => {
const st = subTraces.get(taskId);
const stateValue = data.state ?? data.current ?? '';
// THINKING 且迭代号变化 → 新迭代开始(关闭上一迭代)
if (
stateValue === 'THINKING' &&
data.iteration != null &&
(!st || st.iteration !== data.iteration)
) {
if (st && !st.responded) {
sessionRecorder.recordLLMResponse({
sessionId: taskId,
iteration: st.iteration,
content: st.text,
// v0.8.0 P0-1: 记录真实停止原因(与主管道同口径)
finishReason: st.finishReason ?? 'stop',
tokenUsage: st.usage,
});
sessionRecorder.recordIterationEnd(taskId, {
iteration: st.iteration,
durationMs: Date.now() - st.startedAt,
});
}
subTraces.set(taskId, {
iteration: data.iteration,
startedAt: Date.now(),
text: '',
usage: { input: 0, output: 0, total: 0 },
responded: false,
});
sessionRecorder.recordIterationStart(taskId, data.iteration);
const provider = configService.get<string>('llm.provider') ?? '';
const model = configService.get<string>('llm.model') ?? '';
sessionRecorder.recordLLMRequest({
sessionId: taskId,
iteration: data.iteration,
provider,
model,
messageCount: data.iteration + 1,
});
return;
}
// TERMINATED → 补记最终迭代(正常流程只在下一轮 THINKING 补记,最终轮无后续)
if (stateValue === 'TERMINATED' && st) {
if (!st.responded) {
sessionRecorder.recordLLMResponse({
sessionId: taskId,
iteration: st.iteration,
content: st.text,
// v0.8.0 P0-1: 记录真实停止原因(与主管道同口径)
finishReason: st.finishReason ?? 'stop',
tokenUsage: st.usage,
});
}
sessionRecorder.recordIterationEnd(taskId, {
iteration: st.iteration,
durationMs: Date.now() - st.startedAt,
});
}
},
);
// SubEngine 流事件 → TRACE 内容录制(文本累积 + 工具调用/结果)
orchestrator.on(
'subStreamEvent',
({ taskId, event }: { taskId: string; event: MetonaStreamEvent }) => {
const st = subTraces.get(taskId);
if (!st) return;
switch (event.type) {
case MetonaStreamEventType.TEXT_DELTA:
if (event.delta) st.text += event.delta;
break;
case MetonaStreamEventType.USAGE:
if (event.usage) {
st.usage = {
input: event.usage.inputTokens ?? 0,
output: event.usage.outputTokens ?? 0,
total: event.usage.totalTokens ?? 0,
};
}
break;
case MetonaStreamEventType.DONE:
// v0.8.0 P0-1: 捕获 Provider 原生停止原因(SubAgent TRACE 真值)
if (event.finishReason) st.finishReason = event.finishReason;
break;
case MetonaStreamEventType.TOOL_CALL_COMPLETE:
if (event.toolCall) {
sessionRecorder.recordToolCall({
sessionId: taskId,
iteration: event.iteration,
toolName: event.toolCall.name,
args: event.toolCall.args,
});
}
break;
case MetonaStreamEventType.TOOL_RESULT:
if (event.toolResult) {
const resultPreview =
typeof event.toolResult.result === 'string'
? event.toolResult.result
: JSON.stringify(event.toolResult.result);
sessionRecorder.recordToolResult({
sessionId: taskId,
iteration: event.iteration,
toolName: event.toolResult.toolName,
success: event.toolResult.success,
durationMs: event.toolResult.durationMs ?? 0,
resultPreview,
error: event.toolResult.error,
});
}
break;
default:
break;
}
},
);
}