diff --git a/Agent.md b/Agent.md index 59a2a91..2de8698 100644 --- a/Agent.md +++ b/Agent.md @@ -66,7 +66,7 @@ EMRG is a self-evolving AI agent architecture experiment. Python implementation, - Streaming chat with delta rendering (16ms batching), markdown on done (marked + DOMPurify + local highlight.js subset), tool call status cards (2000-char truncation + expand) - Session list/switch/new/delete + right-click rename (context menu, #423) synced with daemon; own-stream busy lock (G65); broadcast streams from other clients tagged "来自其他客户端" - Disconnect/reconnect: red status dot, auto daemon respawn (stale-port detection), session resume, input bar restored on disconnect (no 30s fake-timeout) - - Unit tests `npm test` (127: 37 daemon_client + 8 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config); RESPONSE_TYPES mirror daemon protocol verified against `daemon.py` + - Unit tests `npm test` (142: 43 daemon_client + 17 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config); RESPONSE_TYPES mirror daemon protocol verified against `daemon.py` - **Auto project tracking** — Automatically detects and records working directories; project-scoped sessions - **Rant-driven evolution** — User feedback via `/rant` drives automatic self-improvement cycles - **Headless GitHub auth** — Non-interactive evolution auto-extracts `GH_TOKEN` from git credential store (osxkeychain / credential helper); PR comment/LGTM queries fall back to REST API (GraphQL needs `read:org` scope) @@ -94,7 +94,7 @@ pkill -f "emrg.server"; rm -f ~/.emrg/emrgd.port; python -m emrg ``` Python: `uv run pytest tests/ -v` (680) — import check: `uv run python -c "from emrg.client.app import run_client"` -GUI: `cd emrg/gui && npm test` (127: 37 daemon_client + 8 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config) — syntax: `node --check main.js preload.js daemon_client.js renderer/js/*.js` +GUI: `cd emrg/gui && npm test` (142: 43 daemon_client + 17 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config) — syntax: `node --check main.js preload.js daemon_client.js renderer/js/*.js` CI: `uv run pytest` + GUI tests + **actionlint workflow lint** (`rhysd/actionlint@v1.7.12` gate, #444 — workflow 解析错误在 PR CI 即失败,如 `if:` secrets 上下文) Re-trigger: `scripts/re-trigger-ci.sh [branch]` (workflow_dispatch, #527 — 替代空 commit 重触发:Actions outage 会整段丢弃 push 事件,dispatch 走 API 路径不受影响) diff --git a/README.cn.md b/README.cn.md index 758a60b..7878a73 100644 --- a/README.cn.md +++ b/README.cn.md @@ -282,7 +282,7 @@ uv run python -m emrg # 启动 TUI cd emrg/gui npm ci # 安装依赖(生产模式可 --omit=dev) npm start # 启动 GUI(自动拉起 daemon) -npm test # 运行 Node 测试(127 项:37 daemon_client + 8 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config;集成测试在 CI 跑,本地可 npm run test:integration) +npm test # 运行 Node 测试(142 项:43 daemon_client + 17 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config;集成测试在 CI 跑,本地可 npm run test:integration) ``` CI 通过 GitHub Actions 自动运行测试并检查冲突标记(`.github/workflows/test.yml`)。 diff --git a/README.md b/README.md index 6021510..4cd2e58 100644 --- a/README.md +++ b/README.md @@ -281,7 +281,7 @@ uv run python -m emrg # launch TUI cd emrg/gui npm ci # install deps (production: --omit=dev) npm start # launch GUI (auto-starts daemon) -npm test # run Node tests (127: 37 daemon_client + 8 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config; integration runs in CI, local: npm run test:integration) +npm test # run Node tests (142: 43 daemon_client + 17 conn-manager + 22 app-commands + 32 renderer smoke + 15 i18n + 7 integration + 3 commands + 3 build-config; integration runs in CI, local: npm run test:integration) ``` CI runs tests and checks for conflict markers automatically via GitHub Actions (`.github/workflows/test.yml`). diff --git a/emrg/gui/conn-manager.js b/emrg/gui/conn-manager.js index 2986ff1..e01af96 100644 --- a/emrg/gui/conn-manager.js +++ b/emrg/gui/conn-manager.js @@ -4,38 +4,50 @@ // open session (each session = one independent websocket connection, aligned // with the TUI multi-open model). // -// This slice establishes the open/close/get contract and the daemon-ownership -// bootstrap. main.js rewiring to use this manager (single-session no-regression -// target) lands in a later slice — until then main.js keeps its single client. +// This slice completes the P2 wiring contract: +// - `ensureDaemon()` keeps a daemon-level connection (`_daemonConn`) used for +// non-session commands (ping / list_sessions / set_model / github_* / ...). +// Session connections only connect to the already-running daemon +// (ensureConnected({ skipStart: true }) — never spawn). +// - `open(sid, projectPath)` creates the per-session connection with +// per-connection delta batching (deltaBatchMs=16, #626) and resumes the +// session (auto-subscribe). `{ resume: false }` skips resume for brand-new +// sessions (daemon implicitly subscribes on first task message). +// - restart recovery: short-window all-drop → recoverAll (re-open + re-subscribe). +// - `onOpen` / `onRecovered` hooks let main.js attach the renderer event +// bridge (sid-tagged) and refresh UI state after recovery. // // Design notes (from the rant): // - Each open session = one independent ws connection → natural isolation, no // event routing. -// - connManager ensures the daemon is ready (spawn if missing) before opening; -// session connections then use ensureConnected({ skipStart: true }) — they -// only connect to the already-running daemon, never spawn. // - resume_session(sid, cwd=projectPath) auto-subscribes the connection. // - Already-open sid → reuse the existing connection (no duplicate). const { DaemonClient } = require("./daemon_client.js"); class ConnManager { - constructor({ projectDir, logger = console, isPackaged = false, restartWindowMs = 1000 } = {}) { + constructor({ projectDir, logger = console, isPackaged = false, restartWindowMs = 1000, singleRetryDelayMs = 1000 } = {}) { this.projectDir = projectDir; this.logger = logger; this.isPackaged = isPackaged; this._conns = new Map(); // sid -> { conn, projectPath } + this._daemonConn = null; // daemon 级连接(ping/list_sessions 等非会话命令) + this._openHooks = new Set(); // (sid, conn) => void(新会话连接建立时) + this._recoverHooks = new Set(); // () => void(recoverAll 完成后) // daemon 重启恢复(rant 15:07:19 P2):短窗口内所有连接同时断 → 判定 daemon 重启 - // → 全部重连重订阅;单条断 → 不做全量恢复(留给独立退避重试)。 + // → 全部重连重订阅;单条断 → 独立退避重试(多会话场景,单会话全断走恢复)。 this._restartWindowMs = restartWindowMs; this._disconnects = new Map(); // sid -> timestamp(最近一次断连) this._recovering = false; // 恢复中守卫(防 close→disconnect→recoverAll 递归) + this._singleRetryDelayMs = singleRetryDelayMs; + this._singleRetries = new Map(); // sid -> timer(单连接独立退避在途) } // 确保 daemon 已运行(connManager = daemon 生命周期唯一 owner)。 - // 引导 client ensureConnected():port 文件缺失 → spawn;已运行 → 直连。 - // 连接后立即关闭引导连接——会话连接统一走 skipStart(只连不拉)。 - async _ensureDaemon() { + // 引导连接保留为 _daemonConn:供 ping/list_sessions 等非会话命令使用; + // 会话连接统一 skipStart(只连不拉)。已连接 → 直接复用。 + async ensureDaemon() { + if (this._daemonConn && this._daemonConn.connected) return this._daemonConn; const boot = new DaemonClient({ projectDir: this.projectDir, logger: this.logger, @@ -43,39 +55,67 @@ class ConnManager { }); try { await boot.ensureConnected(); - } finally { + } catch (e) { boot.close(); + throw e; } + this._daemonConn = boot; + return boot; + } + + // daemon 级连接访问器(main.js 非会话命令用;未建立返回 null) + daemonConn() { + return this._daemonConn || null; } // open(sid, projectPath):创建 DaemonClient → ensureConnected(skipStart) → - // resume_session(自动订阅)。已打开的 sid 直接复用现有连接。 - async open(sid, projectPath) { + // resume_session(自动订阅;resume:false 跳过,供新会话首条消息前用)。 + // 已打开的 sid 直接复用现有连接。失败时关闭连接不泄漏。 + async open(sid, projectPath, { resume = true } = {}) { const existing = this._conns.get(sid); - if (existing) return existing.conn; - await this._ensureDaemon(); + if (existing) { + if (existing.conn.connected) return existing.conn; + this.close(sid); // 残留断连连接 → 关闭重开(_intentionalClose 标记抑制断线横幅) + } + await this.ensureDaemon(); const conn = new DaemonClient({ projectDir: this.projectDir, logger: this.logger, isPackaged: this.isPackaged, + deltaBatchMs: 16, // P2:delta 批量(G122 16ms)每连接一份(#626) }); - await conn.ensureConnected({ skipStart: true }); // daemon 已就绪 → 只连不拉 - await conn.sendCommandAndWait("resume_session", { session_id: sid, cwd: projectPath }, 5000); + try { + await conn.ensureConnected({ skipStart: true }); // daemon 已就绪 → 只连不拉 + if (resume) { + await conn.sendCommandAndWait("resume_session", { session_id: sid, cwd: projectPath }, 5000); + } + } catch (e) { + conn.close(); // resume 失败(会话已删等)→ 不泄漏连接 + throw e; + } // 断开监听 → 重启恢复判定(仅当所有打开会话在同一短窗口内断开) conn.onEvent((type) => { if (type === "disconnected") { + if (conn._intentionalClose) return; // 主动关闭(切走/删除)不参与重启判定 this._disconnects.set(sid, Date.now()); this._onDisconnect(sid); } }); this._conns.set(sid, { conn, projectPath }); + for (const cb of this._openHooks) { + try { cb(sid, conn); } catch (e) { this.logger.warn(`[gui] connManager onOpen hook error: ${e.message}`); } + } return conn; } // close(sid):conn.close(断开 ws)→ 移除。返回是否有关闭对象。 + // 标记 _intentionalClose:主动关闭(切走/删除)不触发 renderer 断线横幅 + // (桥检查该标记;真断连/daemon 重启的 disconnected 照常转发)。 close(sid) { const entry = this._conns.get(sid); if (!entry) return false; + this._cancelSingleRetry(sid); // 主动关闭 → 取消该会话的独立退避 + entry.conn._intentionalClose = true; entry.conn.close(); this._conns.delete(sid); this._disconnects.delete(sid); @@ -94,6 +134,20 @@ class ConnManager { closeAll() { for (const sid of [...this._conns.keys()]) this.close(sid); + if (this._daemonConn) { + this._daemonConn.close(); + this._daemonConn = null; + } + } + + // 新会话连接建立钩子(main.js 挂 renderer 事件桥;recoverAll 重开路径同样触发) + onOpen(callback) { + this._openHooks.add(callback); + } + + // recoverAll 完成钩子(main.js 刷新 UI 状态:status connected + sessions + pong) + onRecovered(callback) { + this._recoverHooks.add(callback); } // ── daemon 重启恢复(rant 15:07:19 P2)────────────────────────────── @@ -118,10 +172,45 @@ class ConnManager { this.recoverAll().catch((e) => this.logger.warn(`[gui] connManager recover failed: ${e.message}`) ); + } else { + // 单条断(多会话场景,非重启)→ 独立退避重试 + this._scheduleSingleRetry(sid); } } - // 全部重连重订阅(复用 open 序列:引导 daemon 就绪 → skipStart 会话连接 → + // 单连接独立退避(rant 15:07:19 P2:单条断 → 独立退避重试)。 + // 退避到期后若该会话连接仍断 → open() 重开(stale → 关闭重开 + resume 重订阅)。 + _scheduleSingleRetry(sid) { + if (this._singleRetries.has(sid)) return; // 已有退避在途 + this.logger.info( + `[gui] connManager: session ${sid} dropped (not all) — independent backoff retry in ${this._singleRetryDelayMs}ms` + ); + const timer = setTimeout(() => { + this._singleRetries.delete(sid); + this._retrySingle(sid).catch((e) => + this.logger.warn(`[gui] connManager: session ${sid} retry failed: ${e.message}`) + ); + }, this._singleRetryDelayMs); + timer.unref?.(); + this._singleRetries.set(sid, timer); + } + + _cancelSingleRetry(sid) { + const t = this._singleRetries.get(sid); + if (t) { + clearTimeout(t); + this._singleRetries.delete(sid); + } + } + + async _retrySingle(sid) { + const entry = this._conns.get(sid); + if (!entry || entry.conn.connected) return; // 已重开/已主动关闭 + this.logger.info(`[gui] connManager retrying session ${sid}`); + await this.open(sid, entry.projectPath); // stale → close + reopen(含 resume 重订阅) + } + + // 全部重连重订阅(复用 open 序列:ensureDaemon → skipStart 会话连接 → // resume_session)。单会话恢复失败跳过不阻塞其余(写盘/重试由后续片处理)。 async recoverAll() { if (this._recovering) return; @@ -132,6 +221,8 @@ class ConnManager { })); try { for (const { sid } of sessions) this.close(sid); + for (const t of this._singleRetries.values()) clearTimeout(t); + this._singleRetries.clear(); this._disconnects.clear(); for (const { sid, projectPath } of sessions) { try { @@ -144,6 +235,9 @@ class ConnManager { } finally { this._recovering = false; } + for (const cb of this._recoverHooks) { + try { cb(); } catch (e) { this.logger.warn(`[gui] connManager onRecovered hook error: ${e.message}`); } + } } } diff --git a/emrg/gui/daemon_client.js b/emrg/gui/daemon_client.js index fc011c7..1f0495e 100644 --- a/emrg/gui/daemon_client.js +++ b/emrg/gui/daemon_client.js @@ -93,6 +93,18 @@ class DaemonClient { this._deltaBatchMs = deltaBatchMs; this._deltaBuf = []; this._deltaTimer = null; + // P2 connManager(rant 2026-08-10T15:07:19):G65 自有流锁每连接一份。 + // 从 main.js 全局移入——多会话各自独立:本连接发出的 task 流运行中 → + // ownStream=true,切会话/关连接前必须释放。 + this.ownStream = false; + this.ownStreamRequestId = null; + } + + // 释放自有流锁(G65)。done(request 匹配或 timeout 兜底)、session busy 即发 + // 错误、cancelled(request 匹配)与断连时调用;sendTask 抛异常时由调用方清理。 + _releaseOwnStream() { + this.ownStream = false; + this.ownStreamRequestId = null; } // ── 生命周期 ──────────────────────────────────────────── @@ -492,6 +504,12 @@ class DaemonClient { }; if (mode && mode !== "auto") payload.mode = mode; this._setCurrentStream(rid); + // G65:自有流锁——本连接发出流式 task 即标记,done/error/cancelled/断连释放 + // (多会话各自独立;main.js emrg:sendMessage 的 G65 切会话检查读本字段) + if (stream) { + this.ownStream = true; + this.ownStreamRequestId = rid; + } this.ws.send(JSON.stringify(payload)); return rid; } @@ -579,6 +597,8 @@ class DaemonClient { } if (frame.type === "cancelled") { this._flushDeltaBuf(); // 终态前冲刷(rant 14:11 同源:delta 不晚于终态) + // 自有流取消 → 释放 G65 锁(带 request_id 的 cancelled 明确是本流的终态) + if (frame.request_id === this.ownStreamRequestId) this._releaseOwnStream(); this._emit("cancelled", frame); return; } @@ -610,6 +630,9 @@ class DaemonClient { } if (frame.error) { this._flushDeltaBuf(); // 终态前冲刷(rant 14:11 同源) + // session busy 是即发错误(daemon 返回后无 done 跟随)——释放自有流锁,防 G65 锁泄漏 + // (流式错误如 LLM error 则有 done 跟随,由 done 分支释放,不在此处理) + if (frame.error && String(frame.error).includes("session busy")) this._releaseOwnStream(); this._emit("error", frame); return; } @@ -644,6 +667,10 @@ class DaemonClient { if (rid && this._currentStream && this._currentStream.requestId === rid) { this.clearActiveStream(); } + // G65:仅自有流的 done 释放锁(广播 done 不影响);timeout 兜底同样只清自有 + if (rid === this.ownStreamRequestId || (frame.timeout && this.ownStream)) { + this._releaseOwnStream(); + } // G83:done 清理分组缓存(DOM 保留) if (rid) this._cleanupGroup(rid, true); } @@ -717,6 +744,7 @@ class DaemonClient { // G41/G89/G97:断连处理 _onClose() { this.connected = false; + this._releaseOwnStream(); // 断连即释放 G65 自有流锁(防锁泄漏) this._rejectAllPending("connection closed"); this.clearActiveStream(); // G94 timer 清理:断连后 30s 超时 timer 不应再触发(防虚假"响应超时"提示) this.clearGroups(); // G97:断连清空广播分组缓存(含 10 分钟 timer),防"幽灵"分组残留 diff --git a/emrg/gui/main.js b/emrg/gui/main.js index 4735f26..06ba355 100644 --- a/emrg/gui/main.js +++ b/emrg/gui/main.js @@ -11,7 +11,8 @@ const os = require("os"); const path = require("path"); const { spawn } = require("child_process"); const { parse: parseToml, stringify: stringifyToml } = require("smol-toml"); -const { DaemonClient, generateSessionId, SESSION_ID_RE, PORT_FILE } = require("./daemon_client"); +const { generateSessionId, SESSION_ID_RE } = require("./daemon_client"); +const { ConnManager } = require("./conn-manager"); const APP_VERSION = require("./package.json").version; // ── 单实例锁(G85/G120:第二个实例退出并 focus 已有窗口)── @@ -24,12 +25,10 @@ if (!app.requestSingleInstanceLock()) { function main() { const logger = createLogger(); let win = null; - let client = null; + let connManager = null; // P2(rant 15:07:19):连接管理器 = daemon 生命周期唯一 owner let projectDir = os.homedir(); let configExists = false; let currentSessionId = null; - let ownStream = false; // 自有流运行中(G65:禁止切会话) - let ownStreamRequestId = null; // 自有流 request_id(广播 done 不清锁) let reconnectTimer = null; // Rant 2026-08-09T13:16:36 ③/⑤:重连指数退避(1s→2s→4s→…封顶 60s)。 // 之前固定 1s——daemon 缺失时每 5s 一轮 spawn,弹窗/日志风暴。成功连接后复位。 @@ -159,7 +158,7 @@ vision = false resolve(fs.readFileSync(configPath(), "utf8")); return; } - const python = client?._findPython() || "python3"; + const python = connManager?.daemonConn()?._findPython() || "python3"; const child = spawn(python, ["-c", "from emrg.config import ensure_config; ensure_config()"], { cwd: projectDir, stdio: "ignore", @@ -265,19 +264,21 @@ vision = false if (requestId !== undefined && (typeof requestId !== "string" || requestId.length < 8 || requestId.length > 64)) { throw new Error("invalid request_id"); // G143:renderer 预生成 id 的格式护栏 } - if (!client || !client.connected) throw new Error("daemon not connected"); - ownStream = true; + // P2:每会话独立连接——首条消息前自动打开(新会话不 resume,daemon 隐式订阅) + let conn = connManager?.get(sessionId); + if (!conn || !conn.connected) { + conn = await openSession(sessionId, projectDir, { resume: false }); + } let rid; try { // G143:renderer 预生成 requestId(send 前标记自有流,消除 IPC 往返竞态窗口) // WorkBuddy P2:mode="ask" → 纯对话(daemon 不启用工具) - rid = client.sendTask({ sessionId, cwd: projectDir, prompt: text, stream: true, requestId, mode }); + // G65:conn.sendTask 内部标记 ownStream(每连接独立锁) + rid = conn.sendTask({ sessionId, cwd: projectDir, prompt: text, stream: true, requestId, mode }); } catch (e) { - ownStream = false; // sendTask 抛异常(ws.send 失败)→ 释放锁,防 G65 锁泄漏 - ownStreamRequestId = null; + conn._releaseOwnStream(); // sendTask 抛异常(ws.send 失败)→ 释放锁,防 G65 锁泄漏 throw e; } - ownStreamRequestId = rid; // 追踪自有流(G65 锁仅由自有 done 释放) return { ok: true, requestId: rid }; // G124:回传 requestId → renderer 识别自有流 }); @@ -285,33 +286,35 @@ vision = false ipcMain.handle("emrg:switchSession", async (_e, { sessionId }) => { if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); - if (ownStream) throw new Error("stream in progress — cannot switch"); // G65 - if (!client?.connected) throw new Error("daemon not connected"); - client.clearGroups(); // G110:切会话清空旧分组缓存(含 timer),防广播"幽灵"残留 - const meta = await client.sendCommandAndWait("resume_session", { session_id: sessionId, cwd: projectDir }, 5000) - .catch((e) => { - // G106:被动删除恢复——resume error → 刷新列表 + 切最近 - if (/not found|error/i.test(e.message)) { - return listSessions().then((sessions) => { - const next = sessions[0]?.session_id || null; - return { error: "session_not_found", sessions, next_session: next }; - }); - } - throw e; - }); - if (meta.error === "session_not_found") { - // G106:resume 失败(会话已删)→ 不更新 currentSessionId(renderer 会切到 next_session, - // 但 main 侧保持旧值,重连 resume 也不指向已删会话;renderer 随后 switchSession(next) 会纠正) - return meta; + // G65:自有流运行中禁止切会话(每连接独立锁,查当前激活连接) + if (connManager?.get(currentSessionId)?.ownStream) throw new Error("stream in progress — cannot switch"); + const prevSid = currentSessionId; + // G110:切会话清空旧连接分组缓存(含 timer),防广播"幽灵"残留 + connManager?.get(prevSid)?.clearGroups(); + try { + await openSession(sessionId, projectDir); // 打开(新)会话连接 + resume_session 自动订阅 + } catch (e) { + // G106:被动删除恢复——resume error → 刷新列表 + 切最近 + if (/not found|error/i.test(e.message)) { + connManager?.close(sessionId); // 防御:失败连接已由 open 内部关闭 + return listSessions().then((sessions) => { + const next = sessions[0]?.session_id || null; + return { error: "session_not_found", sessions, next_session: next }; + }); + } + throw e; } currentSessionId = sessionId; win.setTitle(`EMRG — ${sessionId}`); // G109 - return meta; + // 单会话模型:切走即关闭旧会话连接(每会话一条连接;多会话 openSessions 由 P4 管理) + if (prevSid && prevSid !== sessionId) connManager?.close(prevSid); + return {}; }); ipcMain.handle("emrg:deleteSession", async (_e, { sessionId }) => { if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); - await client.sendCommandAndWait("delete_session", { session_id: sessionId, cwd: projectDir }, 5000); + await requireConn().sendCommandAndWait("delete_session", { session_id: sessionId, cwd: projectDir }, 5000); + connManager?.close(sessionId); // P2:删除会话 → 关闭该会话连接(若打开) return { ok: true }; }); @@ -319,7 +322,7 @@ vision = false if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); const clean = String(title || "").trim().slice(0, 80); // 截断超长标题 if (!clean) throw new Error("empty title"); - const frame = await client.sendCommandAndWait("rename_session", { session_id: sessionId, cwd: projectDir, title: clean }, 5000); + const frame = await requireConn().sendCommandAndWait("rename_session", { session_id: sessionId, cwd: projectDir, title: clean }, 5000); return { ok: true, title: frame.title || clean }; }); @@ -335,21 +338,21 @@ vision = false ipcMain.handle("emrg:clearSession", async (_e, { sessionId }) => { // GUI / 指令 P1:/clear — 清空当前会话(daemon 协议 clear_session 已存在) if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); - await client.sendCommandAndWait("clear_session", { session_id: sessionId, cwd: projectDir }, 5000); + await requireConn().sendCommandAndWait("clear_session", { session_id: sessionId, cwd: projectDir }, 5000); return { ok: true }; }); ipcMain.handle("emrg:compactSession", async (_e, { sessionId }) => { // GUI / 指令 P1:/compact — 压缩当前会话历史(daemon 协议 compact 已存在) if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); - await client.sendCommandAndWait("compact", { session_id: sessionId, cwd: projectDir }, 5000); + await requireConn().sendCommandAndWait("compact", { session_id: sessionId, cwd: projectDir }, 5000); return { ok: true }; }); ipcMain.handle("emrg:listHistory", async (_e, { sessionId }) => { // GUI / 指令 P2:/rewind — 获取会话历史消息点(daemon 协议 list_history 已存在) if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); - const frame = await client.sendCommandAndWait("list_history", { session_id: sessionId, cwd: projectDir }, 5000); + const frame = await requireConn().sendCommandAndWait("list_history", { session_id: sessionId, cwd: projectDir }, 5000); return { messages: frame.messages || [] }; }); @@ -359,7 +362,7 @@ vision = false if (typeof recordIndex !== "number" || !Number.isInteger(recordIndex) || recordIndex < 0) { throw new Error("invalid record_index"); } - const frame = await client.sendCommandAndWait( + const frame = await requireConn().sendCommandAndWait( "rewind_session", { session_id: sessionId, cwd: projectDir, record_index: recordIndex }, 5000 @@ -374,7 +377,7 @@ vision = false if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); params.session_id = sessionId; } - const frame = await client.sendCommandAndWait("list_memories", params, 5000); + const frame = await requireConn().sendCommandAndWait("list_memories", params, 5000); return frame.memories || []; }); @@ -386,7 +389,7 @@ vision = false if (!validateSessionId(sessionId)) throw new Error("invalid session_id"); params.session_id = sessionId; } - const frame = await client.sendCommandAndWait("read_memory", params, 5000); + const frame = await requireConn().sendCommandAndWait("read_memory", params, 5000); return frame.memory || { id: memoryId, content: "" }; }); @@ -417,27 +420,27 @@ vision = false ipcMain.handle("emrg:listProjects", async () => { // GUI / 指令 P4:/rant 项目下拉 — daemon list_projects → projects_list - const frame = await client.sendCommandAndWait("list_projects", {}, 5000); + const frame = await requireConn().sendCommandAndWait("list_projects", {}, 5000); return frame.projects || []; }); ipcMain.handle("emrg:listTasks", async () => { // GUI / 指令 P4:/trigger — daemon list_tasks → tasks_list - const frame = await client.sendCommandAndWait("list_tasks", {}, 5000); + const frame = await requireConn().sendCommandAndWait("list_tasks", {}, 5000); return frame.tasks || []; }); ipcMain.handle("emrg:triggerTask", async (_e, { name }) => { // GUI / 指令 P4:/trigger — daemon trigger_task → trigger_result if (typeof name !== "string" || !name.trim()) throw new Error("invalid task name"); - const frame = await client.sendCommandAndWait("trigger_task", { name: name.trim() }, 5000); + const frame = await requireConn().sendCommandAndWait("trigger_task", { name: name.trim() }, 5000); return frame; }); ipcMain.handle("emrg:sendRant", async (_e, { message, project = "" } = {}) => { // GUI / 指令 P4:/rant — 提交反馈到演化系统(daemon rant 协议,字段序与 rants.jsonl 一致) if (typeof message !== "string" || !message.trim()) throw new Error("invalid rant message"); - const frame = await client.sendCommandAndWait("rant", { + const frame = await requireConn().sendCommandAndWait("rant", { message: message.trim().slice(0, 10000), project: String(project || "").trim(), // timestamp deliberately NOT sent: daemon stamps rants with local time @@ -449,13 +452,13 @@ vision = false ipcMain.handle("emrg:evolutionSummary", async (_e, { limit = 5 } = {}) => { // GUI / 指令 P3:自进化可见化 — daemon evolution_summary(count + 最近改进) - const frame = await client.sendCommandAndWait("evolution_summary", { limit }, 5000); + const frame = await requireConn().sendCommandAndWait("evolution_summary", { limit }, 5000); return { count: frame.count ?? 0, recent: frame.recent || [] }; }); ipcMain.handle("emrg:githubStatus", async () => { // Windows GCM rant Stage 2:设置页 GitHub 连接状态(daemon github_status) - const frame = await client.sendCommandAndWait("github_status", {}, 10000); + const frame = await requireConn().sendCommandAndWait("github_status", {}, 10000); return { authenticated: Boolean(frame.authenticated), user: frame.user || null }; }); @@ -463,7 +466,7 @@ vision = false // Auto update-check prompt (rant 2026-08-10T07:12:12): query daemon's // cached latest release; display-only, no auto download/install. try { - const frame = await client.sendCommandAndWait("update_check", {}, 10000); + const frame = await requireConn().sendCommandAndWait("update_check", {}, 10000); return { current_version: frame.current_version || "", latest_version: frame.latest_version || "", @@ -480,26 +483,26 @@ vision = false // Idempotency (rant 07:12:12 §4): record that the GUI showed the prompt // for this version — same version never re-prompted. try { - await client.sendCommandAndWait("update_check_prompted", { version: String(version || "") }, 5000); + await requireConn().sendCommandAndWait("update_check_prompted", { version: String(version || "") }, 5000); } catch { /* best-effort */ } return { ok: true }; }); ipcMain.handle("emrg:githubConnect", async (_e, { token }) => { // Windows GCM rant Stage 2:PAT 授权 + setup-git(daemon github_connect) - const frame = await client.sendCommandAndWait("github_connect", { token: String(token || "").trim() }, 40000); + const frame = await requireConn().sendCommandAndWait("github_connect", { token: String(token || "").trim() }, 40000); return { ok: Boolean(frame.ok), user: frame.user || null, error: frame.error || null }; }); ipcMain.handle("emrg:githubDisconnect", async () => { // Windows GCM rant Stage 2:断开 GitHub(daemon github_disconnect) - const frame = await client.sendCommandAndWait("github_disconnect", {}, 40000); + const frame = await requireConn().sendCommandAndWait("github_disconnect", {}, 40000); return { ok: Boolean(frame.ok), error: frame.error || null }; }); ipcMain.handle("emrg:githubConnectWeb", async () => { // Windows GCM rant Stage 2b:device flow 启动(daemon github_connect_web) - const frame = await client.sendCommandAndWait("github_connect_web", {}, 15000); + const frame = await requireConn().sendCommandAndWait("github_connect_web", {}, 15000); return { ok: Boolean(frame.ok), code: frame.code || null, url: frame.url || null, error: frame.error || null }; }); @@ -511,7 +514,7 @@ vision = false }); ipcMain.handle("emrg:setModel", async (_e, { model }) => { - await client.sendCommandAndWait("set_model", { model }, 5000); + await requireConn().sendCommandAndWait("set_model", { model }, 5000); return { ok: true }; }); @@ -525,7 +528,7 @@ vision = false }); ipcMain.handle("emrg:listModels", async () => { - const frame = await client.sendCommandAndWait("list_models", {}, 5000); + const frame = await requireConn().sendCommandAndWait("list_models", {}, 5000); return frame.models || []; }); @@ -553,7 +556,7 @@ vision = false ipcMain.handle("emrg:saveSettings", async (_e, rawConfig) => { const cfg = validateConfig(rawConfig || {}); - const wasRunning = await client?.isRunning() || false; // G123 + const wasRunning = Boolean(connManager?.daemonConn()?.connected); // G123:daemon 已连 = 运行中 let text; if (!fs.existsSync(configPath())) { text = await ensureConfigTemplate(); // G116 @@ -599,11 +602,12 @@ vision = false // G24:无参数 // G141:断连边界——ws 可能已 null/closed(_onClose 后 connected=false),sendCommand 抛异常 // 不能让它泄漏为 IPC reject → renderer unhandled rejection(对比 sendMessage 的 try-catch 防护) - if (client?.ws) { - try { await client.sendCommand("cancel"); } catch { /* 断连时忽略 */ } + // P2:cancel 发到当前激活连接(自有流所在连接),并释放其 G65 锁 + const c = activeConn(); + if (c?.ws) { + try { await c.sendCommand("cancel"); } catch { /* 断连时忽略 */ } } - ownStream = false; - ownStreamRequestId = null; + c?._releaseOwnStream(); return { ok: true }; }); @@ -614,77 +618,73 @@ vision = false }); } - // ── daemon 生命周期 ───────────────────────────────────── + // ── daemon 生命周期(P2:connManager 为 daemon 唯一 owner)────────── - async function ensureConnected() { - if (!client) { - // Phase 4(rant #12 §4 R7):打包模式传 app.isPackaged → daemon_client 走 - // 捆绑 emrgd 分支(_findDaemonExecutable)。 - client = new DaemonClient({ projectDir, logger, isPackaged: app.isPackaged }); - // G122:message_delta 16ms 批量推送 - let deltaBuf = []; - let deltaTimer = null; - // rant 14:11:冲刷 delta 缓冲——终态事件(done/error/cancelled)直通不走缓冲, - // 若残留 delta 在 16ms 定时器之后才 flush,会晚于终态到达渲染层 → - // handleDelta 找不到 group 节点 → 建孤儿节点(误标"来自其他客户端")+ 光标永不消失。 - const flushDeltaBuf = () => { - if (deltaTimer) { - clearTimeout(deltaTimer); - deltaTimer = null; - } - if (deltaBuf.length && win && !win.isDestroyed()) { - const chunks = deltaBuf; - deltaBuf = []; - win.webContents.send("emrg:event", { type: "message_delta", data: { chunks } }); - } - }; - client.onEvent((type, data) => { - if (type === "message_delta") { - deltaBuf.push(data); - if (!deltaTimer) { - deltaTimer = setTimeout(flushDeltaBuf, 16); - } - return; - } - if (type === "done" || type === "error" || type === "cancelled") { - flushDeltaBuf(); // 终态前先清空缓冲:delta 保证不晚于终态(webContents.send 保序) - } - if (type === "done") { - // 仅自有流的 done 释放 G65 锁(广播 done 不影响);timeout 兜底同样只清自有 - if (data.request_id === ownStreamRequestId || (data.timeout && ownStream)) { - ownStream = false; - ownStreamRequestId = null; - } - } - if (type === "error") { - // session busy 是即发错误(daemon 返回后无 done 跟随)——释放 ownStream,防 G65 锁泄漏 - // (流式错误如 LLM error 则有 done 跟随,由 done 分支释放,不在此处理) - if (data.error && String(data.error).includes("session busy")) { - ownStream = false; - ownStreamRequestId = null; - } - } + // 惰性初始化 connManager(挂事件桥 + 恢复钩子)。 + function ensureConnManager() { + if (connManager) return connManager; + connManager = new ConnManager({ projectDir, logger, isPackaged: app.isPackaged }); + // 每个新会话连接建立时挂 renderer 事件桥(附带 sid;含 recoverAll 重开路径) + connManager.onOpen((sid, conn) => { + conn.onEvent((type, data) => { if (win && !win.isDestroyed()) { - win.webContents.send("emrg:event", { type, data }); - } - }); - client.onEvent((type) => { - if (type === "disconnected") { - ownStream = false; - ownStreamRequestId = null; - scheduleReconnect(); + // 主动关闭(切走/删除)不触发断线横幅——真断连/daemon 重启照常转发 + if (type === "disconnected" && conn._intentionalClose) return; + win.webContents.send("emrg:event", { type, data, sid }); } }); + }); + // daemon 重启恢复完成后刷新 UI 状态(对齐旧 G41 重连成功块) + connManager.onRecovered(async () => { + try { + const sessions = await listSessions(); + sendToRenderer("sessions", { sessions }); + const pong = await waitForPong(); + sendToRenderer("status", { connected: true, server_id: pong?.identity?.instance_id, model: pong?.model }); + logger.info("[gui] connManager recovery complete"); + } catch (e) { + logger.warn(`[gui] post-recovery refresh failed: ${e.message}`); + } + }); + return connManager; + } + + // 当前激活连接:有会话连接用会话连接(同一 daemon,命令通用);否则 daemon 级连接。 + function activeConn() { + if (!connManager) return null; + if (currentSessionId) { + const c = connManager.get(currentSessionId); + if (c) return c; } + return connManager.daemonConn(); + } + + // 同步取可用连接(未连 → 抛错,与旧 `if (!client?.connected) throw` 语义一致)。 + function requireConn() { + const c = activeConn(); + if (!c || !c.connected) throw new Error("daemon not connected"); + return c; + } + + // 打开(或复用)会话连接;事件桥由 onOpen 钩子统一挂(含 recoverAll 重开)。 + async function openSession(sid, projectPath, { resume = true } = {}) { + const existing = connManager.get(sid); + if (existing && existing.connected) return existing; + return connManager.open(sid, projectPath, { resume }); + } + + async function ensureConnected() { + ensureConnManager(); try { - await client.ensureConnected(); + await connManager.ensureDaemon(); logger.info("[gui] connected to emrgd"); cancelReconnect(); reconnectDelayMs = 1000; // 退避复位 daemonStoppedNotified = false; // 节流提示复位(下个生命周期可再提示) sendToRenderer("status", { connected: true }); } catch (e) { - if (client._authFailed) { + const dm = connManager.daemonConn(); + if (dm?._authFailed) { // G88:认证失败 → 停止自动重试 sendToRenderer("status", { connected: false, auth_failed: true, error: e.message }); return; @@ -701,6 +701,8 @@ vision = false } } + // daemon 级重连退避(connManager 重启恢复覆盖会话连接;此处覆盖"无会话连接 + // 时 daemon 连接不可用"的初始/空闲场景)。 function scheduleReconnect() { if (stopping || reconnectTimer) return; const delay = reconnectDelayMs; @@ -709,19 +711,17 @@ vision = false reconnectTimer = null; sendToRenderer("status", { connected: false, reconnecting: true }); await ensureConnected(); - if (client?.connected) { - // G41:重连成功 → list_sessions + 重新 resume 当前会话 + if (connManager.daemonConn()?.connected) { + // G41(P2 改写):恢复当前会话连接(若 daemon 重启后未由 recoverAll 重开) + if (currentSessionId && !connManager.get(currentSessionId)) { + try { await openSession(currentSessionId, projectDir); } catch { /* 会话可能已删 */ } + } const sessions = await listSessions(); sendToRenderer("sessions", { sessions }); - if (currentSessionId) { - try { - await client.sendCommandAndWait("resume_session", { session_id: currentSessionId, cwd: projectDir }, 5000); - } catch { /* 会话可能已删 */ } - } const pong = await waitForPong(); sendToRenderer("status", { connected: true, server_id: pong?.identity?.instance_id, model: pong?.model }); } - }, 1000); + }, delay); } function cancelReconnect() { @@ -729,8 +729,10 @@ vision = false } async function waitForPong(timeoutMs = 3000) { + const conn = activeConn(); + if (!conn || !conn.connected) return null; // 未连接 → 直接超时语义(不抛) return new Promise((resolve) => { - const off = client.onEvent((type, data) => { + const off = conn.onEvent((type, data) => { if (type === "pong") { off(); clearTimeout(timer); @@ -738,13 +740,13 @@ vision = false } }); const timer = setTimeout(() => { off(); resolve(null); }, timeoutMs); - client.sendCommand("ping"); + conn.sendCommand("ping"); }); } async function listSessions() { try { - const frame = await client.sendCommandAndWait("list_sessions", { cwd: projectDir }, 5000); + const frame = await requireConn().sendCommandAndWait("list_sessions", { cwd: projectDir }, 5000); return frame.sessions || []; } catch (e) { logger.warn(`[gui] list_sessions failed: ${e.message}`); @@ -812,7 +814,7 @@ vision = false app.on("window-all-closed", () => { stopping = true; cancelReconnect(); - if (client) client.close(); + connManager?.closeAll(); // P2:关闭全部会话连接 + daemon 级连接 if (process.platform !== "darwin") app.quit(); }); diff --git a/emrg/gui/test/conn-manager.test.js b/emrg/gui/test/conn-manager.test.js index bba0d15..b008cb6 100644 --- a/emrg/gui/test/conn-manager.test.js +++ b/emrg/gui/test/conn-manager.test.js @@ -168,6 +168,87 @@ test("P2 close: 断开并移除;再 get → null;重复 close → false", as assert.strictEqual(manager.close("sess-1"), false, "second close returns false"); }); +test("P2 close: 主动关闭标记 _intentionalClose 且不触发重启恢复", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const { conn } = await driveOpen(manager, "sess-1", "/proj/a"); + let recoverCalls = 0; + manager.recoverAll = async () => { recoverCalls += 1; }; + assert.strictEqual(manager.close("sess-1"), true); + assert.strictEqual(conn._intentionalClose, true, "close must mark intentional close"); + assert.strictEqual(manager.get("sess-1"), null); + await new Promise((r) => setTimeout(r, 30)); + assert.strictEqual(recoverCalls, 0, "intentional close must not trigger restart recovery"); +}); + +test("P2 open: 断连残留连接 → 关闭重开(不返回 stale conn)", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const s1 = await driveOpen(manager, "sess-1", "/proj/a"); + // 停用自动恢复(本测试专注 open 的 stale 处理;自动恢复由前序测试覆盖)—— + // 否则单会话全断会被判定为 daemon 重启自动 recoverAll + const origRecover = manager.recoverAll.bind(manager); + manager.recoverAll = async () => {}; + // 模拟真断连(非主动关闭):ws close → conn.connected=false + s1.sessionWs.emit("close"); + await new Promise((r) => setTimeout(r, 30)); + assert.strictEqual(manager.get("sess-1").connected, false, "conn dropped"); + + // 再次 open → 必须重开新连接(stale conn 不可复用) + const p = manager.open("sess-1", "/proj/a"); + const deadline = Date.now() + 2000; + while (Date.now() < deadline && currentMockWs === s1.sessionWs) { + await new Promise((r) => setTimeout(r, 5)); + } + const newWs = await waitForWs(() => currentMockWs !== s1.sessionWs); + await driveAuth(newWs); + const d2 = Date.now() + 2000; + while (newWs.sent.length < 2 && Date.now() < d2) await new Promise((r) => setTimeout(r, 5)); + const f = JSON.parse(newWs.sent.at(-1)); + if (f.type === "resume_session") { + newWs.emit("message", Buffer.from(JSON.stringify({ type: "resume_result", session_id: f.session_id }))); + } + const conn = await p; + manager.recoverAll = origRecover; + assert.strictEqual(conn.connected, true, "reopened conn connected"); + assert.notStrictEqual(conn, s1.conn, "must be a fresh conn, not the stale one"); + assert.strictEqual(manager.get("sess-1"), conn); +}); + +test("P2 单连接退避: 多会话中单条断 → 独立退避重开该会话(不触发全量恢复)", async () => { + const manager = new ConnManager({ projectDir: tmpHome, singleRetryDelayMs: 20 }); + const s1 = await driveOpen(manager, "sess-1", "/proj/a"); + const s2 = await driveOpen(manager, "sess-2", "/proj/b"); + let recoverCalls = 0; + const origRecover = manager.recoverAll.bind(manager); + manager.recoverAll = async () => { recoverCalls += 1; }; + + // 只有 s1 断(s2 仍连)→ 非重启 → 不 recoverAll,走单连接退避 + s1.sessionWs.emit("close"); + await new Promise((r) => setTimeout(r, 10)); + assert.strictEqual(recoverCalls, 0, "single drop must not trigger recoverAll"); + + // 退避到期 → 重开 s1(stale 关闭重开 + resume 重订阅) + const deadline = Date.now() + 4000; + let newWs = null; + while (Date.now() < deadline && !newWs) { + if (currentMockWs && currentMockWs !== s1.sessionWs && currentMockWs !== s2.sessionWs) newWs = currentMockWs; + await new Promise((r) => setTimeout(r, 5)); + } + assert.ok(newWs, "retry must open a new ws for the dropped session"); + await driveAuth(newWs); + const d2 = Date.now() + 2000; + while (newWs.sent.length < 2 && Date.now() < d2) await new Promise((r) => setTimeout(r, 5)); + const f = JSON.parse(newWs.sent.at(-1)); + if (f.type === "resume_session") { + newWs.emit("message", Buffer.from(JSON.stringify({ type: "resume_result", session_id: f.session_id }))); + } + // 等重开完成 + const d3 = Date.now() + 4000; + while (Date.now() < d3 && !manager.get("sess-1")?.connected) await new Promise((r) => setTimeout(r, 5)); + assert.strictEqual(manager.get("sess-1").connected, true, "session 1 reconnected via single retry"); + assert.strictEqual(manager.get("sess-2").connected, true, "session 2 untouched"); + manager.recoverAll = origRecover; +}); + test("P2 closeAll: 全部关闭", async () => { const manager = new ConnManager({ projectDir: tmpHome }); await driveOpen(manager, "sess-1", "/proj/a"); @@ -256,3 +337,164 @@ test("P2 重启恢复: recoverAll 重连重订阅全部会话(复用 open 序 assert.strictEqual(manager.get(sid).connected, true, `session ${sid} reconnected`); } }); + +// ── P2 rewire(rant 15:07:19):daemon 级连接 + onOpen 钩子 + resume:false ── + +test("P2 ensureDaemon: 保留 daemon 级连接(复用,不重复建连)", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const p = manager.ensureDaemon(); + const bootWs = await waitForWs(); + await driveAuth(bootWs); + const dm = await p; + assert.ok(dm instanceof DaemonClient, "ensureDaemon returns a DaemonClient"); + assert.strictEqual(manager.daemonConn(), dm, "daemon conn accessible via daemonConn()"); + assert.strictEqual(dm.connected, true); + // 再次 ensureDaemon → 同一实例(不新开连接) + const dm2 = await manager.ensureDaemon(); + assert.strictEqual(dm2, dm, "second ensureDaemon reuses the daemon conn"); +}); + +test("P2 open: 复用 ensureDaemon 连接(无重复引导)+ resume_session 照常", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const p0 = manager.ensureDaemon(); + const bootWs = await waitForWs(); + await driveAuth(bootWs); + const dm = await p0; + assert.strictEqual(manager.daemonConn(), dm); + const p = manager.open("sess-1", "/proj/a"); + let sessionWs = await waitForWs(() => currentMockWs !== bootWs); + await driveAuth(sessionWs); + const deadline = Date.now() + 2000; + while (sessionWs.sent.length < 2 && Date.now() < deadline) { + await new Promise((r) => setTimeout(r, 5)); + } + const resumeFrame = JSON.parse(sessionWs.sent.at(-1)); + assert.strictEqual(resumeFrame.type, "resume_session"); + sessionWs.emit("message", Buffer.from(JSON.stringify({ type: "resume_result", session_id: "sess-1" }))); + const conn = await p; + assert.strictEqual(manager.get("sess-1"), conn); + // 仅一个引导连接(daemon 级连接未被重复创建) + assert.strictEqual(currentMockWs, sessionWs); +}); + +test("P2 open({resume:false}): 新会话首条消息前不 resume_session(daemon 隐式订阅)", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const p = manager.open("sess-new", "/proj/a", { resume: false }); + const bootWs = await waitForWs(); + await driveAuth(bootWs); + await new Promise((r) => setTimeout(r, 10)); + let sessionWs = await waitForWs(() => currentMockWs !== bootWs); + await driveAuth(sessionWs); + // 等一小段时间确认无 resume_session 发出(只有 auth 一条) + await new Promise((r) => setTimeout(r, 50)); + const frames = sessionWs.sent.map((s) => JSON.parse(s)); + assert.ok(!frames.some((f) => f.type === "resume_session"), "resume:false must not resume"); + const conn = await p; + assert.strictEqual(conn.connected, true); +}); + +test("P2 open 失败: resume 报错(会话已删)→ 连接关闭不泄漏 + 抛出", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const p = manager.open("sess-gone", "/proj/a"); + const bootWs = await waitForWs(); + await driveAuth(bootWs); + await new Promise((r) => setTimeout(r, 10)); + let sessionWs = await waitForWs(() => currentMockWs !== bootWs); + await driveAuth(sessionWs); + const deadline = Date.now() + 2000; + while (sessionWs.sent.length < 2 && Date.now() < deadline) { + await new Promise((r) => setTimeout(r, 5)); + } + // 监视 DaemonClient.close —— open 失败必须关闭半开连接(防泄漏) + const closeCalls = []; + const origClose = DaemonClient.prototype.close; + DaemonClient.prototype.close = function () { closeCalls.push(this); return origClose.apply(this, arguments); }; + try { + sessionWs.emit("message", Buffer.from(JSON.stringify({ type: "resume_result", error: "session not found" }))); + await assert.rejects(() => p, /session not found/); + } finally { + DaemonClient.prototype.close = origClose; + } + assert.strictEqual(closeCalls.length, 1, "failed open must close the half-open conn"); + assert.strictEqual(manager.get("sess-gone"), null, "failed open must not be registered"); +}); + +test("P2 onOpen 钩子: 新会话连接建立时触发(含 recoverAll 重开路径)", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + const opened = []; + manager.onOpen((sid, conn) => opened.push([sid, conn])); + const s1 = await driveOpen(manager, "sess-1", "/proj/a"); + const s2 = await driveOpen(manager, "sess-2", "/proj/b"); + assert.deepStrictEqual(opened.map(([sid]) => sid), ["sess-1", "sess-2"]); + + // recoverAll 重开 → 钩子再次触发(重开路径也要挂事件桥) + const origRecover = manager.recoverAll.bind(manager); + manager.recoverAll = async () => {}; + s1.sessionWs.emit("close"); + s2.sessionWs.emit("close"); + await new Promise((r) => setTimeout(r, 30)); + manager.recoverAll = origRecover; + + let reopened = 0; + const processed = new Set(); + const stale = new Set([s1.sessionWs, s2.sessionWs]); + const deadline = Date.now() + 4000; + const driver = (async () => { + while (reopened < 2 && Date.now() < deadline) { + const ws = await waitForWs(() => !processed.has(currentMockWs) && !stale.has(currentMockWs)); + processed.add(ws); + await driveAuth(ws); + const d2 = Date.now() + 300; + while (ws.sent.length < 2 && Date.now() < d2) await new Promise((r) => setTimeout(r, 5)); + if (ws.sent.length >= 2) { + const f = JSON.parse(ws.sent.at(-1)); + if (f.type === "resume_session") { + ws.emit("message", Buffer.from(JSON.stringify({ type: "resume_result", session_id: f.session_id }))); + reopened += 1; + } + } + } + })(); + await manager.recoverAll(); + await driver; + assert.strictEqual(reopened, 2, "both sessions reopened"); + assert.deepStrictEqual(opened.map(([sid]) => sid), ["sess-1", "sess-2", "sess-1", "sess-2"], + "onOpen fires again for recovered conns"); +}); + +test("P2 onRecovered 钩子: recoverAll 完成后触发(main.js 刷新 UI 状态用)", async () => { + const manager = new ConnManager({ projectDir: tmpHome }); + let recovered = 0; + manager.onRecovered(() => { recovered += 1; }); + const s1 = await driveOpen(manager, "sess-1", "/proj/a"); + const origRecover = manager.recoverAll.bind(manager); + manager.recoverAll = async () => {}; + s1.sessionWs.emit("close"); + await new Promise((r) => setTimeout(r, 30)); + manager.recoverAll = origRecover; + + let reopened = 0; + const processed = new Set(); + const stale = new Set([s1.sessionWs]); + const deadline = Date.now() + 4000; + const driver = (async () => { + while (reopened < 1 && Date.now() < deadline) { + const ws = await waitForWs(() => !processed.has(currentMockWs) && !stale.has(currentMockWs)); + processed.add(ws); + await driveAuth(ws); + const d2 = Date.now() + 300; + while (ws.sent.length < 2 && Date.now() < d2) await new Promise((r) => setTimeout(r, 5)); + if (ws.sent.length >= 2) { + const f = JSON.parse(ws.sent.at(-1)); + if (f.type === "resume_session") { + ws.emit("message", Buffer.from(JSON.stringify({ type: "resume_result", session_id: f.session_id }))); + reopened += 1; + } + } + } + })(); + await manager.recoverAll(); + await driver; + assert.strictEqual(reopened, 1); + assert.strictEqual(recovered, 1, "onRecovered must fire after recoverAll completes"); +}); diff --git a/emrg/gui/test/daemon_client.test.js b/emrg/gui/test/daemon_client.test.js index a26c567..6c87872 100644 --- a/emrg/gui/test/daemon_client.test.js +++ b/emrg/gui/test/daemon_client.test.js @@ -826,3 +826,64 @@ test("18:47:37: stale projectDir port + spawn 节流失败 → probe 复用 cano assert.match(probeLine, /failed\(daemon failed to start after 3 attempts/); assert.ok(logs.some((l) => l.includes("existing daemon detected at port=41234, reusing")), "复用日志"); }); + +// ── P2 自有流锁(G65 每连接独立;rant 15:07:19)────────────────────────── + +test("P2 ownStream: sendTask(stream:true) 标记 ownStream + requestId", async () => { + const client = new DaemonClient({ projectDir: tmpHome }); + await connectClient(client); + const rid = client.sendTask({ sessionId: "s_260803_1730_abcd1234", cwd: "/proj", prompt: "hi", stream: true, requestId: "req-own-1" }); + assert.strictEqual(client.ownStream, true, "stream:true must set ownStream"); + assert.strictEqual(client.ownStreamRequestId, "req-own-1"); + assert.strictEqual(rid, "req-own-1"); +}); + +test("P2 ownStream: 自有 done(request 匹配)→ 释放锁;广播 done(不匹配)→ 保持", async () => { + const client = new DaemonClient({ projectDir: tmpHome }); + await connectClient(client); + client.sendTask({ sessionId: "s_260803_1730_abcd1234", cwd: "/proj", prompt: "hi", stream: true, requestId: "req-own-2" }); + const send = (obj) => currentMockWs.emit("message", Buffer.from(JSON.stringify(obj))); + // 广播 done(其他客户端/其他流)→ 锁保持 + send({ request_id: "req-other", done: true, delta: false }); + assert.strictEqual(client.ownStream, true, "broadcast done must not release own lock"); + // 自有 done → 释放 + send({ request_id: "req-own-2", done: true, delta: false }); + assert.strictEqual(client.ownStream, false, "own done must release lock"); + assert.strictEqual(client.ownStreamRequestId, null); +}); + +test("P2 ownStream: timeout 兜底 done(无匹配 request)→ 释放锁", async () => { + const client = new DaemonClient({ projectDir: tmpHome }); + await connectClient(client); + client.sendTask({ sessionId: "s_260803_1730_abcd1234", cwd: "/proj", prompt: "hi", stream: true, requestId: "req-own-3" }); + const send = (obj) => currentMockWs.emit("message", Buffer.from(JSON.stringify(obj))); + send({ request_id: "req-stale", done: true, delta: false, timeout: true }); + assert.strictEqual(client.ownStream, false, "timeout done must release own lock"); +}); + +test("P2 ownStream: session busy 即发 error → 释放锁(防 G65 锁泄漏)", async () => { + const client = new DaemonClient({ projectDir: tmpHome }); + await connectClient(client); + client.sendTask({ sessionId: "s_260803_1730_abcd1234", cwd: "/proj", prompt: "hi", stream: true, requestId: "req-own-4" }); + const send = (obj) => currentMockWs.emit("message", Buffer.from(JSON.stringify(obj))); + send({ error: "session busy: another stream running" }); + assert.strictEqual(client.ownStream, false, "session busy error must release lock"); +}); + +test("P2 ownStream: cancelled(request 匹配)→ 释放锁", async () => { + const client = new DaemonClient({ projectDir: tmpHome }); + await connectClient(client); + client.sendTask({ sessionId: "s_260803_1730_abcd1234", cwd: "/proj", prompt: "hi", stream: true, requestId: "req-own-5" }); + const send = (obj) => currentMockWs.emit("message", Buffer.from(JSON.stringify(obj))); + send({ type: "cancelled", request_id: "req-own-5" }); + assert.strictEqual(client.ownStream, false, "own cancelled must release lock"); +}); + +test("P2 ownStream: 断连 → 释放锁", async () => { + const client = new DaemonClient({ projectDir: tmpHome }); + await connectClient(client); + client.sendTask({ sessionId: "s_260803_1730_abcd1234", cwd: "/proj", prompt: "hi", stream: true, requestId: "req-own-6" }); + assert.strictEqual(client.ownStream, true); + client.close(); + assert.strictEqual(client.ownStream, false, "disconnect must release lock"); +});