From 77d73d0b84cad0b0379c9d04e2f333cb40dd7741 Mon Sep 17 00:00:00 2001 From: Mr-Maidong <15557111679@163.com> Date: Wed, 6 May 2026 11:25:10 +0800 Subject: [PATCH 01/10] feat(base): persist channel sessions across restarts Session routing was lost on channel restart because clearAll() deleted sessions.json and restoreSessions() was never called at startup. Additionally, AcpBridge.loadSession() always returned undefined due to reading a non-existent field on LoadSessionResponse. Co-authored-by: Qwen-Coder --- packages/channels/base/src/AcpBridge.ts | 4 ++-- packages/channels/base/src/SessionRouter.ts | 12 +++--------- packages/cli/src/commands/channel/start.ts | 18 ++++++++++++++++++ 3 files changed, 23 insertions(+), 11 deletions(-) diff --git a/packages/channels/base/src/AcpBridge.ts b/packages/channels/base/src/AcpBridge.ts index 9f5638b4506..afcfe11333b 100644 --- a/packages/channels/base/src/AcpBridge.ts +++ b/packages/channels/base/src/AcpBridge.ts @@ -132,12 +132,12 @@ export class AcpBridge extends EventEmitter { async loadSession(sessionId: string, cwd: string): Promise { const conn = this.ensureConnection(); - const response = await conn.loadSession({ + await conn.loadSession({ sessionId, cwd, mcpServers: [], }); - return response.sessionId; + return sessionId; } async prompt( diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index bf08301bc69..944581884fc 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -1,4 +1,4 @@ -import { existsSync, readFileSync, writeFileSync, unlinkSync } from 'node:fs'; +import { existsSync, readFileSync, writeFileSync } from 'node:fs'; import type { SessionScope, SessionTarget } from './types.js'; import type { AcpBridge } from './AcpBridge.js'; @@ -196,18 +196,12 @@ export class SessionRouter { return { restored, failed }; } - /** Clear in-memory state and delete persist file. Used on clean shutdown. */ + /** Clear in-memory state on shutdown. Persist file is kept so sessions + * can be restored on next start via {@link restoreSessions}. */ clearAll(): void { this.toSession.clear(); this.toTarget.clear(); this.toCwd.clear(); - if (this.persistPath && existsSync(this.persistPath)) { - try { - unlinkSync(this.persistPath); - } catch { - // best-effort - } - } } private persist(): void { diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index 8e11c31933a..734f203f282 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -215,6 +215,15 @@ async function startSingle(name: string, proxy?: string): Promise { channels.set(name, channel); registerToolCallDispatch(bridge, router, channels); + // Restore sessions from previous run + const restoreResult = await router.restoreSessions(); + if (restoreResult.restored > 0) { + writeStdoutLine( + `[Channel] Sessions restored: ${restoreResult.restored}` + + (restoreResult.failed > 0 ? `, failed: ${restoreResult.failed}` : ''), + ); + } + try { await channel.connect(); } catch (err) { @@ -365,6 +374,15 @@ async function startAll(proxy?: string): Promise { } registerToolCallDispatch(bridge, router, channels); + // Restore sessions from previous run + const restoreResult = await router.restoreSessions(); + if (restoreResult.restored > 0) { + writeStdoutLine( + `[Channel] Sessions restored: ${restoreResult.restored}` + + (restoreResult.failed > 0 ? `, failed: ${restoreResult.failed}` : ''), + ); + } + // Connect all channels let connectedCount = 0; for (const [name, channel] of channels) { From c4cf65e1e82dcc2f7161e62b68d43fdfe7f2a881 Mon Sep 17 00:00:00 2001 From: Mr-Maidong <15557111679@163.com> Date: Wed, 6 May 2026 15:19:33 +0800 Subject: [PATCH 02/10] fix(base): address review feedback on session persistence PR - Add restoreSessions() tests: persist file restoration, stale entry skipping, missing file, and empty file cases - Fix log consistency: log restore results when failed > 0 even if restored is 0 (both startSingle and startAll) - Throw in acpAgent.loadSession when session does not exist on disk, preventing stale session accumulation in the persist file - Fix clearAll() JSDoc wording Co-authored-by: Qwen-Coder --- .../channels/base/src/SessionRouter.test.ts | 134 ++++++++++++++++++ packages/channels/base/src/SessionRouter.ts | 4 +- packages/cli/src/acp-integration/acpAgent.ts | 6 + packages/cli/src/commands/channel/start.ts | 4 +- 4 files changed, 144 insertions(+), 4 deletions(-) diff --git a/packages/channels/base/src/SessionRouter.test.ts b/packages/channels/base/src/SessionRouter.test.ts index d33f67fd77a..a32e76855d2 100644 --- a/packages/channels/base/src/SessionRouter.test.ts +++ b/packages/channels/base/src/SessionRouter.test.ts @@ -229,4 +229,138 @@ describe('SessionRouter', () => { expect(bridge.newSession).not.toHaveBeenCalled(); }); }); + + describe('restoreSessions', () => { + it('restores sessions from persist file', async () => { + const { mkdirSync, writeFileSync, rmSync } = await import('node:fs'); + const { join } = await import('node:path'); + const tmpDir = join('/tmp', `test-restore-${Date.now()}`); + mkdirSync(tmpDir, { recursive: true }); + + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'old-session-1', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', + }, + 'telegram:bob:chat2': { + sessionId: 'old-session-2', + target: { + channelName: 'telegram', + senderId: 'bob', + chatId: 'chat2', + }, + cwd: '/workspace', + }, + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); + + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const result = await router.restoreSessions(); + + expect(result.restored).toBe(2); + expect(result.failed).toBe(0); + expect(bridge.loadSession).toHaveBeenCalledTimes(2); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); + expect(router.hasSession('telegram', 'bob', 'chat2')).toBe(true); + + const target = router.getTarget('old-session-1'); + expect(target).toEqual({ + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }); + + rmSync(tmpDir, { recursive: true }); + }); + + it('skips entries whose loadSession throws', async () => { + const { mkdirSync, writeFileSync, rmSync } = await import('node:fs'); + const { join } = await import('node:path'); + const tmpDir = join('/tmp', `test-restore-fail-${Date.now()}`); + mkdirSync(tmpDir, { recursive: true }); + + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'good-session', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', + }, + 'telegram:bob:chat2': { + sessionId: 'stale-session', + target: { + channelName: 'telegram', + senderId: 'bob', + chatId: 'chat2', + }, + cwd: '/workspace', + }, + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); + + // Make loadSession throw for the stale session + (bridge.loadSession as ReturnType).mockImplementation( + (id: string) => { + if (id === 'stale-session') { + throw new Error('Session not found'); + } + return id; + }, + ); + + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const result = await router.restoreSessions(); + + expect(result.restored).toBe(1); + expect(result.failed).toBe(1); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); + expect(router.hasSession('telegram', 'bob', 'chat2')).toBe(false); + + // Persist file should be updated to remove the failed entry + const { readFileSync } = await import('node:fs'); + const updated = JSON.parse(readFileSync(persistFile, 'utf-8')); + expect(Object.keys(updated)).toEqual(['telegram:alice:chat1']); + + rmSync(tmpDir, { recursive: true }); + }); + + it('returns zeros when no persist file exists', async () => { + const router = new SessionRouter( + bridge, + '/tmp', + 'user', + '/nonexistent/sessions.json', + ); + const result = await router.restoreSessions(); + expect(result).toEqual({ restored: 0, failed: 0 }); + }); + + it('returns zeros when persist file is empty', async () => { + const { mkdirSync, writeFileSync, rmSync } = await import('node:fs'); + const { join } = await import('node:path'); + const tmpDir = join('/tmp', `test-restore-empty-${Date.now()}`); + mkdirSync(tmpDir, { recursive: true }); + + const persistFile = join(tmpDir, 'sessions.json'); + writeFileSync(persistFile, '{}'); + + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const result = await router.restoreSessions(); + + expect(result).toEqual({ restored: 0, failed: 0 }); + expect(bridge.loadSession).not.toHaveBeenCalled(); + + rmSync(tmpDir, { recursive: true }); + }); + }); }); diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index 944581884fc..29a6b9da388 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -196,8 +196,8 @@ export class SessionRouter { return { restored, failed }; } - /** Clear in-memory state on shutdown. Persist file is kept so sessions - * can be restored on next start via {@link restoreSessions}. */ + /** Clear in-memory state. Persist file is left intact for the next start + * via {@link restoreSessions}. */ clearAll(): void { this.toSession.clear(); this.toTarget.clear(); diff --git a/packages/cli/src/acp-integration/acpAgent.ts b/packages/cli/src/acp-integration/acpAgent.ts index 41905996f09..effeadb8eed 100644 --- a/packages/cli/src/acp-integration/acpAgent.ts +++ b/packages/cli/src/acp-integration/acpAgent.ts @@ -293,6 +293,12 @@ class QwenAgent implements Agent { }, ); + if (!exists) { + throw new Error( + `Session ${params.sessionId} does not exist at ${params.cwd}`, + ); + } + const config = await this.newSessionConfig( params.cwd, params.mcpServers, diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index 734f203f282..bef953a2c38 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -217,7 +217,7 @@ async function startSingle(name: string, proxy?: string): Promise { // Restore sessions from previous run const restoreResult = await router.restoreSessions(); - if (restoreResult.restored > 0) { + if (restoreResult.restored > 0 || restoreResult.failed > 0) { writeStdoutLine( `[Channel] Sessions restored: ${restoreResult.restored}` + (restoreResult.failed > 0 ? `, failed: ${restoreResult.failed}` : ''), @@ -376,7 +376,7 @@ async function startAll(proxy?: string): Promise { // Restore sessions from previous run const restoreResult = await router.restoreSessions(); - if (restoreResult.restored > 0) { + if (restoreResult.restored > 0 || restoreResult.failed > 0) { writeStdoutLine( `[Channel] Sessions restored: ${restoreResult.restored}` + (restoreResult.failed > 0 ? `, failed: ${restoreResult.failed}` : ''), From 9573868a4acd1851b3b95a452c54c2656a02932f Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Thu, 7 May 2026 23:45:39 +0800 Subject: [PATCH 03/10] test: add regression tests for clearAll persist and loadSession rejection Co-authored-by: Qwen-Coder --- .../channels/base/src/SessionRouter.test.ts | 61 +++++++++++++++++++ .../cli/src/acp-integration/acpAgent.test.ts | 51 +++++++++++++++- 2 files changed, 111 insertions(+), 1 deletion(-) diff --git a/packages/channels/base/src/SessionRouter.test.ts b/packages/channels/base/src/SessionRouter.test.ts index a32e76855d2..71a37ec91a1 100644 --- a/packages/channels/base/src/SessionRouter.test.ts +++ b/packages/channels/base/src/SessionRouter.test.ts @@ -217,6 +217,67 @@ describe('SessionRouter', () => { expect(router.hasSession('ch', 'alice', 'chat1')).toBe(false); expect(router.getAll()).toEqual([]); }); + + it('preserves persist file for restoration on next start', async () => { + const { mkdirSync, writeFileSync, rmSync, readFileSync, existsSync } = + await import('node:fs'); + const { join } = await import('node:path'); + const tmpDir = join('/tmp', `test-cleara11-persist-${Date.now()}`); + mkdirSync(tmpDir, { recursive: true }); + + const persistFile = join(tmpDir, 'sessions.json'); + const originalEntries = { + 'telegram:alice:chat1': { + sessionId: 'session-xyz', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', + }, + }; + + try { + // Create a router, write a session, then clearAll + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + await router.resolve('telegram', 'alice', 'chat1'); + + // Populate the persist file with known data + writeFileSync(persistFile, JSON.stringify(originalEntries, null, 2)); + + router.clearAll(); + + // Memory should be cleared + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(false); + expect(router.getAll()).toEqual([]); + + // Persist file should still exist with original content + expect(existsSync(persistFile)).toBe(true); + const onDisk = JSON.parse(readFileSync(persistFile, 'utf-8')); + expect(onDisk).toEqual(originalEntries); + + // A fresh router should be able to restore from it + const freshRouter = new SessionRouter( + bridge, + '/tmp', + 'user', + persistFile, + ); + const result = await freshRouter.restoreSessions(); + + expect(result.restored).toBe(1); + expect(result.failed).toBe(0); + expect(freshRouter.hasSession('telegram', 'alice', 'chat1')).toBe(true); + expect(freshRouter.getTarget('session-xyz')).toEqual({ + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }); + } finally { + rmSync(tmpDir, { recursive: true }); + } + }); }); describe('setBridge', () => { diff --git a/packages/cli/src/acp-integration/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index a5d92a30311..b148cdea545 100644 --- a/packages/cli/src/acp-integration/acpAgent.test.ts +++ b/packages/cli/src/acp-integration/acpAgent.test.ts @@ -92,6 +92,11 @@ vi.mock('@qwen-code/qwen-code-core', () => ({ _args: args, })), SessionService: vi.fn(), + Storage: { + runWithRuntimeBaseDir: vi.fn( + (_dir: unknown, _cwd: unknown, fn: () => unknown) => fn(), + ), + }, tokenLimit: vi.fn(), SessionStartSource: { Startup: 'startup', @@ -126,7 +131,11 @@ import { import type { Config } from '@qwen-code/qwen-code-core'; import type { LoadedSettings } from '../config/settings.js'; import type { CliArgs } from '../config/config.js'; -import { SessionEndReason, MCPServerConfig } from '@qwen-code/qwen-code-core'; +import { + SessionEndReason, + MCPServerConfig, + SessionService, +} from '@qwen-code/qwen-code-core'; import type { McpServer } from '@agentclientprotocol/sdk'; import { AgentSideConnection } from '@agentclientprotocol/sdk'; import { loadSettings } from '../config/settings.js'; @@ -631,6 +640,7 @@ describe('QwenAgent MCP SSE/HTTP support', () => { type AgentLike = { initialize: (args: Record) => Promise; newSession: (args: Record) => Promise; + loadSession: (args: Record) => Promise; }; let mockConfig: Config; @@ -893,4 +903,43 @@ describe('QwenAgent MCP SSE/HTTP support', () => { mockConnectionState.resolve(); await agentPromise; }); + + it('loadSession rejects when sessionExists returns false, skipping config load', async () => { + // Mock SessionService to report session does not exist + const sessionExistsStub = vi.fn().mockResolvedValue(false); + vi.mocked(SessionService).mockImplementation( + () => + ({ + sessionExists: sessionExistsStub, + }) as unknown as InstanceType, + ); + + await setupSessionMocks('session-nonexistent'); + + const agentPromise = runAcpAgent( + mockConfig, + makeSessionSettings(), + mockArgv, + ); + await vi.waitFor(() => expect(capturedAgentFactory).toBeDefined()); + + const agent = capturedAgentFactory!({ + get closed() { + return mockConnectionState.promise; + }, + }) as AgentLike; + + await expect( + agent.loadSession({ + cwd: '/tmp', + sessionId: 'nonexistent', + }), + ).rejects.toThrow('Session nonexistent does not exist at /tmp'); + + // Verify sessionExists was queried for the correct session + expect(sessionExistsStub).toHaveBeenCalledWith('nonexistent'); + + mockConnectionState.resolve(); + await agentPromise; + }); }); From 2e3a52430940f3d5452890515708b680f92ff734 Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Fri, 8 May 2026 15:22:22 +0800 Subject: [PATCH 04/10] refactor: address review suggestions for session-restore PR --- .../channels/base/src/SessionRouter.test.ts | 178 +++++++++--------- .../cli/src/acp-integration/acpAgent.test.ts | 61 +++--- packages/cli/src/commands/channel/start.ts | 26 ++- 3 files changed, 139 insertions(+), 126 deletions(-) diff --git a/packages/channels/base/src/SessionRouter.test.ts b/packages/channels/base/src/SessionRouter.test.ts index 71a37ec91a1..fa1a6dcf19e 100644 --- a/packages/channels/base/src/SessionRouter.test.ts +++ b/packages/channels/base/src/SessionRouter.test.ts @@ -298,46 +298,48 @@ describe('SessionRouter', () => { const tmpDir = join('/tmp', `test-restore-${Date.now()}`); mkdirSync(tmpDir, { recursive: true }); - const persistFile = join(tmpDir, 'sessions.json'); - const entries = { - 'telegram:alice:chat1': { - sessionId: 'old-session-1', - target: { - channelName: 'telegram', - senderId: 'alice', - chatId: 'chat1', + try { + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'old-session-1', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', }, - cwd: '/workspace', - }, - 'telegram:bob:chat2': { - sessionId: 'old-session-2', - target: { - channelName: 'telegram', - senderId: 'bob', - chatId: 'chat2', + 'telegram:bob:chat2': { + sessionId: 'old-session-2', + target: { + channelName: 'telegram', + senderId: 'bob', + chatId: 'chat2', + }, + cwd: '/workspace', }, - cwd: '/workspace', - }, - }; - writeFileSync(persistFile, JSON.stringify(entries, null, 2)); - - const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); - const result = await router.restoreSessions(); + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); - expect(result.restored).toBe(2); - expect(result.failed).toBe(0); - expect(bridge.loadSession).toHaveBeenCalledTimes(2); - expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); - expect(router.hasSession('telegram', 'bob', 'chat2')).toBe(true); + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const result = await router.restoreSessions(); - const target = router.getTarget('old-session-1'); - expect(target).toEqual({ - channelName: 'telegram', - senderId: 'alice', - chatId: 'chat1', - }); + expect(result.restored).toBe(2); + expect(result.failed).toBe(0); + expect(bridge.loadSession).toHaveBeenCalledTimes(2); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); + expect(router.hasSession('telegram', 'bob', 'chat2')).toBe(true); - rmSync(tmpDir, { recursive: true }); + const target = router.getTarget('old-session-1'); + expect(target).toEqual({ + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }); + } finally { + rmSync(tmpDir, { recursive: true }); + } }); it('skips entries whose loadSession throws', async () => { @@ -346,53 +348,55 @@ describe('SessionRouter', () => { const tmpDir = join('/tmp', `test-restore-fail-${Date.now()}`); mkdirSync(tmpDir, { recursive: true }); - const persistFile = join(tmpDir, 'sessions.json'); - const entries = { - 'telegram:alice:chat1': { - sessionId: 'good-session', - target: { - channelName: 'telegram', - senderId: 'alice', - chatId: 'chat1', + try { + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'good-session', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', }, - cwd: '/workspace', - }, - 'telegram:bob:chat2': { - sessionId: 'stale-session', - target: { - channelName: 'telegram', - senderId: 'bob', - chatId: 'chat2', + 'telegram:bob:chat2': { + sessionId: 'stale-session', + target: { + channelName: 'telegram', + senderId: 'bob', + chatId: 'chat2', + }, + cwd: '/workspace', }, - cwd: '/workspace', - }, - }; - writeFileSync(persistFile, JSON.stringify(entries, null, 2)); - - // Make loadSession throw for the stale session - (bridge.loadSession as ReturnType).mockImplementation( - (id: string) => { - if (id === 'stale-session') { - throw new Error('Session not found'); - } - return id; - }, - ); - - const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); - const result = await router.restoreSessions(); - - expect(result.restored).toBe(1); - expect(result.failed).toBe(1); - expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); - expect(router.hasSession('telegram', 'bob', 'chat2')).toBe(false); + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); + + // Make loadSession throw for the stale session + (bridge.loadSession as ReturnType).mockImplementation( + (id: string) => { + if (id === 'stale-session') { + throw new Error('Session not found'); + } + return id; + }, + ); - // Persist file should be updated to remove the failed entry - const { readFileSync } = await import('node:fs'); - const updated = JSON.parse(readFileSync(persistFile, 'utf-8')); - expect(Object.keys(updated)).toEqual(['telegram:alice:chat1']); + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const result = await router.restoreSessions(); - rmSync(tmpDir, { recursive: true }); + expect(result.restored).toBe(1); + expect(result.failed).toBe(1); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); + expect(router.hasSession('telegram', 'bob', 'chat2')).toBe(false); + + // Persist file should be updated to remove the failed entry + const { readFileSync } = await import('node:fs'); + const updated = JSON.parse(readFileSync(persistFile, 'utf-8')); + expect(Object.keys(updated)).toEqual(['telegram:alice:chat1']); + } finally { + rmSync(tmpDir, { recursive: true }); + } }); it('returns zeros when no persist file exists', async () => { @@ -412,16 +416,18 @@ describe('SessionRouter', () => { const tmpDir = join('/tmp', `test-restore-empty-${Date.now()}`); mkdirSync(tmpDir, { recursive: true }); - const persistFile = join(tmpDir, 'sessions.json'); - writeFileSync(persistFile, '{}'); - - const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); - const result = await router.restoreSessions(); + try { + const persistFile = join(tmpDir, 'sessions.json'); + writeFileSync(persistFile, '{}'); - expect(result).toEqual({ restored: 0, failed: 0 }); - expect(bridge.loadSession).not.toHaveBeenCalled(); + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const result = await router.restoreSessions(); - rmSync(tmpDir, { recursive: true }); + expect(result).toEqual({ restored: 0, failed: 0 }); + expect(bridge.loadSession).not.toHaveBeenCalled(); + } finally { + rmSync(tmpDir, { recursive: true }); + } }); }); }); diff --git a/packages/cli/src/acp-integration/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index b148cdea545..da459e1601f 100644 --- a/packages/cli/src/acp-integration/acpAgent.test.ts +++ b/packages/cli/src/acp-integration/acpAgent.test.ts @@ -905,7 +905,8 @@ describe('QwenAgent MCP SSE/HTTP support', () => { }); it('loadSession rejects when sessionExists returns false, skipping config load', async () => { - // Mock SessionService to report session does not exist + // Save the default SessionService mock and override for this test + const savedImpl = vi.mocked(SessionService).getMockImplementation(); const sessionExistsStub = vi.fn().mockResolvedValue(false); vi.mocked(SessionService).mockImplementation( () => @@ -914,32 +915,40 @@ describe('QwenAgent MCP SSE/HTTP support', () => { }) as unknown as InstanceType, ); - await setupSessionMocks('session-nonexistent'); + try { + await setupSessionMocks('session-nonexistent'); - const agentPromise = runAcpAgent( - mockConfig, - makeSessionSettings(), - mockArgv, - ); - await vi.waitFor(() => expect(capturedAgentFactory).toBeDefined()); - - const agent = capturedAgentFactory!({ - get closed() { - return mockConnectionState.promise; - }, - }) as AgentLike; - - await expect( - agent.loadSession({ - cwd: '/tmp', - sessionId: 'nonexistent', - }), - ).rejects.toThrow('Session nonexistent does not exist at /tmp'); - - // Verify sessionExists was queried for the correct session - expect(sessionExistsStub).toHaveBeenCalledWith('nonexistent'); + const agentPromise = runAcpAgent( + mockConfig, + makeSessionSettings(), + mockArgv, + ); + await vi.waitFor(() => expect(capturedAgentFactory).toBeDefined()); - mockConnectionState.resolve(); - await agentPromise; + const agent = capturedAgentFactory!({ + get closed() { + return mockConnectionState.promise; + }, + }) as AgentLike; + + await expect( + agent.loadSession({ + cwd: '/tmp', + sessionId: 'nonexistent', + }), + ).rejects.toThrow('Session nonexistent does not exist at /tmp'); + + // Verify sessionExists was queried for the correct session + expect(sessionExistsStub).toHaveBeenCalledWith('nonexistent'); + + mockConnectionState.resolve(); + await agentPromise; + } finally { + // Restore the default mock to avoid polluting subsequent tests + vi.mocked(SessionService).mockImplementation( + savedImpl as typeof vi.mocked, + ); + // Also let beforeEach clearAllMocks handle the sessionExists stub + } }); }); diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index bef953a2c38..b9b82da2e5d 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -57,6 +57,16 @@ function sessionsPath(): string { return path.join(os.homedir(), '.qwen', 'channels', 'sessions.json'); } +async function restoreAndLogSessions(router: SessionRouter): Promise { + const result = await router.restoreSessions(); + if (result.restored > 0 || result.failed > 0) { + writeStdoutLine( + `[Channel] Sessions restored: ${result.restored}` + + (result.failed > 0 ? `, failed: ${result.failed}` : ''), + ); + } +} + function loadChannelsConfig(): Record { const settings = loadSettings(process.cwd()); const channels = ( @@ -216,13 +226,7 @@ async function startSingle(name: string, proxy?: string): Promise { registerToolCallDispatch(bridge, router, channels); // Restore sessions from previous run - const restoreResult = await router.restoreSessions(); - if (restoreResult.restored > 0 || restoreResult.failed > 0) { - writeStdoutLine( - `[Channel] Sessions restored: ${restoreResult.restored}` + - (restoreResult.failed > 0 ? `, failed: ${restoreResult.failed}` : ''), - ); - } + await restoreAndLogSessions(router); try { await channel.connect(); @@ -375,13 +379,7 @@ async function startAll(proxy?: string): Promise { registerToolCallDispatch(bridge, router, channels); // Restore sessions from previous run - const restoreResult = await router.restoreSessions(); - if (restoreResult.restored > 0 || restoreResult.failed > 0) { - writeStdoutLine( - `[Channel] Sessions restored: ${restoreResult.restored}` + - (restoreResult.failed > 0 ? `, failed: ${restoreResult.failed}` : ''), - ); - } + await restoreAndLogSessions(router); // Connect all channels let connectedCount = 0; From cb0242628ab24be1b13dc0728ca46d831002d8cd Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Fri, 8 May 2026 23:04:09 +0800 Subject: [PATCH 05/10] fix: prevent data loss and crashes in session restore --- packages/channels/base/src/SessionRouter.ts | 35 ++++++++++++++++++--- packages/cli/src/commands/channel/start.ts | 12 ++----- 2 files changed, 33 insertions(+), 14 deletions(-) diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index 29a6b9da388..2b504582359 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -164,8 +164,28 @@ export class SessionRouter { let entries: Record; try { - entries = JSON.parse(readFileSync(this.persistPath, 'utf-8')); + const raw = readFileSync(this.persistPath, 'utf-8'); + const parsed: unknown = JSON.parse(raw); + + // Guard against JSON.parse(null) and other non-object results + if ( + parsed === null || + typeof parsed !== 'object' || + Array.isArray(parsed) + ) { + process.stderr.write( + `[SessionRouter] Corrupted persist file (expected object, got ${parsed === null ? 'null' : Array.isArray(parsed) ? 'array' : typeof parsed}).\n`, + ); + return { restored: 0, failed: 0 }; + } + + entries = parsed as Record; } catch { + // Corrupted file — warn and return empty so the next persist() + // can write a clean file after a successful restore. + process.stderr.write( + `[SessionRouter] Failed to parse persist file: ${this.persistPath}\n`, + ); return { restored: 0, failed: 0 }; } @@ -182,14 +202,21 @@ export class SessionRouter { this.toTarget.set(sessionId, entry.target); this.toCwd.set(sessionId, entry.cwd); restored++; - } catch { + } catch (err: unknown) { // Session can't be loaded — will create fresh on next message + const msg = err instanceof Error ? err.message : String(err); + process.stderr.write( + `[SessionRouter] Failed to restore session ${entry.sessionId}: ${msg}\n`, + ); failed++; } } - // Update persist file to only include successfully restored sessions - if (failed > 0) { + // Update persist file to remove failed entries. + // Only persist if at least one session was successfully restored — + // an all-failure persist() would write empty {} and permanently delete + // routing data that might be recoverable on a subsequent attempt. + if (failed > 0 && restored > 0) { this.persist(); } diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index b9b82da2e5d..435130c1b4f 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -271,14 +271,10 @@ async function startSingle(name: string, proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); + await router.restoreSessions(); channel.setBridge(bridge); registerToolCallDispatch(bridge, router, channels); attachDisconnectHandler(bridge); - - const result = await router.restoreSessions(); - writeStdoutLine( - `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, - ); } catch (err) { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, @@ -441,16 +437,12 @@ async function startAll(proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); + await router.restoreSessions(); for (const channel of channels.values()) { channel.setBridge(bridge); } registerToolCallDispatch(bridge, router, channels); attachDisconnectHandler(bridge); - - const result = await router.restoreSessions(); - writeStdoutLine( - `[Channel] Bridge restarted. Sessions restored: ${result.restored}, failed: ${result.failed}`, - ); } catch (err) { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, From c832ca5ea3eb76aaa455ff9e3ab8f821106798a9 Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Sat, 9 May 2026 08:31:21 +0800 Subject: [PATCH 06/10] fix(base): address latest review feedback on session persist - Add restoreAndLogSessions() in bridge crash-recovery paths (startSingle/startAll) - Add TODO for all-failure accumulation when restored === 0 - Update JSDoc to reflect stderr logging of failed loads - Validate entry shape before calling loadSession to guard corrupted persist files - Fix test temp dir: typo 'cleara11' and use mkdtempSync+tmpdir for cross-platform - Add comment explaining AcpBridge.loadSession returns caller-supplied sessionId - Fix TS cast in acpAgent.test.ts: savedImpl type error --- packages/channels/base/src/AcpBridge.ts | 2 ++ .../channels/base/src/SessionRouter.test.ts | 24 +++++++++---------- packages/channels/base/src/SessionRouter.ts | 21 +++++++++++++++- .../cli/src/acp-integration/acpAgent.test.ts | 2 +- packages/cli/src/commands/channel/start.ts | 4 ++-- 5 files changed, 37 insertions(+), 16 deletions(-) diff --git a/packages/channels/base/src/AcpBridge.ts b/packages/channels/base/src/AcpBridge.ts index afcfe11333b..38de54189a9 100644 --- a/packages/channels/base/src/AcpBridge.ts +++ b/packages/channels/base/src/AcpBridge.ts @@ -137,6 +137,8 @@ export class AcpBridge extends EventEmitter { cwd, mcpServers: [], }); + // LoadSessionResponse has no sessionId field — the ACP agent loads + // the session under the caller-supplied ID, so we echo it back. return sessionId; } diff --git a/packages/channels/base/src/SessionRouter.test.ts b/packages/channels/base/src/SessionRouter.test.ts index fa1a6dcf19e..0216b0ff026 100644 --- a/packages/channels/base/src/SessionRouter.test.ts +++ b/packages/channels/base/src/SessionRouter.test.ts @@ -219,11 +219,11 @@ describe('SessionRouter', () => { }); it('preserves persist file for restoration on next start', async () => { - const { mkdirSync, writeFileSync, rmSync, readFileSync, existsSync } = + const { writeFileSync, rmSync, readFileSync, existsSync, mkdtempSync } = await import('node:fs'); const { join } = await import('node:path'); - const tmpDir = join('/tmp', `test-cleara11-persist-${Date.now()}`); - mkdirSync(tmpDir, { recursive: true }); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-clearAll-persist-')); const persistFile = join(tmpDir, 'sessions.json'); const originalEntries = { @@ -293,10 +293,10 @@ describe('SessionRouter', () => { describe('restoreSessions', () => { it('restores sessions from persist file', async () => { - const { mkdirSync, writeFileSync, rmSync } = await import('node:fs'); + const { writeFileSync, rmSync, mkdtempSync } = await import('node:fs'); const { join } = await import('node:path'); - const tmpDir = join('/tmp', `test-restore-${Date.now()}`); - mkdirSync(tmpDir, { recursive: true }); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-restore-')); try { const persistFile = join(tmpDir, 'sessions.json'); @@ -343,10 +343,10 @@ describe('SessionRouter', () => { }); it('skips entries whose loadSession throws', async () => { - const { mkdirSync, writeFileSync, rmSync } = await import('node:fs'); + const { writeFileSync, rmSync, mkdtempSync } = await import('node:fs'); const { join } = await import('node:path'); - const tmpDir = join('/tmp', `test-restore-fail-${Date.now()}`); - mkdirSync(tmpDir, { recursive: true }); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-restore-fail-')); try { const persistFile = join(tmpDir, 'sessions.json'); @@ -411,10 +411,10 @@ describe('SessionRouter', () => { }); it('returns zeros when persist file is empty', async () => { - const { mkdirSync, writeFileSync, rmSync } = await import('node:fs'); + const { writeFileSync, rmSync, mkdtempSync } = await import('node:fs'); const { join } = await import('node:path'); - const tmpDir = join('/tmp', `test-restore-empty-${Date.now()}`); - mkdirSync(tmpDir, { recursive: true }); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-restore-empty-')); try { const persistFile = join(tmpDir, 'sessions.json'); diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index 2b504582359..2caf3420e0c 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -152,7 +152,8 @@ export class SessionRouter { /** * Restore session mappings from a previous bridge. * Called after bridge restart — attempts loadSession for each saved mapping. - * Failed loads are silently dropped (new session on next message). + * Failed loads are logged to stderr and dropped from the persist file + * (a new session will be created on the next incoming message). */ async restoreSessions(): Promise<{ restored: number; @@ -193,6 +194,19 @@ export class SessionRouter { let failed = 0; for (const [key, entry] of Object.entries(entries)) { + // Guard against malformed entries from a corrupted persist file. + if ( + typeof entry?.sessionId !== 'string' || + typeof entry?.cwd !== 'string' || + typeof entry?.target?.channelName !== 'string' + ) { + process.stderr.write( + `[SessionRouter] Skipping malformed entry for key ${key}\n`, + ); + failed++; + continue; + } + try { const sessionId = await this.bridge.loadSession( entry.sessionId, @@ -216,6 +230,11 @@ export class SessionRouter { // Only persist if at least one session was successfully restored — // an all-failure persist() would write empty {} and permanently delete // routing data that might be recoverable on a subsequent attempt. + // + // TODO: When failed > 0 && restored === 0, the persist file is never + // trimmed — every restart retries the same dead session IDs. Consider + // an escape hatch: after N consecutive all-failure cycles, force-prune + // entries that the ACP agent explicitly rejects (sessionExists=false). if (failed > 0 && restored > 0) { this.persist(); } diff --git a/packages/cli/src/acp-integration/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index da459e1601f..e1d17d68498 100644 --- a/packages/cli/src/acp-integration/acpAgent.test.ts +++ b/packages/cli/src/acp-integration/acpAgent.test.ts @@ -946,7 +946,7 @@ describe('QwenAgent MCP SSE/HTTP support', () => { } finally { // Restore the default mock to avoid polluting subsequent tests vi.mocked(SessionService).mockImplementation( - savedImpl as typeof vi.mocked, + (savedImpl ?? vi.fn()) as (cwd: string) => SessionService, ); // Also let beforeEach clearAllMocks handle the sessionExists stub } diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index 435130c1b4f..e9099f552d9 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -271,7 +271,7 @@ async function startSingle(name: string, proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); - await router.restoreSessions(); + await restoreAndLogSessions(router); channel.setBridge(bridge); registerToolCallDispatch(bridge, router, channels); attachDisconnectHandler(bridge); @@ -437,7 +437,7 @@ async function startAll(proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); - await router.restoreSessions(); + await restoreAndLogSessions(router); for (const channel of channels.values()) { channel.setBridge(bridge); } From 06036bd6babe407577a4b495125226b8804cd310 Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Sat, 9 May 2026 14:51:21 +0800 Subject: [PATCH 07/10] fix(base): clean stale mappings on failed restore and fix setBridge ordering Two crash-recovery issues: 1. restoreSessions() now deletes in-memory toSession/toTarget/toCwd entries when loadSession fails, preventing stale session IDs from being re-persisted and silently dropping messages after bridge restart. 2. Reorder crash-recovery in startSingle/startAll: channel.setBridge(), registerToolCallDispatch(), and attachDisconnectHandler() now run before restoreAndLogSessions(), closing a race where the channel held a dead bridge while the router already pointed at the new one. --- .../channels/base/src/SessionRouter.test.ts | 51 +++++++++++++++++++ packages/channels/base/src/SessionRouter.ts | 6 ++- packages/cli/src/commands/channel/start.ts | 4 +- 3 files changed, 58 insertions(+), 3 deletions(-) diff --git a/packages/channels/base/src/SessionRouter.test.ts b/packages/channels/base/src/SessionRouter.test.ts index 0216b0ff026..e3f8bd414ab 100644 --- a/packages/channels/base/src/SessionRouter.test.ts +++ b/packages/channels/base/src/SessionRouter.test.ts @@ -399,6 +399,57 @@ describe('SessionRouter', () => { } }); + it('removes stale in-memory mappings on failed restore (crash recovery)', async () => { + const { writeFileSync, rmSync, mkdtempSync } = await import('node:fs'); + const { join } = await import('node:path'); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-restore-crash-')); + + try { + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'will-fail', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', + }, + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); + + // Simulate pre-crash in-memory state: the router already has a mapping + // for this key from before the bridge crash. + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + (bridge.newSession as ReturnType).mockResolvedValueOnce( + 'will-fail', + ); + await router.resolve( + 'telegram', + 'alice', + 'chat1', + undefined, + '/workspace', + ); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); + + // Now make loadSession throw so restoreSessions hits the catch path. + (bridge.loadSession as ReturnType).mockRejectedValueOnce( + new Error('Session not found'), + ); + + const result = await router.restoreSessions(); + expect(result.restored).toBe(0); + expect(result.failed).toBe(1); + // The stale in-memory mapping should be cleaned up + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(false); + } finally { + rmSync(tmpDir, { recursive: true }); + } + }); + it('returns zeros when no persist file exists', async () => { const router = new SessionRouter( bridge, diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index 2caf3420e0c..ec1b3c276f0 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -217,11 +217,15 @@ export class SessionRouter { this.toCwd.set(sessionId, entry.cwd); restored++; } catch (err: unknown) { - // Session can't be loaded — will create fresh on next message + // Session can't be loaded — remove stale in-memory mapping so + // subsequent resolve() / persist() don't reference a dead session. const msg = err instanceof Error ? err.message : String(err); process.stderr.write( `[SessionRouter] Failed to restore session ${entry.sessionId}: ${msg}\n`, ); + this.toSession.delete(key); + this.toTarget.delete(entry.sessionId); + this.toCwd.delete(entry.sessionId); failed++; } } diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index e9099f552d9..adbe9181c08 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -271,10 +271,10 @@ async function startSingle(name: string, proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); - await restoreAndLogSessions(router); channel.setBridge(bridge); registerToolCallDispatch(bridge, router, channels); attachDisconnectHandler(bridge); + await restoreAndLogSessions(router); } catch (err) { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, @@ -437,12 +437,12 @@ async function startAll(proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); - await restoreAndLogSessions(router); for (const channel of channels.values()) { channel.setBridge(bridge); } registerToolCallDispatch(bridge, router, channels); attachDisconnectHandler(bridge); + await restoreAndLogSessions(router); } catch (err) { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, From ac56d92b9fc04da81a7dcfa056f4000fd12ff44c Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Sun, 10 May 2026 18:31:26 +0800 Subject: [PATCH 08/10] fix(base): add loadSession timeout to prevent indefinite blocking during session restore --- packages/channels/base/src/SessionRouter.ts | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index ec1b3c276f0..612f05d41af 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -8,6 +8,8 @@ interface PersistedEntry { cwd: string; } +const LOAD_SESSION_TIMEOUT_MS = 30_000; + export class SessionRouter { private toSession: Map = new Map(); // routing key → session ID private toTarget: Map = new Map(); // session ID → target @@ -208,10 +210,15 @@ export class SessionRouter { } try { - const sessionId = await this.bridge.loadSession( - entry.sessionId, - entry.cwd, - ); + const sessionId = await Promise.race([ + this.bridge.loadSession(entry.sessionId, entry.cwd), + new Promise((_, reject) => + setTimeout( + () => reject(new Error('loadSession timed out')), + LOAD_SESSION_TIMEOUT_MS, + ), + ), + ]); this.toSession.set(key, sessionId); this.toTarget.set(sessionId, entry.target); this.toCwd.set(sessionId, entry.cwd); From 04f24d11af14dddb1db0a19eeb405984fddc7973 Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Mon, 11 May 2026 22:41:41 +0800 Subject: [PATCH 09/10] fix(base): address critical review feedback on session persistence - Atomic write: use temp file + rename for persist() to prevent data corruption on crash or disk full - Add error logging in persist() catch block instead of silent failure - Fix timeout handling: clear setTimeout on success and attach no-op catch to loadSession promise to prevent unhandled rejection - Add try/catch around restoreAndLogSessions in startSingle/startAll initial startup paths for consistent error handling - Add tests for loadSession timeout and timer cleanup scenarios --- .../channels/base/src/SessionRouter.test.ts | 107 ++++++++++++++++++ packages/channels/base/src/SessionRouter.ts | 50 ++++++-- packages/cli/src/commands/channel/start.ts | 16 ++- 3 files changed, 162 insertions(+), 11 deletions(-) diff --git a/packages/channels/base/src/SessionRouter.test.ts b/packages/channels/base/src/SessionRouter.test.ts index e3f8bd414ab..3d4e07ee0d5 100644 --- a/packages/channels/base/src/SessionRouter.test.ts +++ b/packages/channels/base/src/SessionRouter.test.ts @@ -480,5 +480,112 @@ describe('SessionRouter', () => { rmSync(tmpDir, { recursive: true }); } }); + + it('times out when loadSession hangs and logs error to stderr', async () => { + const { writeFileSync, rmSync, mkdtempSync } = await import('node:fs'); + const { join } = await import('node:path'); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-restore-timeout-')); + + vi.useFakeTimers(); + + try { + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'hanging-session', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', + }, + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); + + // Make loadSession hang indefinitely (never resolves/rejects) + let hangingResolve: (value: string) => void; + const hangingPromise = new Promise((resolve) => { + hangingResolve = resolve; + }); + (bridge.loadSession as ReturnType).mockReturnValue( + hangingPromise, + ); + + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + + // Start restoreSessions - it will hang waiting for loadSession + const resultPromise = router.restoreSessions(); + + // Advance time past the 30s timeout + await vi.advanceTimersByTimeAsync(30_000); + + const result = await resultPromise; + + // Should have recorded a failure due to timeout + expect(result.restored).toBe(0); + expect(result.failed).toBe(1); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(false); + + // The hanging loadSession should have a no-op catch attached + // to prevent unhandled rejection. Resolve it now to verify no crash. + hangingResolve!('hanging-session'); + await Promise.resolve(); // let microtasks settle + } finally { + vi.useRealTimers(); + rmSync(tmpDir, { recursive: true }); + } + }); + + it('clears timeout when loadSession completes before timeout', async () => { + const { writeFileSync, rmSync, mkdtempSync } = await import('node:fs'); + const { join } = await import('node:path'); + const { tmpdir } = await import('node:os'); + const tmpDir = mkdtempSync(join(tmpdir(), 'test-restore-notimeout-')); + + vi.useFakeTimers(); + + try { + const persistFile = join(tmpDir, 'sessions.json'); + const entries = { + 'telegram:alice:chat1': { + sessionId: 'fast-session', + target: { + channelName: 'telegram', + senderId: 'alice', + chatId: 'chat1', + }, + cwd: '/workspace', + }, + }; + writeFileSync(persistFile, JSON.stringify(entries, null, 2)); + + // Make loadSession resolve quickly (within timeout) + (bridge.loadSession as ReturnType).mockImplementation( + (id: string) => Promise.resolve(id), + ); + + const router = new SessionRouter(bridge, '/tmp', 'user', persistFile); + const resultPromise = router.restoreSessions(); + + // Advance time a bit (but not past timeout) + await vi.advanceTimersByTimeAsync(1000); + + const result = await resultPromise; + + // Should have succeeded + expect(result.restored).toBe(1); + expect(result.failed).toBe(0); + expect(router.hasSession('telegram', 'alice', 'chat1')).toBe(true); + + // Advance past the original timeout to ensure no stray timer fires + await vi.advanceTimersByTimeAsync(30_000); + // No error should occur (timeout was cleared) + } finally { + vi.useRealTimers(); + rmSync(tmpDir, { recursive: true }); + } + }); }); }); diff --git a/packages/channels/base/src/SessionRouter.ts b/packages/channels/base/src/SessionRouter.ts index 612f05d41af..e291bd944ff 100644 --- a/packages/channels/base/src/SessionRouter.ts +++ b/packages/channels/base/src/SessionRouter.ts @@ -1,4 +1,13 @@ -import { existsSync, readFileSync, writeFileSync } from 'node:fs'; +import { + existsSync, + mkdtempSync, + readFileSync, + renameSync, + rmSync, + writeFileSync, +} from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; import type { SessionScope, SessionTarget } from './types.js'; import type { AcpBridge } from './AcpBridge.js'; @@ -209,16 +218,23 @@ export class SessionRouter { continue; } + // Race loadSession against a timeout to avoid indefinite hangs. + // We keep a reference to the loadSession promise so we can attach + // a no-op catch handler after timeout, preventing unhandled rejection. + let loadPromise: Promise | undefined; try { + loadPromise = this.bridge.loadSession(entry.sessionId, entry.cwd); + let timeoutId: ReturnType; const sessionId = await Promise.race([ - this.bridge.loadSession(entry.sessionId, entry.cwd), - new Promise((_, reject) => - setTimeout( + loadPromise, + new Promise((_, reject) => { + timeoutId = setTimeout( () => reject(new Error('loadSession timed out')), LOAD_SESSION_TIMEOUT_MS, - ), - ), + ); + }), ]); + clearTimeout(timeoutId!); this.toSession.set(key, sessionId); this.toTarget.set(sessionId, entry.target); this.toCwd.set(sessionId, entry.cwd); @@ -226,6 +242,9 @@ export class SessionRouter { } catch (err: unknown) { // Session can't be loaded — remove stale in-memory mapping so // subsequent resolve() / persist() don't reference a dead session. + // Note: loadSession may still be running in background after timeout; + // attach a no-op catch to prevent unhandled rejection warnings. + loadPromise?.catch(() => {}); const msg = err instanceof Error ? err.message : String(err); process.stderr.write( `[SessionRouter] Failed to restore session ${entry.sessionId}: ${msg}\n`, @@ -277,9 +296,22 @@ export class SessionRouter { } try { - writeFileSync(this.persistPath, JSON.stringify(data, null, 2), 'utf-8'); - } catch { - // best-effort — don't break message flow for persistence failure + // Atomic write: write to temp file in a temp directory then rename. + // Using a temp directory (not just a temp file) ensures the temp file + // is on the same filesystem as the target for atomic rename. + // rename is atomic on POSIX and works reliably on Windows when + // source and target are on the same filesystem. + const tmpDir = mkdtempSync(join(tmpdir(), 'session-persist-')); + const tmpFile = join(tmpDir, 'sessions.json'); + writeFileSync(tmpFile, JSON.stringify(data, null, 2), 'utf-8'); + renameSync(tmpFile, this.persistPath); + // Clean up the now-empty temp directory (file was moved by rename) + rmSync(tmpDir, { recursive: true }); + } catch (err: unknown) { + const msg = err instanceof Error ? err.message : String(err); + process.stderr.write( + `[SessionRouter] Failed to persist sessions to ${this.persistPath}: ${msg}\n`, + ); } } } diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index adbe9181c08..7b1aae25f5d 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -226,7 +226,13 @@ async function startSingle(name: string, proxy?: string): Promise { registerToolCallDispatch(bridge, router, channels); // Restore sessions from previous run - await restoreAndLogSessions(router); + try { + await restoreAndLogSessions(router); + } catch (err) { + writeStderrLine( + `[Channel] Failed to restore sessions: ${err instanceof Error ? err.message : String(err)}`, + ); + } try { await channel.connect(); @@ -375,7 +381,13 @@ async function startAll(proxy?: string): Promise { registerToolCallDispatch(bridge, router, channels); // Restore sessions from previous run - await restoreAndLogSessions(router); + try { + await restoreAndLogSessions(router); + } catch (err) { + writeStderrLine( + `[Channel] Failed to restore sessions: ${err instanceof Error ? err.message : String(err)}`, + ); + } // Connect all channels let connectedCount = 0; From 7e40926e330a8c0433343a34622553abc6f843aa Mon Sep 17 00:00:00 2001 From: Maidong <408097061@qq.com> Date: Tue, 12 May 2026 22:10:42 +0800 Subject: [PATCH 10/10] fix(base): improve crash recovery ordering and error messages - Move restoreAndLogSessions before channel.setBridge to prevent new sessions from being overwritten by stale restore data - Split error handling to distinguish bridge startup failures from session restore failures with appropriate error messages --- packages/cli/src/commands/channel/start.ts | 38 ++++++++++++++++------ 1 file changed, 28 insertions(+), 10 deletions(-) diff --git a/packages/cli/src/commands/channel/start.ts b/packages/cli/src/commands/channel/start.ts index 7b1aae25f5d..7ba383dda47 100644 --- a/packages/cli/src/commands/channel/start.ts +++ b/packages/cli/src/commands/channel/start.ts @@ -277,15 +277,24 @@ async function startSingle(name: string, proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); - channel.setBridge(bridge); - registerToolCallDispatch(bridge, router, channels); - attachDisconnectHandler(bridge); - await restoreAndLogSessions(router); } catch (err) { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, ); + return; + } + + try { + await restoreAndLogSessions(router); + } catch (err) { + writeStderrLine( + `[Channel] Bridge restarted but session restore failed: ${err instanceof Error ? err.message : String(err)}`, + ); } + + channel.setBridge(bridge); + registerToolCallDispatch(bridge, router, channels); + attachDisconnectHandler(bridge); }); }; attachDisconnectHandler(bridge); @@ -449,17 +458,26 @@ async function startAll(proxy?: string): Promise { bridge = new AcpBridge(bridgeOpts); await bridge.start(); router.setBridge(bridge); - for (const channel of channels.values()) { - channel.setBridge(bridge); - } - registerToolCallDispatch(bridge, router, channels); - attachDisconnectHandler(bridge); - await restoreAndLogSessions(router); } catch (err) { writeStderrLine( `[Channel] Failed to restart bridge: ${err instanceof Error ? err.message : String(err)}`, ); + return; + } + + try { + await restoreAndLogSessions(router); + } catch (err) { + writeStderrLine( + `[Channel] Bridge restarted but session restore failed: ${err instanceof Error ? err.message : String(err)}`, + ); + } + + for (const channel of channels.values()) { + channel.setBridge(bridge); } + registerToolCallDispatch(bridge, router, channels); + attachDisconnectHandler(bridge); }); }; attachDisconnectHandler(bridge);