Files
metona-ai-desktop/electron/services/agent-engine-manager.service.ts
T
thzxx 2230bcec3f feat: v0.4.0 四阶段迭代 — 安全加固 + 工程基线 + 架构重构 + 双 Provider 扩展
P0 安全修复:
- API Key 加密存储(safeStorage 密钥链,版本化前缀,历史明文平滑兼容)
- 间接提示注入防护(SecurityScanHook 工具结果深扫描,网络工具脱敏/本地工具警示分级)
- error:report IPC 断链修复(渲染进程错误上报落 electron-log + 审计)
- abort 信号贯通工具层(run_command/dev-tools 子进程随会话中断终止)
- run_command 沙箱加固(cd 系统目录/敏感文件读取拦截 + chcp 前缀剥离防解析退化)
- .env 真实生效(dotenv 回退加载,应用内配置优先)

P1 工程基础:
- ESLint 9 flat config + 全部 34 条存量 warnings 清零(零容忍基线)
- 测试基线 118 用例 11 文件(token/文件防护/权限/沙箱/注入/命令/引擎/注册表/审计链/摘要分层)
- test:electron 双模式(ELECTRON_RUN_AS_NODE 跑 Electron ABI,SQLite 套件全执行)
- SessionRecorder 多会话隔离 + 9 种 TRACE 事件补全(含最终轮 iteration_end)
- Provider 故障转移(重试耗尽/不可重试一次性切换 fallback + 前端通知)
- MCP 真就绪(等待全部连接完成再广播 tools:ready)
- SLO/HealthChecker 真实接入(60s 巡检 + 托盘状态)
- CONFIG_DEFAULTS 单一来源(消除 SEED 双源漂移)

P2 架构升级:
- handlers.ts 1940 行拆分为 13 个 IPC 域模块(防重入注册 + 多窗口广播)
- AgentEngineManager 每会话独立引擎(LRU 30 + adapter 工厂隔离 abort 信号)
- TaskOrchestrator EngineProvider 改造 + abortByParent 联动中断 SubAgent
- 会话摘要分层上下文(session_summaries 滚动摘要 + 截断游标清理防因果污染)
- 消息编辑重发/重新生成(truncateAfter IPC + store 动作 + UI)
- Markdown 导出 / WebSearch 并行抓取(并发 3)/ 记忆 TF 缓存 / 版本构建期注入

P3 能力扩展:
- OpenAI Adapter(o 系列推理模型 reasoning_effort/max_completion_tokens)
- Anthropic Adapter(原生 Messages API:tool_use 块/角色合并/thinking budget/图片 base64/SSE 事件机)
- 设置页/Onboarding 六 Provider 全链路接入
2026-08-20 23:17:02 +08:00

205 lines
7.5 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.
/**
* Agent Engine Manager — 每会话独立引擎管理器(P2-10)
*
* 解决原"全局单引擎"的两个缺陷:
* 1. 全局串行锁:原 AgentLoopEngine.currentRunPromise 使所有会话共享一把锁,
* 上一会话未结束时新会话消息需排队等待(最长卡 120s 工具超时)。
* 现在每个会话持有独立引擎实例,多会话可并行运行。
* 2. adapter abort 信号互踩:原所有引擎/SubAgent 共享一个 adapter 实例,
* setAbortSignal 单槽位导致并发时中断信号错乱。
* 现在创建引擎时通过 adapter 工厂为每个引擎生成独立 adapter 实例
* (adapter 是无状态的配置包装,实例化成本可忽略)。
*
* 引擎生命周期:
* - 按需创建(首次 sendMessage 时),事件统一转发到 manager(附加 sessionId
* - LRU 淘汰:缓存超过 30 个引擎时,淘汰最旧的非运行中引擎
*
* @see electron/harness/agent-loop/engine.ts — 引擎实现
*/
import { EventEmitter } from 'events';
import { AgentLoopEngine } from '../harness/agent-loop';
import type { AgentLoopConfig } from '../harness/agent-loop/types';
import type { IMetonaProviderAdapter, MetonaToolDef } from '../harness/types';
import type { ToolRegistry } from '../harness/tools/registry';
import type { PreToolHook } from '../harness/hooks/pre-tool';
import type { PostToolHook } from '../harness/hooks/post-tool';
import log from 'electron-log';
/** 引擎缓存上限(超过后淘汰最旧的非运行中引擎) */
const MAX_ENGINES = 30;
export class AgentEngineManager extends EventEmitter {
private engines = new Map<string, AgentLoopEngine>();
/** 运行中的会话(stateChange INIT 添加 / TERMINATED 移除),用于 LRU 淘汰保护 */
private running = new Set<string>();
/** 主 adapter(供 MemoryConsolidator 等共享组件使用) */
private primaryAdapter: IMetonaProviderAdapter;
private fallbackAdapter: IMetonaProviderAdapter | null = null;
private baseConfig: Partial<AgentLoopConfig>;
private workspacePath = '';
constructor(private opts: {
/** adapter 工厂(每次调用返回新实例;闭包内读取最新配置) */
buildAdapter: () => IMetonaProviderAdapter;
baseConfig: Partial<AgentLoopConfig>;
toolRegistry?: ToolRegistry;
preToolHooks?: PreToolHook[];
postToolHooks?: PostToolHook[];
}) {
super();
this.primaryAdapter = opts.buildAdapter();
this.baseConfig = { ...opts.baseConfig };
}
/** 主 adapter(供 consolidator / orchestrator 等共享使用) */
getAdapter(): IMetonaProviderAdapter {
return this.primaryAdapter;
}
/** 创建独立 adapter 实例(每引擎/SubAgent 独享,避免 abort 信号互踩) */
createAdapter(): IMetonaProviderAdapter {
return this.opts.buildAdapter();
}
getFallbackAdapter(): IMetonaProviderAdapter | null {
return this.fallbackAdapter;
}
/** 设置故障转移 Provider(同步到所有引擎) */
setFallbackAdapter(adapter: IMetonaProviderAdapter | null): void {
this.fallbackAdapter = adapter;
for (const engine of this.engines.values()) {
engine.setFallbackAdapter(adapter);
}
}
getWorkspacePath(): string {
return this.workspacePath;
}
setWorkspacePath(path: string): void {
this.workspacePath = path;
for (const engine of this.engines.values()) {
engine.setWorkspacePath(path);
}
}
/** 同步工具列表到所有引擎(工具开关变更 / MCP 注册完成时) */
setToolsAll(tools: MetonaToolDef[]): void {
for (const engine of this.engines.values()) {
engine.setTools(tools);
}
}
/** 热更新所有引擎配置(设置变更时) */
updateConfigAll(partial: Partial<AgentLoopConfig>): void {
this.baseConfig = { ...this.baseConfig, ...partial };
for (const engine of this.engines.values()) {
engine.updateConfig(partial);
}
}
/**
* 重建所有 adapterLLM 配置变更时由 reloadAdapter 调用)
* 工厂闭包读取最新配置,重建 primary + 各引擎独立实例。
*/
refreshAdapters(): void {
this.primaryAdapter = this.opts.buildAdapter();
for (const engine of this.engines.values()) {
engine.setAdapter(this.createAdapter());
engine.setFallbackAdapter(this.fallbackAdapter);
}
}
/** 获取(或创建)会话引擎 */
getEngine(sessionId: string): AgentLoopEngine {
let engine = this.engines.get(sessionId);
if (!engine) {
engine = this.createEngine(sessionId);
this.engines.set(sessionId, engine);
this.evictIdleEngines();
}
return engine;
}
/** 中断指定会话引擎 */
abort(sessionId: string): void {
this.engines.get(sessionId)?.abort();
}
/** 等待指定会话当前 run 结束 */
async waitForAbort(sessionId: string, timeoutMs = 5_000): Promise<boolean> {
return this.engines.get(sessionId)?.waitForAbort(timeoutMs) ?? true;
}
/** 当前引擎数量(测试/诊断用) */
get size(): number {
return this.engines.size;
}
// ===== 私有方法 =====
private createEngine(sessionId: string): AgentLoopEngine {
const engine = new AgentLoopEngine(
{
...this.baseConfig,
contextWindow: this.primaryAdapter.getContextWindow(),
},
this.createAdapter(),
this.opts.toolRegistry,
this.opts.preToolHooks ?? [],
this.opts.postToolHooks ?? [],
);
engine.setFallbackAdapter(this.fallbackAdapter);
if (this.workspacePath) engine.setWorkspacePath(this.workspacePath);
this.forwardEngineEvents(engine, sessionId);
log.debug(`[EngineManager] engine created for session ${sessionId} (total: ${this.engines.size + 1})`);
return engine;
}
/** 将引擎事件转发到 manager(统一附加 sessionId,供常驻监听器消费) */
private forwardEngineEvents(engine: AgentLoopEngine, sessionId: string): void {
engine.on('streamEvent', (event) => {
this.emit('streamEvent', { ...event, sessionId: event.sessionId || sessionId });
});
engine.on('stateChange', (data) => {
const payload = { ...data, sessionId: data.sessionId || sessionId };
// 维护运行中集合(LRU 淘汰保护)
if (payload.state === 'INIT' || payload.current === 'INIT') this.running.add(sessionId);
if (payload.state === 'TERMINATED' || payload.current === 'TERMINATED') this.running.delete(sessionId);
this.emit('stateChange', payload);
});
engine.on('complete', (data) => {
this.running.delete(sessionId);
this.emit('complete', { ...data, sessionId: data.sessionId || sessionId });
});
engine.on('compressed', (data) => {
this.emit('compressed', { ...data, sessionId });
});
engine.on('deadLoop', (data) => {
this.emit('deadLoop', { ...data, sessionId: data.sessionId || sessionId });
});
engine.on('providerSwitched', (data) => {
this.emit('providerSwitched', { ...data, sessionId: data.sessionId || sessionId });
});
engine.on('aborted', () => {
this.emit('aborted', { sessionId });
});
}
/** LRU 淘汰:超过上限时删除最旧的非运行中引擎 */
private evictIdleEngines(): void {
if (this.engines.size <= MAX_ENGINES) return;
let toEvict = this.engines.size - MAX_ENGINES;
for (const [sessionId, engine] of this.engines) {
if (toEvict <= 0) break;
if (this.running.has(sessionId)) continue;
engine.destroy();
this.engines.delete(sessionId);
log.debug(`[EngineManager] evicted idle engine for session ${sessionId}`);
toEvict--;
}
}
}