diff --git a/adapters/common/__tests__/chat-runtime.test.ts b/adapters/common/__tests__/chat-runtime.test.ts index f80e95ec..9c38bf84 100644 --- a/adapters/common/__tests__/chat-runtime.test.ts +++ b/adapters/common/__tests__/chat-runtime.test.ts @@ -266,6 +266,69 @@ describe('ImChatRuntime authorization', () => { }) describe('ImChatRuntime session lifecycle', () => { + it.each(['/status', '/stop', '/clear'])('reports a retryable disconnect for %s without changing the binding', async (command) => { + const { runtime, bridge, httpClient, sessionStore, notices } = createRuntime() + sessionStore.set(CHAT_ID, 'old-session', path.join(tmpDir, 'original-project')) + httpClient.existingSessions.add('old-session') + const original = sessionStore.get(CHAT_ID) + bridge.waitForOpen = async () => false + + await inbound(runtime, command) + + expect(sessionStore.get(CHAT_ID)).toEqual(original) + expect(httpClient.createdSessions).toEqual([]) + expect(bridge.sent).toEqual([]) + expect(bridge.stopped).toEqual([]) + expect(notices).toHaveLength(1) + expect(notices[0]).toContain('重试') + expect(notices[0]).not.toContain('当前没有活动会话') + }) + + it.each(['timeout', 'error'])('preserves the stored session and project after a reconnect %s, then retries it', async (failure) => { + const { runtime, bridge, httpClient, sessionStore, notices } = createRuntime() + const workDir = path.join(tmpDir, 'original-project') + sessionStore.set(CHAT_ID, 'old-session', workDir) + httpClient.existingSessions.add('old-session') + const original = sessionStore.get(CHAT_ID) + let attempts = 0 + bridge.isSessionOpen = () => false + bridge.waitForOpen = async () => { + attempts += 1 + if (attempts === 1) { + if (failure === 'error') throw new Error('connection failed') + return false + } + return true + } + + await inbound(runtime, 'continue') + + expect(sessionStore.get(CHAT_ID)).toEqual(original) + expect(new SessionStore(path.join(tmpDir, 'adapter-sessions.json')).get(CHAT_ID)).toEqual(original) + expect(httpClient.createdSessions).toEqual([]) + expect(bridge.sent).toEqual([]) + expect(notices.join('\n')).toContain('重试') + + await inbound(runtime, 'continue') + + expect(attempts).toBe(2) + expect(bridge.getSessionId(CHAT_ID)).toBe('old-session') + expect(sessionStore.get(CHAT_ID)).toEqual(original) + expect(httpClient.createdSessions).toEqual([]) + expect(bridge.sent.map((item) => item.content)).toEqual(['continue']) + }) + + it('creates a replacement only when the stored session is confirmed missing', async () => { + const { runtime, bridge, httpClient, sessionStore } = createRuntime() + sessionStore.set(CHAT_ID, 'deleted-session', path.join(tmpDir, 'original-project')) + + await inbound(runtime, 'continue') + + expect(httpClient.createdSessions).toEqual([tmpDir]) + expect(sessionStore.get(CHAT_ID)?.sessionId).not.toBe('deleted-session') + expect(bridge.sent.map((item) => item.content)).toEqual(['continue']) + }) + it('creates a session on the first message and reuses it for the second', async () => { const { runtime, httpClient, bridge } = createRuntime() diff --git a/adapters/common/__tests__/session-recovery.test.ts b/adapters/common/__tests__/session-recovery.test.ts index d3d0d956..f42d2863 100644 --- a/adapters/common/__tests__/session-recovery.test.ts +++ b/adapters/common/__tests__/session-recovery.test.ts @@ -49,6 +49,59 @@ function makeBridge(options?: { currentSessionId?: string | null; open?: boolean } describe('restoreStoredSessionBinding', () => { + it.each(['timeout', 'connect error', 'wait error'])('preserves the binding and transient state after a %s', async (failure) => { + const entry = { sessionId: 'stored-session', workDir: '/original-project', updatedAt: 1 } + const store = makeStore(entry) + const bridge = makeBridge() + const connect = bridge.connectSession + const wait = bridge.waitForOpen + if (failure === 'connect error') bridge.connectSession = () => { throw new Error('connect failed') } + else bridge.waitForOpen = async () => { + if (failure === 'wait error') throw new Error('wait failed') + return false + } + let cleared = 0 + const options = { + chatId: 'chat-1', + bridge, + sessionStore: store, + httpClient: { sessionExists: async () => true }, + onServerMessage: () => {}, + logPrefix: '[Test]', + clearTransientState: () => { cleared += 1 }, + } + + expect(await restoreStoredSessionBinding(options)).toEqual({ status: 'unavailable', session: entry }) + expect(store.current()).toEqual(entry) + expect(cleared).toBe(0) + expect(bridge.calls).not.toContain('reset') + + bridge.connectSession = connect + bridge.waitForOpen = wait + expect(await restoreStoredSessionBinding(options)).toEqual({ status: 'restored', session: entry }) + expect(store.current()).toEqual(entry) + }) + + it.each([false, true])('keeps the binding when verification fails and reconnect returns %s', async (opened) => { + const entry = { sessionId: 'stored-session', workDir: '/original-project', updatedAt: 1 } + const store = makeStore(entry) + const bridge = makeBridge() + bridge.waitForOpen = async () => opened + + const result = await restoreStoredSessionBinding({ + chatId: 'chat-1', + bridge, + sessionStore: store, + httpClient: { sessionExists: async () => { throw new Error('HTTP timeout') } }, + onServerMessage: () => {}, + logPrefix: '[Test]', + }) + + expect(result).toEqual({ status: opened ? 'restored' : 'unavailable', session: entry }) + expect(store.current()).toEqual(entry) + expect(bridge.calls).toEqual(['connect:stored-session', 'handler']) + }) + it('resets stale bridge memory when server-side delete removed the stored mapping', async () => { const store = makeStore(null) const bridge = makeBridge({ currentSessionId: 'deleted-session', hasSession: true }) @@ -66,7 +119,7 @@ describe('restoreStoredSessionBinding', () => { }, }) - expect(restored).toBeNull() + expect(restored).toEqual({ status: 'missing' }) expect(bridge.calls).toEqual(['reset']) expect(cleared).toBe(1) }) @@ -89,7 +142,7 @@ describe('restoreStoredSessionBinding', () => { }, }) - expect(restored).toBeNull() + expect(restored).toEqual({ status: 'missing' }) expect(store.current()).toBeNull() expect(bridge.calls).toEqual([]) expect(cleared).toBe(1) @@ -113,7 +166,7 @@ describe('restoreStoredSessionBinding', () => { }, }) - expect(restored).toEqual(entry) + expect(restored).toEqual({ status: 'restored', session: entry }) expect(bridge.calls).toEqual(['reset', 'connect:stored-session', 'handler', 'wait']) expect(cleared).toBe(1) }) @@ -138,7 +191,7 @@ describe('restoreStoredSessionBinding', () => { logPrefix: '[Test]', }) - expect(restored).toEqual(entry) + expect(restored).toEqual({ status: 'restored', session: entry }) expect(checked).toBe(0) expect(bridge.calls).toEqual([]) }) @@ -157,7 +210,7 @@ describe('restoreStoredSessionBinding', () => { logPrefix: '[Test]', }) - expect(restored).toEqual(entry) + expect(restored).toEqual({ status: 'restored', session: entry }) expect(bridge.calls).toEqual(['connect:stored-session', 'handler', 'wait']) }) }) diff --git a/adapters/common/chat-runtime.ts b/adapters/common/chat-runtime.ts index 3db366d6..56d22ffd 100644 --- a/adapters/common/chat-runtime.ts +++ b/adapters/common/chat-runtime.ts @@ -34,7 +34,7 @@ import { parsePermissionCommand, type PermissionDecision, } from './permission.js' -import { restoreStoredSessionBinding } from './session-recovery.js' +import { restoreStoredSessionBinding, SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from './session-recovery.js' import { SessionSelectionController } from './session-selection.js' import { syncImPermissionState } from './permission-sync.js' import type { SessionStore } from './session-store.js' @@ -300,9 +300,9 @@ export class ImChatRuntime { return true } if (STOP_ALIASES.has(text)) { - const stored = await this.ensureExistingSession(chatId) - if (!stored) { - await this.port.sendNotice(chatId, formatImStatus(null)) + const result = await this.ensureExistingSession(chatId) + if (result.status !== 'restored') { + await this.port.sendNotice(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return true } this.bridge.sendStopGeneration(chatId) @@ -310,9 +310,9 @@ export class ImChatRuntime { return true } if (CLEAR_ALIASES.has(text)) { - const stored = await this.ensureExistingSession(chatId) - if (!stored) { - await this.port.sendNotice(chatId, formatImStatus(null)) + const result = await this.ensureExistingSession(chatId) + if (result.status !== 'restored') { + await this.port.sendNotice(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return true } this.clearTransientChatState(chatId) @@ -502,7 +502,7 @@ export class ImChatRuntime { // ---------- sessions ---------- - async ensureExistingSession(chatId: string): Promise<{ sessionId: string; workDir: string } | null> { + async ensureExistingSession(chatId: string): Promise { return await restoreStoredSessionBinding({ chatId, bridge: this.bridge, @@ -515,8 +515,12 @@ export class ImChatRuntime { } private async ensureSession(chatId: string): Promise { - const stored = await this.ensureExistingSession(chatId) - if (stored) return true + const result = await this.ensureExistingSession(chatId) + if (result.status === 'restored') return true + if (result.status === 'unavailable') { + await this.port.sendNotice(chatId, SESSION_RECONNECT_NOTICE) + return false + } if (this.defaultWorkDir) { return await this.createSessionForChat(chatId, this.defaultWorkDir) @@ -621,8 +625,10 @@ export class ImChatRuntime { } async buildStatusText(chatId: string): Promise { - const stored = await this.ensureExistingSession(chatId) - if (!stored) return formatImStatus(null) + const result = await this.ensureExistingSession(chatId) + if (result.status === 'unavailable') return SESSION_RECONNECT_NOTICE + if (result.status === 'missing') return formatImStatus(null) + const stored = result.session const runtime = this.getRuntimeState(chatId) let projectName = path.basename(stored.workDir) || stored.workDir diff --git a/adapters/common/session-recovery.ts b/adapters/common/session-recovery.ts index fa44f6d2..48bb1dc3 100644 --- a/adapters/common/session-recovery.ts +++ b/adapters/common/session-recovery.ts @@ -2,6 +2,13 @@ import type { AdapterHttpClient } from './http-client.js' import type { SessionEntry, SessionStore } from './session-store.js' import type { ServerMessage, WsBridge } from './ws-bridge.js' +export type SessionRestoreResult = + | { status: 'restored'; session: SessionEntry } + | { status: 'missing' } + | { status: 'unavailable'; session: SessionEntry } + +export const SESSION_RECONNECT_NOTICE = '暂时无法连接原会话,已保留会话和工作目录,请稍后重试。' + type BridgeSessionOps = Pick< WsBridge, | 'connectSession' @@ -41,11 +48,11 @@ export async function restoreStoredSessionBinding({ onServerMessage, logPrefix, clearTransientState, -}: RestoreStoredSessionBindingOptions): Promise { +}: RestoreStoredSessionBindingOptions): Promise { const stored = sessionStore.get(chatId) if (!stored) { resetStaleBridge(chatId, bridge, clearTransientState) - return null + return { status: 'missing' } } const currentSessionId = bridge.getSessionId(chatId) @@ -54,7 +61,7 @@ export async function restoreStoredSessionBinding({ } if (bridge.isSessionOpen(chatId, stored.sessionId)) { - return stored + return { status: 'restored', session: stored } } let exists = true @@ -73,11 +80,23 @@ export async function restoreStoredSessionBinding({ const hadBridgeSession = bridge.hasSession(chatId) resetStaleBridge(chatId, bridge, clearTransientState) if (!hadBridgeSession) clearTransientState?.() - return null + return { status: 'missing' } } - bridge.connectSession(chatId, stored.sessionId) - bridge.onServerMessage(chatId, onServerMessage) - const opened = await bridge.waitForOpen(chatId) - return opened ? stored : null + // A transport failure is not evidence that the session was deleted. Keep + // its binding so the next inbound message can retry the same session. + try { + bridge.connectSession(chatId, stored.sessionId) + bridge.onServerMessage(chatId, onServerMessage) + if (await bridge.waitForOpen(chatId)) { + return { status: 'restored', session: stored } + } + } catch (err) { + console.warn( + `${logPrefix} Failed to reconnect stored session ${stored.sessionId}: ${ + err instanceof Error ? err.message : String(err) + }`, + ) + } + return { status: 'unavailable', session: stored } } diff --git a/adapters/dingtalk/index.ts b/adapters/dingtalk/index.ts index 181286a3..b7c4b975 100644 --- a/adapters/dingtalk/index.ts +++ b/adapters/dingtalk/index.ts @@ -28,7 +28,7 @@ import { formatProjectSelectionOutcome, ProjectSelectionController, } from '../common/project-selection-router.js' -import { restoreStoredSessionBinding } from '../common/session-recovery.js' +import { restoreStoredSessionBinding, SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from '../common/session-recovery.js' import { isAllowedUser, tryPair } from '../common/pairing.js' import { AttachmentStore } from '../common/attachment/attachment-store.js' import { checkAttachmentLimit } from '../common/attachment/attachment-limits.js' @@ -248,7 +248,7 @@ function clearPendingPermissions(chatId: string): void { pendingPermissions.delete(chatId) } -async function ensureExistingSession(chatId: string): Promise<{ sessionId: string; workDir: string } | null> { +async function ensureExistingSession(chatId: string): Promise { return await restoreStoredSessionBinding({ chatId, bridge, @@ -261,8 +261,10 @@ async function ensureExistingSession(chatId: string): Promise<{ sessionId: strin } async function buildStatusText(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (!stored) return formatImStatus(null) + const result = await ensureExistingSession(chatId) + if (result.status === 'unavailable') return SESSION_RECONNECT_NOTICE + if (result.status === 'missing') return formatImStatus(null) + const stored = result.session const runtime = getRuntimeState(chatId) let projectName = path.basename(stored.workDir) || stored.workDir @@ -312,8 +314,12 @@ async function buildStatusText(chatId: string): Promise { } async function ensureSession(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (stored) return true + const result = await ensureExistingSession(chatId) + if (result.status === 'restored') return true + if (result.status === 'unavailable') { + await sendText(chatId, SESSION_RECONNECT_NOTICE) + return false + } return await createSessionForChat(chatId, defaultWorkDir) } @@ -510,9 +516,9 @@ async function routeUserMessage(chatId: string, text: string, attachments: Attac return } if (!hasAttachments && (trimmed === '/clear' || trimmed === '清空')) { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } clearTransientChatState(chatId) @@ -525,9 +531,9 @@ async function routeUserMessage(chatId: string, text: string, attachments: Attac return } if (!hasAttachments && (trimmed === '/stop' || trimmed === '停止')) { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } bridge.sendStopGeneration(chatId) diff --git a/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts b/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts index da7e2bce..2ea87f61 100644 --- a/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts +++ b/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts @@ -187,6 +187,44 @@ async function send(platform: Platform, text: string, options: { unauthorized?: for (const platform of platforms) { describe(`${platform} actual module session selection`, () => { + it('keeps the original session and project after reconnect timeout and resumes it on retry', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + const originalProject = path.join(temporaryRoot, `${platform}-original`) + fs.mkdirSync(originalProject) + adapter.bridge.resetSession(chatId) + adapter.clearTransientChatState(chatId) + adapter.sessionStore.set(chatId, 'original-session', originalProject) + const originalBinding = adapter.sessionStore.get(chatId) + const creationsBefore = newSessionCount + const failedPrompt = `${platform}: Continue after reconnect` + const failedOpen = spyOn(adapter.bridge, 'waitForOpen').mockImplementationOnce(async () => { + adapter.bridge.resetSession(chatId) + return false + }) + const sendSpy = spyOn(adapter.bridge, 'sendUserMessage') + try { + await send(platform, failedPrompt) + expect(adapter.sessionStore.get(chatId)).toEqual(originalBinding) + expect(newSessionCount).toBe(creationsBefore) + expect(sentPrompts.some((prompt) => prompt.content === failedPrompt)).toBe(false) + expect(sendSpy).not.toHaveBeenCalled() + expect(notices.at(-1)).toContain('已保留会话和工作目录') + expect(notices.at(-1)).not.toContain('/new') + + await send(platform, 'Retry the original conversation') + expect(adapter.bridge.getSessionId(chatId)).toBe('original-session') + expect(adapter.sessionStore.get(chatId)).toEqual(originalBinding) + expect(newSessionCount).toBe(creationsBefore) + expect(sendSpy).toHaveBeenCalledWith(chatId, 'Retry the original conversation', undefined) + } finally { + failedOpen.mockRestore() + sendSpy.mockRestore() + adapter.bridge.resetSession(chatId) + adapter.clearTransientChatState(chatId) + } + }) + it('restores original history and honors the real initial active-turn snapshot', async () => { const adapter = adapterFor(platform) const chatId = chatFor(platform) diff --git a/adapters/feishu/index.ts b/adapters/feishu/index.ts index ba088149..a3c513c8 100644 --- a/adapters/feishu/index.ts +++ b/adapters/feishu/index.ts @@ -34,7 +34,7 @@ import { formatProjectSelectionOutcome, ProjectSelectionController, } from '../common/project-selection-router.js' -import { restoreStoredSessionBinding } from '../common/session-recovery.js' +import { restoreStoredSessionBinding, SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from '../common/session-recovery.js' import { isAllowedUser, tryPair } from '../common/pairing.js' import { extractInboundPayload } from './extract-payload.js' import { FeishuMediaService } from './media.js' @@ -220,7 +220,7 @@ function clearTransientChatState(chatId: string): void { pendingPermissions.delete(chatId) } -async function ensureExistingSession(chatId: string): Promise<{ sessionId: string; workDir: string } | null> { +async function ensureExistingSession(chatId: string): Promise { return await restoreStoredSessionBinding({ chatId, bridge, @@ -233,8 +233,10 @@ async function ensureExistingSession(chatId: string): Promise<{ sessionId: strin } async function buildStatusText(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (!stored) return formatImStatus(null) + const result = await ensureExistingSession(chatId) + if (result.status === 'unavailable') return SESSION_RECONNECT_NOTICE + if (result.status === 'missing') return formatImStatus(null) + const stored = result.session const runtime = getRuntimeState(chatId) let projectName = path.basename(stored.workDir) || stored.workDir @@ -647,8 +649,12 @@ function buildPermissionCard( // ---------- session management ---------- async function ensureSession(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (stored) return true + const result = await ensureExistingSession(chatId) + if (result.status === 'restored') return true + if (result.status === 'unavailable') { + await sendText(chatId, SESSION_RECONNECT_NOTICE) + return false + } const workDir = defaultWorkDir if (workDir) { @@ -1000,9 +1006,9 @@ async function handleMessage(data: any): Promise { return } if (!hasAttachments && (msgText === '/clear' || msgText === '清空')) { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } clearTransientChatState(chatId) @@ -1016,9 +1022,9 @@ async function handleMessage(data: any): Promise { return } if (!hasAttachments && (msgText === '/stop' || msgText === '停止')) { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } bridge.sendStopGeneration(chatId) diff --git a/adapters/telegram/__tests__/commands.test.ts b/adapters/telegram/__tests__/commands.test.ts index 9621b35e..cc57eb5b 100644 --- a/adapters/telegram/__tests__/commands.test.ts +++ b/adapters/telegram/__tests__/commands.test.ts @@ -131,7 +131,7 @@ function createController(overrides?: Record) { }, defaultWorkDir: '/work/repo', isAllowedUser: mock(() => true), - ensureExistingSession: mock(async () => ({ sessionId: 'active', workDir: '/work/repo' })), + ensureExistingSession: mock(async () => ({ status: 'restored' as const, session: { sessionId: 'active', workDir: '/work/repo', updatedAt: 1 } })), clearTransientChatState: mock((chatId: string) => bridgeEvents.push(`clear:${chatId}`)), clearOtherSelections: mock(() => {}), isBusy: mock(() => false), @@ -372,7 +372,7 @@ describe('Telegram command controller helpers', () => { it('reports empty lists and command failures without throwing', async () => { const { controller, sent } = createController({ defaultWorkDir: '', - ensureExistingSession: mock(async () => null), + ensureExistingSession: mock(async () => ({ status: 'missing' })), httpClient: { listProviders: mock(async () => { throw new Error('providers down') }), activateOfficialProvider: mock(async () => {}), @@ -516,7 +516,7 @@ describe('Telegram command controller helpers', () => { it('does not claim a selected skill ran when the agent session is unavailable', async () => { const unavailable = createController({ - ensureExistingSession: mock(async () => null), + ensureExistingSession: mock(async () => ({ status: 'missing' })), }) await unavailable.controller.handleSkillsCommand(createCommandContext().ctx) expect(unavailable.sent.at(-1)?.text).toContain('当前项目可用 Skills') @@ -532,6 +532,46 @@ describe('Telegram command controller helpers', () => { expect(callback.edits[0]).toContain('会话已失效') }) + it('keeps skill listing on the original project after a temporary reconnect failure', async () => { + const session = { sessionId: 'original', workDir: '/work/original', updatedAt: 1 } + const restore = mock(async () => ({ status: 'restored', session })) + .mockResolvedValueOnce({ status: 'unavailable', session }) + const { controller, deps, sent } = createController({ ensureExistingSession: restore }) + + await controller.handleSkillsCommand(createCommandContext().ctx) + expect(sent.at(-1)?.text).toContain('已保留会话和工作目录') + expect(sent.at(-1)?.text).not.toContain('/new') + expect(deps.httpClient.listSkills).not.toHaveBeenCalled() + expect(deps.setStoredSession).not.toHaveBeenCalled() + expect(deps.deleteStoredSession).not.toHaveBeenCalled() + + await controller.handleSkillsCommand(createCommandContext().ctx) + expect(deps.httpClient.listSkills).toHaveBeenCalledWith('/work/original') + expect(sent.at(-1)?.text).toContain('/work/original') + }) + + it('retains the skill selection for retry when the original session temporarily cannot reconnect', async () => { + const session = { sessionId: 'original', workDir: '/work/original', updatedAt: 1 } + const restore = mock(async () => ({ status: 'restored', session })) + const { controller, deps, sent, sentUserMessages } = createController({ ensureExistingSession: restore }) + await controller.handleSkillsCommand(createCommandContext().ctx) + restore.mockResolvedValueOnce({ status: 'unavailable', session }) + + const failed = createCommandContext() + await controller.handleSelectionCallback(failed.ctx, { kind: 'skill', action: 'pick', index: 0 }) + expect(sent.at(-1)?.text).toContain('已保留会话和工作目录') + expect(sent.at(-1)?.text).not.toContain('/new') + expect(failed.edits).toEqual([]) + expect(sentUserMessages).toEqual([]) + expect(deps.setStoredSession).not.toHaveBeenCalled() + expect(deps.deleteStoredSession).not.toHaveBeenCalled() + + const retry = createCommandContext() + await controller.handleSelectionCallback(retry.ctx, { kind: 'skill', action: 'pick', index: 0 }) + expect(sentUserMessages).toEqual([{ chatId: '42', content: '/skill-a' }]) + expect(retry.edits[0]).toContain('已调用 Skill') + }) + it('reports a disconnected bridge instead of dropping a selected skill', async () => { const disconnected = createController({ sendUserMessage: mock(() => false), @@ -827,7 +867,7 @@ describe('Telegram command controller helpers', () => { delete: (chatId) => events.push(`delete:${chatId}`), }, isAllowedUser: () => allowPermissionUser, - ensureExistingSession: mock(async () => ({ sessionId: 'active', workDir: '/work/repo' })), + ensureExistingSession: mock(async () => ({ status: 'restored' as const, session: { sessionId: 'active', workDir: '/work/repo', updatedAt: 1 } })), clearTransientChatState: (chatId) => events.push(`clear:${chatId}`), clearOtherSelections: () => {}, isBusy: () => false, diff --git a/adapters/telegram/__tests__/entrypoint-session-routing.test.ts b/adapters/telegram/__tests__/entrypoint-session-routing.test.ts index 7db606b6..f659ee99 100644 --- a/adapters/telegram/__tests__/entrypoint-session-routing.test.ts +++ b/adapters/telegram/__tests__/entrypoint-session-routing.test.ts @@ -4,6 +4,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import type { ServerWebSocket } from 'bun' import { SessionStore } from '../../common/session-store.js' +import { WsBridge } from '../../common/ws-bridge.js' import { AttachmentStore } from '../../common/attachment/attachment-store.js' // Import the actual entrypoint with isolated configuration. Telegram API calls @@ -149,6 +150,33 @@ describe('Telegram entrypoint session routing', () => { if (directory) rmSync(directory, { recursive: true, force: true }) }) + it('retains the original session and project through reconnect timeout and retries that conversation', async () => { + const chatId = 709 + store.set(String(chatId), 'history', worktree) + const originalBinding = store.get(String(chatId)) + const creationsBefore = requests.filter((request) => request === 'POST /api/sessions').length + const failedOpen = spyOn(WsBridge.prototype, 'waitForOpen').mockImplementationOnce(async function (this: WsBridge, id) { + this.resetSession(id) + return false + }) + try { + await text(chatId, 'Continue after reconnect') + expect(store.get(String(chatId))).toEqual(originalBinding) + expect(requests.filter((request) => request === 'POST /api/sessions').length).toBe(creationsBefore) + expect(messages.some((item) => item.message.content === 'Continue after reconnect')).toBe(false) + expect(texts(chatId).at(-1)).toContain('已保留会话和工作目录') + expect(texts(chatId).at(-1)).not.toContain('/new') + + await text(chatId, 'Retry the original conversation') + await eventually(() => expect(messages.some((item) => item.sessionId === 'history' && item.message.content === 'Retry the original conversation')).toBe(true)) + expect(store.get(String(chatId))).toEqual(originalBinding) + expect(requests.filter((request) => request === 'POST /api/sessions').length).toBe(creationsBefore) + broadcast('history', { type: 'message_complete' }) + } finally { + failedOpen.mockRestore() + } + }) + it('runs registered history commands after authorization and deduplication', async () => { const before = requests.length await text(701, '/sessions', { userId: 99 }) diff --git a/adapters/telegram/commands.ts b/adapters/telegram/commands.ts index 077d4e56..a7312010 100644 --- a/adapters/telegram/commands.ts +++ b/adapters/telegram/commands.ts @@ -1,5 +1,6 @@ import { formatImHelp } from '../common/format.js' import { listProjectSessionHistory, restoreSelectedSession } from '../common/session-selection.js' +import { SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from '../common/session-recovery.js' import type { SessionEntry } from '../common/session-store.js' import type { ServerMessage } from '../common/ws-bridge.js' import { @@ -70,7 +71,7 @@ export type TelegramCommandControllerDeps = { httpClient: AdapterHttpClient defaultWorkDir: string isAllowedUser: (userId: number) => boolean - ensureExistingSession: (chatId: string) => Promise<{ sessionId: string; workDir: string } | null> + ensureExistingSession: (chatId: string) => Promise clearTransientChatState: (chatId: string) => void clearOtherSelections: (chatId: string) => void isBusy: (chatId: string) => boolean @@ -240,7 +241,7 @@ export type TelegramRuntimeCommandControllerDeps = { delete: (chatId: string) => void } isAllowedUser: (userId: number) => boolean - ensureExistingSession: (chatId: string) => Promise<{ sessionId: string; workDir: string } | null> + ensureExistingSession: (chatId: string) => Promise clearTransientChatState: (chatId: string) => void clearOtherSelections: (chatId: string) => void isBusy: (chatId: string) => boolean @@ -527,8 +528,12 @@ export function createTelegramCommandController(deps: TelegramCommandControllerD } const showSkills = async (chatId: string): Promise => { - const stored = await deps.ensureExistingSession(chatId) - const cwd = stored?.workDir || deps.defaultWorkDir + const restored = await deps.ensureExistingSession(chatId) + if (restored.status === 'unavailable') { + await deps.api.sendMessage(Number(chatId), SESSION_RECONNECT_NOTICE) + return + } + const cwd = restored.status === 'restored' ? restored.session.workDir : deps.defaultWorkDir if (!cwd) { await deps.api.sendMessage(Number(chatId), '请先发送 /new 选择项目,再查看 Skills。') return @@ -565,8 +570,12 @@ export function createTelegramCommandController(deps: TelegramCommandControllerD const chatId = getCallbackChatId(ctx) if (!chatId) return - const stored = await deps.ensureExistingSession(chatId) - if (!stored) { + const restored = await deps.ensureExistingSession(chatId) + if (restored.status === 'unavailable') { + await deps.api.sendMessage(Number(chatId), SESSION_RECONNECT_NOTICE) + return + } + if (restored.status === 'missing') { pendingSelections.delete(chatId) await ctx.editMessageText('⚠️ 会话已失效,请发送 /new 重新选择项目后再调用 Skill。') return @@ -574,7 +583,7 @@ export function createTelegramCommandController(deps: TelegramCommandControllerD const invocation = `/${item.value}` if (!deps.sendUserMessage(chatId, invocation)) { - await ctx.editMessageText('⚠️ Skill 发送失败,连接可能已断开。请发送 /new 重新连接会话。') + await ctx.editMessageText(`⚠️ Skill 发送失败。${SESSION_RECONNECT_NOTICE}`) return } diff --git a/adapters/telegram/index.ts b/adapters/telegram/index.ts index 480ab24d..bed48a49 100644 --- a/adapters/telegram/index.ts +++ b/adapters/telegram/index.ts @@ -28,7 +28,7 @@ import { } from '../common/permission.js' import { SessionStore } from '../common/session-store.js' import { createAdapterClient } from '../common/adapter-client.js' -import { restoreStoredSessionBinding } from '../common/session-recovery.js' +import { restoreStoredSessionBinding, SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from '../common/session-recovery.js' import { SessionSelectionController } from '../common/session-selection.js' import { syncImPermissionState } from '../common/permission-sync.js' import { isAllowedUser, tryPair } from '../common/pairing.js' @@ -150,7 +150,7 @@ async function handlePermissionDecision(chatId: string, decision: PermissionDeci ) } -async function ensureExistingSession(chatId: string): Promise<{ sessionId: string; workDir: string } | null> { +async function ensureExistingSession(chatId: string): Promise { return await restoreStoredSessionBinding({ chatId, bridge, @@ -163,8 +163,10 @@ async function ensureExistingSession(chatId: string): Promise<{ sessionId: strin } async function buildStatusText(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (!stored) return formatImStatus(null) + const result = await ensureExistingSession(chatId) + if (result.status === 'unavailable') return SESSION_RECONNECT_NOTICE + if (result.status === 'missing') return formatImStatus(null) + const stored = result.session const runtime = getRuntimeState(chatId) let projectName = path.basename(stored.workDir) || stored.workDir @@ -216,8 +218,12 @@ async function buildStatusText(chatId: string): Promise { // ---------- session management ---------- async function ensureSession(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (stored) return true + const result = await ensureExistingSession(chatId) + if (result.status === 'restored') return true + if (result.status === 'unavailable') { + await bot.api.sendMessage(Number(chatId), SESSION_RECONNECT_NOTICE) + return false + } const workDir = defaultWorkDir if (workDir) { @@ -486,9 +492,9 @@ const isAuthorizedTelegramUser = (userId: number) => isAllowedUser('telegram', u registerAuthorizedTelegramCommand(bot, 'stop', isAuthorizedTelegramUser, (ctx) => { const chatId = String(ctx.chat!.id) void (async () => { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await ctx.reply(formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await ctx.reply(result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } bridge.sendStopGeneration(chatId) @@ -504,9 +510,9 @@ registerAuthorizedTelegramCommand(bot, 'status', isAuthorizedTelegramUser, async registerAuthorizedTelegramCommand(bot, 'clear', isAuthorizedTelegramUser, (ctx) => { const chatId = String(ctx.chat!.id) void (async () => { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await ctx.reply(formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await ctx.reply(result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } clearTransientChatState(chatId) diff --git a/adapters/wechat/index.ts b/adapters/wechat/index.ts index be7f219e..3bd796c1 100644 --- a/adapters/wechat/index.ts +++ b/adapters/wechat/index.ts @@ -19,7 +19,7 @@ import { SessionStore } from '../common/session-store.js' import { syncImPermissionState } from '../common/permission-sync.js' import { SessionSelectionController } from '../common/session-selection.js' import { createAdapterClient } from '../common/adapter-client.js' -import { restoreStoredSessionBinding } from '../common/session-recovery.js' +import { restoreStoredSessionBinding, SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from '../common/session-recovery.js' import { isAllowedUser, tryPair } from '../common/pairing.js' import { AttachmentStore } from '../common/attachment/attachment-store.js' import { checkAttachmentLimit } from '../common/attachment/attachment-limits.js' @@ -219,7 +219,7 @@ function enqueueWechat(chatId: string, task: () => Promise): void { }) } -async function ensureExistingSession(chatId: string): Promise<{ sessionId: string; workDir: string } | null> { +async function ensureExistingSession(chatId: string): Promise { return await restoreStoredSessionBinding({ chatId, bridge, @@ -232,8 +232,10 @@ async function ensureExistingSession(chatId: string): Promise<{ sessionId: strin } async function buildStatusText(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (!stored) return formatImStatus(null) + const result = await ensureExistingSession(chatId) + if (result.status === 'unavailable') return SESSION_RECONNECT_NOTICE + if (result.status === 'missing') return formatImStatus(null) + const stored = result.session const runtime = getRuntimeState(chatId) let projectName = path.basename(stored.workDir) || stored.workDir @@ -283,8 +285,12 @@ async function buildStatusText(chatId: string): Promise { } async function ensureSession(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (stored) return true + const result = await ensureExistingSession(chatId) + if (result.status === 'restored') return true + if (result.status === 'unavailable') { + await sendText(chatId, SESSION_RECONNECT_NOTICE) + return false + } const workDir = defaultWorkDir if (workDir) return await createSessionForChat(chatId, workDir) @@ -496,9 +502,9 @@ async function routeUserMessage(message: WechatMessage): Promise { return } if (!hasAttachments && (text === '/stop' || text === '停止')) { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } bridge.sendStopGeneration(chatId) @@ -506,9 +512,9 @@ async function routeUserMessage(message: WechatMessage): Promise { return } if (!hasAttachments && (text === '/clear' || text === '清空')) { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } clearTransientChatState(chatId) diff --git a/adapters/whatsapp/index.ts b/adapters/whatsapp/index.ts index 84ce03ae..5dc41a70 100644 --- a/adapters/whatsapp/index.ts +++ b/adapters/whatsapp/index.ts @@ -29,7 +29,7 @@ import { SessionStore } from '../common/session-store.js' import { syncImPermissionState } from '../common/permission-sync.js' import { SessionSelectionController } from '../common/session-selection.js' import { createAdapterClient } from '../common/adapter-client.js' -import { restoreStoredSessionBinding } from '../common/session-recovery.js' +import { restoreStoredSessionBinding, SESSION_RECONNECT_NOTICE, type SessionRestoreResult } from '../common/session-recovery.js' import { isAllowedUser, tryPair } from '../common/pairing.js' import { AttachmentStore } from '../common/attachment/attachment-store.js' import { checkAttachmentLimit } from '../common/attachment/attachment-limits.js' @@ -166,7 +166,7 @@ async function handlePermissionDecision(chatId: string, decision: PermissionDeci ) } -async function ensureExistingSession(chatId: string): Promise<{ sessionId: string; workDir: string } | null> { +async function ensureExistingSession(chatId: string): Promise { return await restoreStoredSessionBinding({ chatId, bridge, @@ -179,8 +179,10 @@ async function ensureExistingSession(chatId: string): Promise<{ sessionId: strin } async function buildStatusText(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (!stored) return formatImStatus(null) + const result = await ensureExistingSession(chatId) + if (result.status === 'unavailable') return SESSION_RECONNECT_NOTICE + if (result.status === 'missing') return formatImStatus(null) + const stored = result.session const runtime = getRuntimeState(chatId) let projectName = path.basename(stored.workDir) || stored.workDir @@ -206,8 +208,12 @@ async function buildStatusText(chatId: string): Promise { } async function ensureSession(chatId: string): Promise { - const stored = await ensureExistingSession(chatId) - if (stored) return true + const result = await ensureExistingSession(chatId) + if (result.status === 'restored') return true + if (result.status === 'unavailable') { + await sendWhatsAppText(chatId, SESSION_RECONNECT_NOTICE) + return false + } const workDir = defaultWorkDir if (workDir) { @@ -461,9 +467,9 @@ async function routeUserMessage( return } if (command === '/stop') { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendWhatsAppText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendWhatsAppText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } bridge.sendStopGeneration(chatId) @@ -475,9 +481,9 @@ async function routeUserMessage( return } if (command === '/clear') { - const stored = await ensureExistingSession(chatId) - if (!stored) { - await sendWhatsAppText(chatId, formatImStatus(null)) + const result = await ensureExistingSession(chatId) + if (result.status !== 'restored') { + await sendWhatsAppText(chatId, result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null)) return } clearTransientChatState(chatId)