From 59c7857beb56816c3770e522db67ed9362ada0e1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=A8=8B=E5=BA=8F=E5=91=98=E9=98=BF=E6=B1=9F=28Relakkes?= =?UTF-8?q?=29?= Date: Tue, 8 Sep 2026 13:56:16 +0800 Subject: [PATCH] feat(im): resume project session history (#1286) Add shared project and session selection across IM adapters, preserve bindings on failed restoration, and synchronize permissions with desktop clients. Fixes #1286 --- .../common/__tests__/chat-runtime.test.ts | 100 +++++- adapters/common/__tests__/http-client.test.ts | 53 ++- .../common/__tests__/permission-sync.test.ts | 70 ++++ .../__tests__/session-selection.test.ts | 215 ++++++++++++ adapters/common/__tests__/ws-bridge.test.ts | 78 +++++ adapters/common/chat-runtime.ts | 22 ++ adapters/common/format.ts | 7 +- adapters/common/http-client.ts | 12 +- adapters/common/permission-sync.ts | 46 +++ adapters/common/session-selection.ts | 262 ++++++++++++++ adapters/common/ws-bridge.ts | 12 +- adapters/dingtalk/index.ts | 34 +- ...egacy-session-selection-entrypoint.test.ts | 322 ++++++++++++++++++ adapters/feishu/index.ts | 46 ++- adapters/telegram/__tests__/commands.test.ts | 147 +++++++- .../entrypoint-session-routing.test.ts | 302 ++++++++++++++++ .../telegram/__tests__/session-input.test.ts | 63 ++++ adapters/telegram/commands.ts | 123 +++++-- adapters/telegram/index.ts | 130 ++++--- adapters/telegram/menu.ts | 1 + adapters/wechat/index.ts | 32 +- adapters/whatsapp/index.ts | 42 ++- docs/en/im/index.md | 14 + docs/im/index.md | 14 + 24 files changed, 2023 insertions(+), 124 deletions(-) create mode 100644 adapters/common/__tests__/permission-sync.test.ts create mode 100644 adapters/common/__tests__/session-selection.test.ts create mode 100644 adapters/common/permission-sync.ts create mode 100644 adapters/common/session-selection.ts create mode 100644 adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts create mode 100644 adapters/telegram/__tests__/entrypoint-session-routing.test.ts create mode 100644 adapters/telegram/__tests__/session-input.test.ts diff --git a/adapters/common/__tests__/chat-runtime.test.ts b/adapters/common/__tests__/chat-runtime.test.ts index 6c56e347..f80e95ec 100644 --- a/adapters/common/__tests__/chat-runtime.test.ts +++ b/adapters/common/__tests__/chat-runtime.test.ts @@ -16,7 +16,7 @@ import { ImChatRuntime, type ChatPort, type ResponseStream } from '../chat-runti import { loadConfig } from '../config.js' import { MessageDedup } from '../message-dedup.js' import { SessionStore } from '../session-store.js' -import type { AdapterHttpClient } from '../http-client.js' +import type { AdapterHttpClient, SessionListItem } from '../http-client.js' import type { AttachmentRef, ServerMessage, WsBridge } from '../ws-bridge.js' const CHAT_ID = 'chat-1' @@ -87,6 +87,7 @@ class FakeBridge { } class FakeHttpClient { + sessions: SessionListItem[] = [] createdSessions: string[] = [] existingSessions = new Set() projects: Array<{ projectName: string; realPath: string; branch?: string }> = [] @@ -107,6 +108,10 @@ class FakeHttpClient { return this.projects } + async listSessions() { + return { sessions: this.sessions, total: this.sessions.length } + } + async matchProject(query: string) { const project = this.projects.find((item) => item.projectName === query) return project ? { project } : {} @@ -609,3 +614,96 @@ describe('ImChatRuntime attachment authorization', () => { expect(bridge.sent.map((item) => item.content)).toEqual(['with attachment', 'plain text']) }) }) + +describe('ImChatRuntime historical sessions', () => { + function addHistory(httpClient: FakeHttpClient) { + httpClient.existingSessions.add('desktop-history') + httpClient.sessions = [{ + id: 'desktop-history', title: 'Continue desktop work', createdAt: '2026-08-01', + modifiedAt: '2026-08-31', messageCount: 5, projectPath: 'encoded', + workDir: tmpDir, workDirExists: true, + }] + httpClient.projects = [{ projectName: 'demo', realPath: tmpDir }] + } + + it('continues a desktop session through the real command pipeline and forwards the next message to that ID', async () => { + const { runtime, httpClient, bridge, sessionStore, notices } = createRuntime() + addHistory(httpClient) + await inbound(runtime, '/sessions demo') + expect(notices.at(-1)).toContain('Continue desktop work') + expect(bridge.sent).toEqual([]) + await inbound(runtime, '/resume 1') + expect(sessionStore.get(CHAT_ID)?.sessionId).toBe('desktop-history') + await inbound(runtime, '接着完成刚才的任务') + expect(bridge.getSessionId(CHAT_ID)).toBe('desktop-history') + expect(bridge.sent).toEqual([{ content: '接着完成刚才的任务', attachments: undefined }]) + expect(httpClient.createdSessions).toEqual([]) + }) + + it('refuses switching between sending a prompt and the first server status, then allows it after idle', async () => { + const { runtime, httpClient, bridge, sessionStore, notices } = createRuntime() + addHistory(httpClient) + await inbound(runtime, 'start') + const original = sessionStore.get(CHAT_ID)?.sessionId + await inbound(runtime, '/sessions') + await inbound(runtime, '/resume 1') + expect(notices.at(-1)).toContain('正在运行') + expect(sessionStore.get(CHAT_ID)?.sessionId).toBe(original) + await runtime.handleServerMessage(CHAT_ID, { type: 'status', state: 'idle' }) + await inbound(runtime, '/resume 1') + expect(bridge.getSessionId(CHAT_ID)).toBe('desktop-history') + }) + + it('recognizes a running desktop turn from the initial connection snapshot', async () => { + const { runtime, httpClient, sessionStore, notices } = createRuntime() + addHistory(httpClient) + await inbound(runtime, '/sessions demo') + await inbound(runtime, '/resume 1') + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_requests_snapshot', turnActive: true }) + expect(runtime.getRuntimeState(CHAT_ID).state).toBe('thinking') + await inbound(runtime, '/sessions') + await inbound(runtime, '/resume 1') + expect(notices.at(-1)).toContain('正在运行') + expect(sessionStore.get(CHAT_ID)?.sessionId).toBe('desktop-history') + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_request', requestId: 'pending', toolName: 'Bash', input: {} }) + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_requests_snapshot', turnActive: true }) + expect(runtime.getRuntimeState(CHAT_ID).state).toBe('permission_pending') + }) + + it('allows another history selection after a replayed permission is resolved on Desktop', async () => { + const { runtime, httpClient, notices } = createRuntime() + addHistory(httpClient) + await inbound(runtime, '/sessions demo') + await inbound(runtime, '/resume 1') + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_request', requestId: 'desktop-request', toolName: 'Bash', input: {} }) + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_requests_snapshot', toolRequestIds: ['desktop-request'], turnActive: true }) + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_resolved', requestId: 'desktop-request', permissionType: 'tool', allowed: true }) + await runtime.handleServerMessage(CHAT_ID, { type: 'status', state: 'idle' }) + expect(runtime.getRuntimeState(CHAT_ID).pendingPermissionCount).toBe(0) + await inbound(runtime, '/sessions') + await inbound(runtime, '/resume 1') + expect(notices.at(-1)).toContain('已恢复会话') + }) + + it('keeps permission answers ahead of session numbers and cancels the picker for project switching', async () => { + const { runtime, httpClient, bridge, sessionStore } = createRuntime() + addHistory(httpClient) + await inbound(runtime, '/sessions demo') + await runtime.handleServerMessage(CHAT_ID, { type: 'permission_request', requestId: 'req-1', toolName: 'Bash', input: {} }) + await inbound(runtime, '1') + expect(bridge.permissionResponses).toHaveLength(1) + expect(sessionStore.get(CHAT_ID)).toBeNull() + await inbound(runtime, '/projects') + await inbound(runtime, '/resume 1') + expect(sessionStore.get(CHAT_ID)).toBeNull() + }) + + it('authorizes history commands and does not interpret attachment captions as commands', async () => { + const { runtime, httpClient, notices, bridge } = createRuntime() + addHistory(httpClient) + await inbound(runtime, '/sessions demo', { userId: 'stranger' }) + expect(notices.at(-1)).toContain('未授权') + await inbound(runtime, '/sessions demo', { hasAttachments: true, loadAttachments: async () => [] }) + expect(bridge.sent.at(-1)?.content).toBe('/sessions demo') + }) +}) diff --git a/adapters/common/__tests__/http-client.test.ts b/adapters/common/__tests__/http-client.test.ts index 13f647ac..82fa2358 100644 --- a/adapters/common/__tests__/http-client.test.ts +++ b/adapters/common/__tests__/http-client.test.ts @@ -219,6 +219,33 @@ describe('AdapterHttpClient', () => { ) }) + it('sessionExists accepts the server detail shape for an existing desktop worktree session', async () => { + const rootDir = fs.mkdtempSync(path.join(os.tmpdir(), 'im-root-')) + const workDir = path.join(rootDir, '.claude', 'worktrees', 'existing-task') + fs.mkdirSync(workDir, { recursive: true }) + try { + client = new AdapterHttpClient('ws://127.0.0.1:3456', { allowedProjectRoots: [rootDir] }) + globalThis.fetch = mock(() => Promise.resolve(Response.json({ + id: 'desktop-session', + title: 'Existing desktop task', + projectRoot: rootDir, + workDir, + workDirExists: true, + permissionMode: 'default', + messages: [{ role: 'user', content: 'Previous task context' }], + }))) as any + + await expect(client.sessionExists('desktop-session')).resolves.toBe(true) + + // A removed worktree must not be resumed in a different directory merely + // because its logical project still exists or a cached flag says it does. + fs.rmSync(workDir, { recursive: true, force: true }) + await expect(client.sessionExists('desktop-session')).resolves.toBe(false) + } finally { + fs.rmSync(rootDir, { recursive: true, force: true }) + } + }) + it('sessionExists rejects sessions outside the root or using bypassPermissions', async () => { const rootDir = fs.mkdtempSync(path.join(os.tmpdir(), 'im-root-')) const outsideDir = fs.mkdtempSync(path.join(os.tmpdir(), 'outside-')) @@ -227,10 +254,8 @@ describe('AdapterHttpClient', () => { globalThis.fetch = mock((url: string) => { const unsafe = url.endsWith('/outside') return Promise.resolve(Response.json({ - status: { - workDir: unsafe ? outsideDir : rootDir, - permissionMode: unsafe ? 'default' : 'bypassPermissions', - }, + workDir: unsafe ? outsideDir : rootDir, + permissionMode: unsafe ? 'default' : 'bypassPermissions', })) }) as any @@ -322,7 +347,7 @@ describe('AdapterHttpClient', () => { const result = await client.listSessions({ project: rootDir, limit: 10, offset: 5 }) expect(result.sessions.map((session) => session.id)).toEqual(['session-1']) - expect(result.total).toBe(1) + expect(result.total).toBe(2) expect((globalThis.fetch as any).mock.calls[0][0]).toBe( `http://127.0.0.1:3456/api/sessions?project=${encodeURIComponent(rootDir)}&limit=10&offset=5`, ) @@ -331,6 +356,24 @@ describe('AdapterHttpClient', () => { } }) + it('preserves server pagination when an entire page is filtered out', async () => { + const rootDir = fs.mkdtempSync(path.join(os.tmpdir(), 'im-root-')) + try { + client = new AdapterHttpClient('ws://127.0.0.1:3456', { allowedProjectRoots: [rootDir] }) + globalThis.fetch = mock(() => Promise.resolve(Response.json({ + sessions: [{ id: 'unsafe-session', workDir: rootDir, permissionMode: 'bypassPermissions' }], + total: 6, + }))) as any + + await expect(client.listSessions({ limit: 1, offset: 0 })).resolves.toEqual({ + sessions: [], + total: 6, + }) + } finally { + fs.rmSync(rootDir, { recursive: true, force: true }) + } + }) + it('lists and activates providers through the server provider API', async () => { globalThis.fetch = mock((url: string, init?: RequestInit) => { if (url.endsWith('/api/providers') && !init?.method) { diff --git a/adapters/common/__tests__/permission-sync.test.ts b/adapters/common/__tests__/permission-sync.test.ts new file mode 100644 index 00000000..5fea279b --- /dev/null +++ b/adapters/common/__tests__/permission-sync.test.ts @@ -0,0 +1,70 @@ +import { describe, expect, it } from 'bun:test' +import { syncImPermissionState } from '../permission-sync.js' + +function harness() { + const runtime = { state: 'permission_pending' as 'idle' | 'thinking' | 'streaming' | 'tool_executing' | 'permission_pending', pendingPermissionCount: 2 } + const pending = new Map([['chat', new Set(['first', 'second'])], ['other', new Set(['unrelated'])]]) + return { runtime, pending } +} + +describe('IM approval reconciliation for restored sessions (#1286)', () => { + it('removes approvals answered in Desktop, without decrementing twice', () => { + const { runtime, pending } = harness() + const resolved = { type: 'permission_resolved', permissionType: 'tool', requestId: 'first' } + expect(syncImPermissionState('chat', resolved, runtime, pending)).toBe(true) + expect([...pending.get('chat')!]).toEqual(['second']) + expect(runtime).toEqual({ state: 'permission_pending', pendingPermissionCount: 1 }) + syncImPermissionState('chat', resolved, runtime, pending) + expect(runtime.pendingPermissionCount).toBe(1) + + syncImPermissionState('chat', { ...resolved, requestId: 'second' }, runtime, pending) + expect(runtime).toEqual({ state: 'thinking', pendingPermissionCount: 0 }) + expect(pending.has('chat')).toBe(false) + expect([...pending.get('other')!]).toEqual(['unrelated']) + }) + + it('keeps tool approvals intact when computer-use approval IDs overlap', () => { + const { runtime, pending } = harness() + syncImPermissionState('chat', { type: 'permission_resolved', permissionType: 'computer_use', requestId: 'first' }, runtime, pending) + expect([...pending.get('chat')!]).toEqual(['first', 'second']) + expect(runtime.pendingPermissionCount).toBe(2) + }) + + it('uses the reconnect snapshot to remove stale approvals and deduplicate replay', () => { + const { runtime, pending } = harness() + runtime.pendingPermissionCount = 9 + syncImPermissionState('chat', { type: 'permission_requests_snapshot', toolRequestIds: ['second', 'second', 'new'], turnActive: true }, runtime, pending) + expect([...pending.get('chat')!]).toEqual(['second', 'new']) + expect(runtime).toEqual({ state: 'permission_pending', pendingPermissionCount: 2 }) + }) + + it('unblocks switching after a disconnected approval was resolved and the turn finished', () => { + const { runtime, pending } = harness() + syncImPermissionState('chat', { type: 'permission_requests_snapshot', toolRequestIds: [], turnActive: false }, runtime, pending) + expect(runtime).toEqual({ state: 'idle', pendingPermissionCount: 0 }) + expect(pending.has('chat')).toBe(false) + }) + + it('keeps a continuing turn busy after its last approval disappears', () => { + const { runtime, pending } = harness() + syncImPermissionState('chat', { type: 'permission_requests_snapshot', toolRequestIds: [], turnActive: true }, runtime, pending) + expect(runtime).toEqual({ state: 'thinking', pendingPermissionCount: 0 }) + runtime.state = 'streaming' + syncImPermissionState('chat', { type: 'permission_requests_snapshot', toolRequestIds: [], turnActive: true }, runtime, pending) + expect(runtime.state).toBe('streaming') + }) + + it('preserves approvals when an older snapshot only reports turn activity', () => { + const { runtime, pending } = harness() + syncImPermissionState('chat', { type: 'permission_requests_snapshot', turnActive: true }, runtime, pending) + expect(runtime).toEqual({ state: 'permission_pending', pendingPermissionCount: 2 }) + }) + + it('leaves unrelated event routing and late duplicate resolutions unchanged', () => { + const runtime = { state: 'idle' as const, pendingPermissionCount: 0 } + const pending = new Map>() + expect(syncImPermissionState('chat', { type: 'content_delta', text: 'hello' }, runtime, pending)).toBe(false) + syncImPermissionState('chat', { type: 'permission_resolved', permissionType: 'tool', requestId: 'already-done' }, runtime, pending) + expect(runtime).toEqual({ state: 'idle', pendingPermissionCount: 0 }) + }) +}) diff --git a/adapters/common/__tests__/session-selection.test.ts b/adapters/common/__tests__/session-selection.test.ts new file mode 100644 index 00000000..f17ed3de --- /dev/null +++ b/adapters/common/__tests__/session-selection.test.ts @@ -0,0 +1,215 @@ +import { afterEach, beforeEach, describe, expect, it } from 'bun:test' +import * as fs from 'node:fs' +import * as os from 'node:os' +import * as path from 'node:path' +import { SessionSelectionController, SESSION_SELECTION_TTL_MS, listProjectSessionHistory } from '../session-selection.js' +import type { RecentProject, SessionListItem } from '../http-client.js' +import { SessionStore } from '../session-store.js' +import type { ServerMessage } from '../ws-bridge.js' + +let tmp: string +beforeEach(() => { tmp = fs.mkdtempSync(path.join(os.tmpdir(), 'im-session-picker-')) }) +afterEach(() => { fs.rmSync(tmp, { recursive: true, force: true }) }) + +function session(id: string, workDir: string, extra: Partial = {}): SessionListItem { + return { id, title: `历史 ${id}`, createdAt: '2026-08-01', modifiedAt: '2026-08-31', messageCount: 3, projectPath: 'encoded', workDir, workDirExists: true, ...extra } +} + +function project(name: string, realPath: string): RecentProject { + return { projectName: name, realPath, projectPath: realPath, isGit: true, repoName: name, branch: 'main', modifiedAt: '2026-08-31', sessionCount: 3 } +} + +function harness() { + const store = new SessionStore(path.join(tmp, 'sessions.json')) + store.set('chat', 'current', tmp) + const notices: string[] = [] + const connections: string[] = [] + const resets: string[] = [] + const events: ServerMessage[] = [] + const fetches: number[] = [] + const state = { + sessions: [session('old', tmp)], projects: [project('demo', tmp)], + exists: true, open: true, busy: false, now: 0, clears: 0, projectClears: 0, + preflightError: false, fetchError: false, busyDuringPreflight: false, + handler: undefined as ((message: ServerMessage) => void) | undefined, + } + const controller = new SessionSelectionController({ + sessionStore: store, + httpClient: { + async listSessions(options) { + if (state.fetchError) throw new Error('offline') + const offset = options?.offset ?? 0 + fetches.push(offset) + return { sessions: state.sessions.slice(offset, offset + (options?.limit ?? 20)), total: state.sessions.length } + }, + async listRecentProjects() { return state.projects }, + async matchProject(query) { + const matches = state.projects.filter((p) => p.projectName === query || p.realPath === query) + return matches.length === 1 ? { project: matches[0] } : matches.length ? { ambiguous: matches } : {} + }, + async sessionExists() { + if (state.preflightError) throw new Error('offline') + if (state.busyDuringPreflight) state.busy = true + return state.exists + }, + }, + bridge: { + resetSession(chatId) { resets.push(chatId) }, + connectSession(_chatId, id) { connections.push(id); return true }, + onServerMessage(_chatId, handler) { state.handler = handler }, + async waitForOpen() { state.handler?.({ type: 'status', marker: 'restored' }); return state.open }, + }, + async sendNotice(_chatId, text) { notices.push(text) }, + onServerMessage(_chatId, message) { events.push(message) }, + clearTransientState() { state.clears++ }, + clearProjectSelection() { state.projectClears++ }, + isBusy() { return state.busy }, + now: () => state.now, + }) + return { controller, store, notices, connections, resets, events, fetches, state } +} + +describe('IM history selection (#1286)', () => { + it('lists current project history without touching its binding, then persists the selected original ID', async () => { + const h = harness() + await h.controller.handleInput('chat', '/sessions') + expect(h.notices.at(-1)).toContain('历史 old') + expect(h.store.get('chat')?.sessionId).toBe('current') + expect(h.resets).toEqual([]) + await h.controller.handleInput('chat', '/resume 1') + expect(h.connections).toEqual(['old']) + expect(h.state.clears).toBe(1) + expect(h.events).toEqual([{ type: 'status', marker: 'restored' }]) + expect(new SessionStore(path.join(tmp, 'sessions.json')).get('chat')?.sessionId).toBe('old') + expect(h.notices.at(-1)).toContain('可以继续发送消息') + }) + + it('starts with a project picker when no binding survives an upgrade', async () => { + const h = harness() + h.store.delete('chat') + await h.controller.handleInput('chat', '/sessions') + expect(h.notices.at(-1)).toContain('历史会话所在的项目') + await h.controller.handleInput('chat', '1') + expect(h.notices.at(-1)).toContain('历史 old') + await h.controller.handleInput('chat', '1') + expect(h.store.get('chat')?.sessionId).toBe('old') + }) + + it('includes old desktop worktree sessions under their logical project, across server pages', async () => { + const h = harness() + h.state.sessions = Array.from({ length: 100 }, (_, i) => session(`other-${i}`, '/another')) + h.state.sessions.push(session('worktree-old', path.join(tmp, 'worktrees', 'old'), { projectRoot: tmp })) + await h.controller.handleInput('chat', '/sessions demo') + expect(h.fetches).toEqual([0, 100]) + expect(h.notices.at(-1)).toContain('worktree-old') + expect(h.notices.at(-1)).not.toContain('other-') + await h.controller.handleInput('chat', '/resume 1') + expect(h.store.get('chat')?.workDir).toBe(path.join(tmp, 'worktrees', 'old')) + }) + + it('uses the active worktree logical root for /sessions and identifies the active session', async () => { + const h = harness() + const worktree = path.join(tmp, 'worktrees', 'active') + h.store.set('chat', 'current', worktree) + h.state.sessions.push(session('current', worktree, { projectRoot: tmp })) + await h.controller.handleInput('chat', '会话列表') + expect(h.notices.at(-1)).toContain('历史 current(当前)') + expect(h.notices.at(-1)).toContain('历史 old') + }) + + it('accepts an exact worktree directory as well as its logical project root', async () => { + const worktree = path.join(tmp, 'worktrees', 'active') + const all = [session('main', tmp), session('branch', worktree, { projectRoot: tmp })] + const client = { async listSessions() { return { sessions: all, total: all.length } } } + expect((await listProjectSessionHistory(client, worktree)).map((s) => s.id)).toEqual(['branch']) + expect((await listProjectSessionHistory(client, tmp)).map((s) => s.id)).toEqual(['branch', 'main']) + }) + + it('keeps numbered options stable when the server order changes and pages beyond the first eight', async () => { + const h = harness() + h.state.sessions = Array.from({ length: 12 }, (_, i) => session(`id-${String(i).padStart(2, '0')}`, tmp)) + await h.controller.handleInput('chat', '/sessions') + h.state.sessions.reverse() + await h.controller.handleInput('chat', '/sessions next') + expect(h.notices.at(-1)).toContain('9. 历史 id-08') + await h.controller.handleInput('chat', '/resume 1') + expect(h.notices.at(-1)).toContain('编号无效') + await h.controller.handleInput('chat', '/resume 9') + expect(h.store.get('chat')?.sessionId).toBe('id-08') + }) + + it('lets a bound chat select a different project without creating a new session', async () => { + const h = harness() + h.state.projects = [project('demo', tmp), project('demo', path.join(tmp, 'other'))] + h.state.sessions.push(session('second', path.join(tmp, 'other'))) + await h.controller.handleInput('chat', '/sessions demo') + expect(h.notices.at(-1)).toContain('2. demo') + await h.controller.handleInput('chat', '2') + await h.controller.handleInput('chat', '1') + expect(h.store.get('chat')?.sessionId).toBe('second') + await h.controller.handleInput('chat', '/sessions projects') + expect(h.notices.at(-1)).toContain('历史会话所在的项目') + }) + + for (const fault of ['exists', 'open', 'busy', 'preflightError', 'busyDuringPreflight'] as const) { + it(`preserves the previous binding on ${fault} failure`, async () => { + const h = harness() + await h.controller.handleInput('chat', '/sessions') + h.state[fault] = !['exists', 'open'].includes(fault) + await h.controller.handleInput('chat', '/resume 1') + expect(h.store.get('chat')?.sessionId).toBe('current') + expect(h.state.clears).toBe(0) + expect(h.events).toEqual([]) + expect(h.connections).toEqual(fault === 'open' ? ['old', 'current'] : []) + expect(h.notices.at(-1)).not.toContain('已恢复会话') + }) + } + + it('does not create a binding when the first attempted connection fails', async () => { + const h = harness() + h.store.delete('chat') + await h.controller.handleInput('chat', '/sessions demo') + h.state.open = false + await h.controller.handleInput('chat', '/resume 1') + expect(h.store.get('chat')).toBeNull() + expect(h.resets).toHaveLength(2) + }) + + it('does not consume ordinary chat, permission commands, or another chat numeric replies', async () => { + const h = harness() + await h.controller.handleInput('chat', '/sessions') + expect(await h.controller.handleInput('other', '1')).toBe(false) + expect(await h.controller.handleInput('chat', '/allow request')).toBe(false) + expect(await h.controller.handleInput('chat', '继续修复刚才的代码')).toBe(false) + await h.controller.handleInput('chat', '/resume 1') + expect(h.notices.at(-1)).toContain('不存在或已过期') + expect(h.connections).toEqual([]) + }) + + it('rejects expired selection and lets users explicitly cancel', async () => { + const h = harness() + await h.controller.handleInput('chat', '/sessions') + h.state.now += SESSION_SELECTION_TTL_MS + expect(await h.controller.handleInput('chat', '1')).toBe(true) + expect(h.notices.at(-1)).toContain('已过期') + await h.controller.handleInput('chat', '/sessions') + await h.controller.handleInput('chat', '/cancel') + expect(h.notices.at(-1)).toContain('当前会话未改变') + expect(await h.controller.handleInput('chat', '1')).toBe(false) + expect(h.connections).toEqual([]) + }) + + it('reports missing projects, empty history, and HTTP errors without losing the active session', async () => { + const h = harness() + await h.controller.handleInput('chat', '/sessions unknown') + expect(h.notices.at(-1)).toContain('未找到项目') + h.state.sessions = [session('gone', tmp, { workDirExists: false })] + await h.controller.handleInput('chat', '/sessions') + expect(h.notices.at(-1)).toContain('没有可恢复会话') + h.state.fetchError = true + await h.controller.handleInput('chat', '/sessions') + expect(h.notices.at(-1)).toContain('offline') + expect(h.store.get('chat')?.sessionId).toBe('current') + expect(h.resets).toEqual([]) + }) +}) diff --git a/adapters/common/__tests__/ws-bridge.test.ts b/adapters/common/__tests__/ws-bridge.test.ts index d859fbf3..80a0071a 100644 --- a/adapters/common/__tests__/ws-bridge.test.ts +++ b/adapters/common/__tests__/ws-bridge.test.ts @@ -1,5 +1,6 @@ import { describe, it, expect, beforeEach, afterEach } from 'bun:test' import { WsBridge } from '../ws-bridge.js' +import { restoreSelectedSession } from '../session-selection.js' import { WebSocketServer, type WebSocket as WsServerSocket } from 'ws' async function waitFor( @@ -266,6 +267,83 @@ describe('WsBridge: handler serialization', () => { bridge.destroy() }) + it('drops old permission and status messages already queued when a session is replaced', async () => { + const bridge = new WsBridge(serverUrl, 'test', '') + const events: string[] = [] + let release!: () => void + let started!: () => void + const gate = new Promise((resolve) => { release = resolve }) + const firstStarted = new Promise((resolve) => { started = resolve }) + try { + bridge.onServerMessage('chat-queued', async (message) => { + if (message.tag === 'first') { + started() + await gate + } + events.push(message.tag) + }) + bridge.connectSession('chat-queued', 'old') + expect(await bridge.waitForOpen('chat-queued')).toBe(true) + const oldServerSocket = await waitForServerConnection() + oldServerSocket.send(JSON.stringify({ type: 'message_complete', tag: 'first' })) + await firstStarted + + const oldClientSocket = (bridge as any).sessions.get('chat-queued').ws + let received = 0 + const queued = new Promise((resolve) => { + oldClientSocket.on('message', () => { + received++ + if (received === 2) resolve() + }) + }) + oldServerSocket.send(JSON.stringify({ type: 'permission_request', tag: 'old-permission' })) + oldServerSocket.send(JSON.stringify({ type: 'status', tag: 'old-status' })) + await queued + const oldChain = (bridge as any).handlerChains.get('chat-queued') as Promise + + bridge.resetSession('chat-queued') + bridge.onServerMessage('chat-queued', (message) => { events.push(message.tag) }) + bridge.connectSession('chat-queued', 'selected') + expect(await bridge.waitForOpen('chat-queued')).toBe(true) + connections[1]!.send(JSON.stringify({ type: 'status', tag: 'selected-status' })) + expect(await waitFor(() => events.includes('selected-status'))).toBe(true) + release() + await oldChain + + // An already-running callback may complete; its queued successors must + // not restore old permissions or overwrite the selected session state. + expect(events).toEqual(['selected-status', 'first']) + } finally { + release() + bridge.destroy() + } + }) + + it('delivers immediate connection state after history restoration replaces its buffering handler', async () => { + server.once('connection', (socket) => { + socket.send(JSON.stringify({ type: 'status', state: 'permission_pending' })) + socket.send(JSON.stringify({ type: 'permission_request', requestId: 'selected-permission' })) + }) + const bridge = new WsBridge(serverUrl, 'test', '') + const events: string[] = [] + let savedId: string | undefined + try { + const result = await restoreSelectedSession({ + bridge, + httpClient: { sessionExists: async () => true }, + sessionStore: { get: () => null, set(_chatId, id) { savedId = id }, delete() {} }, + clearTransientState() { events.push('clear') }, + onServerMessage(_chatId, message) { events.push(message.type) }, + }, 'chat-immediate', { id: 'selected', title: 'Selected session', workDir: '/fixture/project' }) + expect(result.ok).toBe(true) + expect(savedId).toBe('selected') + expect(await waitFor(() => events.length === 3)).toBe(true) + expect(events).toEqual(['clear', 'status', 'permission_request']) + } finally { + bridge.destroy() + } + }) + it('does not dispatch stale messages from a socket reset before reconnect', async () => { const bridge = new WsBridge(serverUrl, 'test') const events: string[] = [] diff --git a/adapters/common/chat-runtime.ts b/adapters/common/chat-runtime.ts index 0990be99..3db366d6 100644 --- a/adapters/common/chat-runtime.ts +++ b/adapters/common/chat-runtime.ts @@ -35,6 +35,8 @@ import { type PermissionDecision, } from './permission.js' import { restoreStoredSessionBinding } from './session-recovery.js' +import { SessionSelectionController } from './session-selection.js' +import { syncImPermissionState } from './permission-sync.js' import type { SessionStore } from './session-store.js' import type { AttachmentRef, ServerMessage, WsBridge } from './ws-bridge.js' import { ImageBlockWatcher } from './attachment/image-block-watcher.js' @@ -148,6 +150,7 @@ export class ImChatRuntime { private readonly pendingPermissions = new Map>() private readonly pendingProjectSelection = new Set() private readonly imageWatchers = new Map() + private readonly sessionSelection: SessionSelectionController constructor(options: ImChatRuntimeOptions) { this.port = options.port @@ -160,6 +163,17 @@ export class ImChatRuntime { this.dedup = options.dedup this.flushIntervalMs = options.flushIntervalMs ?? 800 this.flushCharThreshold = options.flushCharThreshold ?? 320 + this.sessionSelection = new SessionSelectionController({ + httpClient: this.httpClient, + bridge: this.bridge, + sessionStore: this.sessionStore, + sendNotice: (chatId, text) => this.port.sendNotice(chatId, text), + onServerMessage: (chatId, message) => this.handleServerMessage(chatId, message), + clearTransientState: (chatId) => this.clearTransientChatState(chatId), + clearProjectSelection: (chatId) => { this.pendingProjectSelection.delete(chatId) }, + isBusy: (chatId) => this.getRuntimeState(chatId).state !== 'idle' + || Boolean(this.pendingPermissions.get(chatId)?.size), + }) } /** The work dir a `/new` with no argument would use. Exposed for adapters @@ -230,6 +244,8 @@ export class ImChatRuntime { return } + if (!hasAttachments && await this.sessionSelection.handleInput(chatId, text)) return + if (!hasAttachments && this.pendingProjectSelection.has(chatId)) { if (text) await this.startNewSession(chatId, text) return @@ -247,6 +263,7 @@ export class ImChatRuntime { if (!effective && attachments.length === 0) return this.port.setBusy?.(chatId, true) + this.getRuntimeState(chatId).state = 'thinking' const sent = this.bridge.sendUserMessage( chatId, effective, @@ -254,6 +271,7 @@ export class ImChatRuntime { ) if (!sent) { this.port.setBusy?.(chatId, false) + this.getRuntimeState(chatId).state = 'idle' await this.port.sendNotice(chatId, '消息发送失败,连接可能已断开。请发送 /new 重新开始。') } } @@ -299,6 +317,7 @@ export class ImChatRuntime { } this.clearTransientChatState(chatId) const sent = this.bridge.sendUserMessage(chatId, '/clear') + if (sent) this.getRuntimeState(chatId).state = 'thinking' await this.port.sendNotice( chatId, sent ? '已清空当前会话上下文。' : '无法发送 /clear,请先发送 /new 重新连接会话。', @@ -338,6 +357,7 @@ export class ImChatRuntime { async handleServerMessage(chatId: string, msg: ServerMessage): Promise { const runtime = this.getRuntimeState(chatId) + if (syncImPermissionState(chatId, msg, runtime, this.pendingPermissions)) return switch (msg.type) { case 'connected': @@ -531,6 +551,7 @@ export class ImChatRuntime { } async showProjectPicker(chatId: string): Promise { + this.sessionSelection.clear(chatId) try { const projects = await this.httpClient.listRecentProjects() if (projects.length === 0) { @@ -557,6 +578,7 @@ export class ImChatRuntime { } async startNewSession(chatId: string, query?: string): Promise { + this.sessionSelection.clear(chatId) this.bridge.resetSession(chatId) this.sessionStore.delete(chatId) this.clearTransientChatState(chatId) diff --git a/adapters/common/format.ts b/adapters/common/format.ts index fc1227c2..b0cd2eb8 100644 --- a/adapters/common/format.ts +++ b/adapters/common/format.ts @@ -28,6 +28,11 @@ type ImStatusSummary = { const IM_HELP_LINES = [ '/new [项目] / 新会话 — 新建会话或切换项目', '/projects / 项目列表 — 查看最近项目', + '/sessions [项目] / 会话列表 — 查看当前项目的旧会话', + '/sessions projects — 选择其他项目的旧会话', + '/resume <编号> / 继续会话 <编号> — 恢复列表中的会话', + '/sessions next / /sessions prev — 翻页', + '/cancel / 取消选择 — 退出选择,保留当前会话', '/status / 状态 — 查看当前会话状态', '/clear / 清空 — 清空当前会话上下文', '/stop / 停止 — 停止当前生成', @@ -316,7 +321,7 @@ export function formatImHelp(): string { export function formatImStatus(summary: ImStatusSummary | null): string { if (!summary?.sessionId) { - return '当前没有活动会话。\n\n发送 /new 新建会话,或发送 /projects 选择项目。' + return '当前没有活动会话。\n\n发送 /sessions 选择旧会话,/new 新建会话,或 /projects 选择项目。' } const lines = ['当前会话状态:'] diff --git a/adapters/common/http-client.ts b/adapters/common/http-client.ts index 52f5ca85..30f4a93c 100644 --- a/adapters/common/http-client.ts +++ b/adapters/common/http-client.ts @@ -151,12 +151,10 @@ export class AdapterHttpClient { throw new Error(`Failed to check session: ${(err as any).message}`) } const data = (await res.json()) as { - status?: { - workDir?: string - permissionMode?: string - } + workDir?: string + permissionMode?: string } - return this.isSafeRemoteSession(data.status) + return this.isSafeRemoteSession(data) } finally { clearTimeout(timer) } @@ -295,7 +293,9 @@ export class AdapterHttpClient { workDir: session.workDir, permissionMode: session.permissionMode, })) - return { sessions, total: sessions.length } + // total counts server candidates before the local safety filter. Keep it + // for offset pagination, including pages with no locally allowed sessions. + return { sessions, total: data.total } } finally { clearTimeout(timer) } diff --git a/adapters/common/permission-sync.ts b/adapters/common/permission-sync.ts new file mode 100644 index 00000000..1efb9e40 --- /dev/null +++ b/adapters/common/permission-sync.ts @@ -0,0 +1,46 @@ +import type { ServerMessage } from './ws-bridge.js' + +type PermissionRuntimeState = { + state: 'idle' | 'thinking' | 'streaming' | 'tool_executing' | 'permission_pending' + pendingPermissionCount: number +} + +/** Reconcile approvals handled by another client or while this IM was offline. */ +export function syncImPermissionState( + chatId: string, + message: ServerMessage, + runtime: PermissionRuntimeState, + pendingPermissions: Map>, +): boolean { + if (message.type === 'permission_resolved') { + if (message.permissionType !== 'tool' || typeof message.requestId !== 'string') return true + const pending = pendingPermissions.get(chatId) + pending?.delete(message.requestId) + runtime.pendingPermissionCount = pending?.size ?? 0 + if (!runtime.pendingPermissionCount) { + pendingPermissions.delete(chatId) + if (runtime.state === 'permission_pending') runtime.state = 'thinking' + } + return true + } + + if (message.type !== 'permission_requests_snapshot') return false + // Only an explicit array is authoritative. Older/partial snapshots may carry + // turnActive alone and must not silently discard a replayed request. + if (Array.isArray(message.toolRequestIds)) { + const pending = new Set(message.toolRequestIds.filter((id: unknown): id is string => + typeof id === 'string' && id.length > 0, + )) + if (pending.size) pendingPermissions.set(chatId, pending) + else pendingPermissions.delete(chatId) + } + runtime.pendingPermissionCount = pendingPermissions.get(chatId)?.size ?? 0 + if (runtime.pendingPermissionCount) { + runtime.state = 'permission_pending' + } else if (!message.turnActive) { + runtime.state = 'idle' + } else if (runtime.state === 'idle' || runtime.state === 'permission_pending') { + runtime.state = 'thinking' + } + return true +} diff --git a/adapters/common/session-selection.ts b/adapters/common/session-selection.ts new file mode 100644 index 00000000..a51b5805 --- /dev/null +++ b/adapters/common/session-selection.ts @@ -0,0 +1,262 @@ +import * as fs from 'node:fs' +import * as path from 'node:path' +import type { AdapterHttpClient, RecentProject, SessionListItem } from './http-client.js' +import type { SessionStore } from './session-store.js' +import type { ServerMessage, WsBridge } from './ws-bridge.js' + +export type SessionRestoreDeps = { + httpClient: Pick + bridge: Pick + sessionStore: Pick + onServerMessage: (chatId: string, message: ServerMessage) => void | Promise + clearTransientState: (chatId: string) => void + isBusy?: (chatId: string) => boolean +} + +export type SessionSelectionDeps = SessionRestoreDeps & { + httpClient: Pick + sendNotice: (chatId: string, text: string) => Promise + clearProjectSelection: (chatId: string) => void + now?: () => number +} + +/** Attaching an IM is a binding change, never a new transcript or a prompt. */ +export async function restoreSelectedSession( + deps: SessionRestoreDeps, + chatId: string, + session: Pick, +): Promise<{ ok: boolean; message: string }> { + const busy = () => deps.isBusy?.(chatId) ?? false + if (busy()) return { ok: false, message: '当前会话正在运行或等待审批,请先处理审批或发送 /stop,等停止后再切换。' } + try { + if (!session.workDir || !await deps.httpClient.sessionExists(session.id)) { + return { ok: false, message: '该会话已不存在、目录不可用或不允许通过 IM 访问。请发送 /sessions 刷新列表。' } + } + } catch (err) { + return { ok: false, message: `无法检查会话,当前绑定未改变:${err instanceof Error ? err.message : String(err)}。请重试。` } + } + // A server event may have started a turn while the preflight was in flight. + if (busy()) return { ok: false, message: '当前会话正在运行,请先发送 /stop,等停止后再切换。' } + + const previous = deps.sessionStore.get(chatId) + const buffered: ServerMessage[] = [] + deps.bridge.resetSession(chatId) + try { + deps.bridge.connectSession(chatId, session.id) + deps.bridge.onServerMessage(chatId, (message) => { buffered.push(message) }) + if (!await deps.bridge.waitForOpen(chatId)) throw new Error('连接服务器超时') + deps.sessionStore.set(chatId, session.id, session.workDir) + } catch (err) { + deps.bridge.resetSession(chatId) + if (previous) { + try { + deps.bridge.connectSession(chatId, previous.sessionId) + deps.bridge.onServerMessage(chatId, (message) => deps.onServerMessage(chatId, message)) + } catch { + // The persisted binding is still intact; normal message recovery retries. + } + } + return { ok: false, message: `恢复失败,已保留原会话绑定:${err instanceof Error ? err.message : String(err)}。请重试。` } + } + deps.clearTransientState(chatId) + deps.bridge.onServerMessage(chatId, (message) => deps.onServerMessage(chatId, message)) + for (const message of buffered) { + try { + await deps.onServerMessage(chatId, message) + } catch (err) { + console.warn('[SessionSelection] Failed to present restored state:', err) + } + } + return { ok: true, message: `已恢复会话:${oneLine(session.title || '未命名会话', 80)}\n${session.workDir}\n会话 ID:${session.id}\n可以继续发送消息。` } +} + +const PAGE_SIZE = 8 +const FETCH_SIZE = 100 +export const SESSION_SELECTION_TTL_MS = 15 * 60 * 1000 +type Picker = { + expiresAt: number + page: number +} & ( + | { kind: 'projects'; projects: RecentProject[] } + | { kind: 'sessions'; project: string; sessions: SessionListItem[] } +) + +function oneLine(value: string, length = 100): string { + return value.replace(/[\r\n\t]+/g, ' ').slice(0, length) +} + +function canonicalPath(value: string): string { + try { + return fs.realpathSync(value) + } catch { + return path.resolve(value) + } +} + +/** The API's project filter addresses one transcript directory, while recent + * projects group worktrees. Filter logical roots after reading all summary pages. */ +export async function listProjectSessionHistory( + httpClient: Pick, + project: string, + all?: SessionListItem[], +): Promise { + const root = canonicalPath(project) + return (all ?? await loadSessionHistory(httpClient)).filter((session) => + session.workDir && session.workDirExists !== false && ( + canonicalPath(session.projectRoot || session.workDir) === root || canonicalPath(session.workDir) === root + ), + ).sort((a, b) => b.modifiedAt.localeCompare(a.modifiedAt) || a.id.localeCompare(b.id)) +} + +async function loadSessionHistory(httpClient: Pick): Promise { + const sessions = new Map() + let offset = 0 + let total: number + do { + const result = await httpClient.listSessions({ limit: FETCH_SIZE, offset }) + for (const session of result.sessions) sessions.set(session.id, session) + total = result.total + offset += FETCH_SIZE + } while (offset < total) + return [...sessions.values()] +} + +/** Text-first picker shared by every IM, with a stable, chat-local snapshot. */ +export class SessionSelectionController { + private readonly pickers = new Map() + private readonly now: () => number + + constructor(private readonly deps: SessionSelectionDeps) { + this.now = deps.now ?? Date.now + } + + clear(chatId: string): void { + this.pickers.delete(chatId) + } + + async handleInput(chatId: string, input: string): Promise { + const text = input.trim() + const list = /^(?:\/sessions|会话列表)(?:\s+(.*))?$/.exec(text) + const resume = /^(?:\/resume|继续会话)(?:\s+(.*))?$/.exec(text) + const cancel = text === '/cancel' || text === '取消选择' + const picker = this.pickers.get(chatId) + const numberReply = /^\d+$/.test(text) && Boolean(picker) + if (!list && !resume && !cancel && !numberReply) { + // Normal chat is never a fuzzy session query. Other commands (especially + // permission approval) retain their original routing and pending picker. + if (text && !text.startsWith('/') && !['帮助', '状态', '停止', '清空'].includes(text)) this.clear(chatId) + return false + } + + try { + if (cancel) { + this.clear(chatId) + this.deps.clearProjectSelection(chatId) + await this.deps.sendNotice(chatId, '已取消选择,当前会话未改变。') + } else if (list && ['next', 'prev'].includes(list[1] ?? '')) { + const active = await this.requirePicker(chatId, picker) + if (!active) return true + const count = active.kind === 'projects' ? active.projects.length : active.sessions.length + const nextPage = active.page + (list[1] === 'next' ? 1 : -1) + active.page = Math.max(0, Math.min(Math.ceil(count / PAGE_SIZE) - 1, nextPage)) + await this.show(chatId, active) + } else if (list || (resume && !resume[1])) { + this.clear(chatId) + this.deps.clearProjectSelection(chatId) + await this.open(chatId, list?.[1]) + } else { + const active = await this.requirePicker(chatId, picker) + if (!active) return true + await this.select(chatId, active, resume?.[1] ?? text) + } + } catch (err) { + await this.deps.sendNotice(chatId, `无法恢复会话:${err instanceof Error ? err.message : String(err)}。请重试 /sessions。`) + } + return true + } + + private async requirePicker(chatId: string, picker: Picker | undefined): Promise { + if (picker && picker.expiresAt > this.now()) return picker + this.clear(chatId) + await this.deps.sendNotice(chatId, '选择列表不存在或已过期,请发送 /sessions 重新选择。') + return null + } + + private async open(chatId: string, query?: string): Promise { + if (query === 'projects') return this.openProjects(chatId, await this.deps.httpClient.listRecentProjects()) + if (query) { + const { project, ambiguous } = await this.deps.httpClient.matchProject(query) + if (project) return this.openProject(chatId, project.realPath) + if (ambiguous?.length) return this.openProjects(chatId, ambiguous) + await this.deps.sendNotice(chatId, `未找到项目“${oneLine(query)}”。发送 /sessions projects 选择项目,或 /sessions <绝对路径>。`) + return + } + const current = this.deps.sessionStore.get(chatId) + if (current) { + // Recent projects collapse worktrees into their root. Find the stored + // session's logical root before choosing the default history list. + const sessions = await loadSessionHistory(this.deps.httpClient) + const active = sessions.find((session) => session.id === current.sessionId) + return this.openProject(chatId, active?.projectRoot || current.workDir, sessions) + } + await this.openProjects(chatId, await this.deps.httpClient.listRecentProjects()) + } + + private async openProjects(chatId: string, projects: RecentProject[]): Promise { + if (!projects.length) { + await this.deps.sendNotice(chatId, '没有可访问的历史项目。发送 /new 新建会话,或 /sessions <项目绝对路径> 查找旧会话。') + return + } + const picker: Picker = { kind: 'projects', projects, page: 0, expiresAt: this.now() + SESSION_SELECTION_TTL_MS } + this.pickers.set(chatId, picker) + await this.show(chatId, picker) + } + + private async openProject(chatId: string, project: string, all?: SessionListItem[]): Promise { + const sessions = await listProjectSessionHistory(this.deps.httpClient, project, all) + if (!sessions.length) { + this.clear(chatId) + await this.deps.sendNotice(chatId, `该项目没有可恢复会话:${project}\n发送 /sessions projects 选择其他项目,或 /new <项目> 新建会话。`) + return + } + const picker: Picker = { kind: 'sessions', project, sessions, page: 0, expiresAt: this.now() + SESSION_SELECTION_TTL_MS } + this.pickers.set(chatId, picker) + await this.show(chatId, picker) + } + + private async select(chatId: string, picker: Picker, query: string): Promise { + const index = /^\d+$/.test(query) ? Number(query) - 1 : -1 + const start = picker.page * PAGE_SIZE + const visible = index >= start && index < start + PAGE_SIZE + if (picker.kind === 'projects') { + const project = visible ? picker.projects[index] : undefined + if (project) return this.openProject(chatId, project.realPath) + } else { + const session = visible ? picker.sessions[index] : undefined + if (session) { + const result = await restoreSelectedSession(this.deps, chatId, session) + if (result.ok) this.clear(chatId) + await this.deps.sendNotice(chatId, result.message) + return + } + } + await this.deps.sendNotice(chatId, '编号无效,请使用当前页显示的编号,或发送 /sessions 刷新列表。') + } + + private async show(chatId: string, picker: Picker): Promise { + const start = picker.page * PAGE_SIZE + const currentId = this.deps.sessionStore.get(chatId)?.sessionId + const items = picker.kind === 'projects' ? picker.projects : picker.sessions + const lines = picker.kind === 'projects' + ? picker.projects.slice(start, start + PAGE_SIZE).map((project, i) => `${start + i + 1}. ${oneLine(project.projectName, 60)}\n${oneLine(project.realPath, 180)}`) + : picker.sessions.slice(start, start + PAGE_SIZE).map((session, i) => `${start + i + 1}. ${oneLine(session.title || '未命名会话', 60)}${session.id === currentId ? '(当前)' : ''}\n${session.modifiedAt.slice(0, 16).replace('T', ' ')} · ${session.messageCount} 条消息 · ${session.id.slice(0, 8)}`) + await this.deps.sendNotice(chatId, [ + picker.kind === 'projects' ? '选择历史会话所在的项目:' : `历史会话:${picker.project}`, + `第 ${picker.page + 1}/${Math.ceil(items.length / PAGE_SIZE)} 页`, + '', ...lines, '', + picker.kind === 'projects' ? '回复编号查看会话。' : '回复编号或 /resume <编号> 继续旧会话。', + '翻页:/sessions next、/sessions prev', + '/cancel 取消;列表 15 分钟内有效。', + ].join('\n')) + } +} diff --git a/adapters/common/ws-bridge.ts b/adapters/common/ws-bridge.ts index 288b710f..f0cab3d5 100644 --- a/adapters/common/ws-bridge.ts +++ b/adapters/common/ws-bridge.ts @@ -202,8 +202,7 @@ export class WsBridge { } if (msg.type === 'pong') return if (this.sessions.get(chatId) !== session) return - const handler = this.handlers.get(chatId) - if (!handler) return + if (!this.handlers.has(chatId)) return // Serialize per-chat handler calls: chain each message onto the previous // one so a slow handler (e.g. one awaiting im.message.create) fully @@ -213,7 +212,14 @@ export class WsBridge { const prev = this.handlerChains.get(chatId) ?? Promise.resolve() const next = prev .catch(() => {}) // upstream errors must not poison the chain - .then(() => Promise.resolve().then(() => handler(msg))) + .then(() => { + // Resetting a chat cannot cancel promises already queued for its old + // socket. Recheck ownership when delivery actually starts, then use + // the current handler so restoration's temporary buffer cannot trap + // initial status/permission frames after the live handler replaces it. + if (this.sessions.get(chatId) !== session) return + return this.handlers.get(chatId)?.(msg) + }) .catch((err) => { console.error(`[WsBridge] Handler error on ${chatId}:`, err) }) diff --git a/adapters/dingtalk/index.ts b/adapters/dingtalk/index.ts index f90e9098..181286a3 100644 --- a/adapters/dingtalk/index.ts +++ b/adapters/dingtalk/index.ts @@ -20,6 +20,8 @@ import { type PermissionDecision, } from '../common/permission.js' import { SessionStore } from '../common/session-store.js' +import { syncImPermissionState } from '../common/permission-sync.js' +import { SessionSelectionController } from '../common/session-selection.js' import { type RecentProject } from '../common/http-client.js' import { createAdapterClient } from '../common/adapter-client.js' import { @@ -88,6 +90,21 @@ attachmentStore.gc().catch((err) => { console.warn('[DingTalk] AttachmentStore.gc failed:', err instanceof Error ? err.message : err) }) +const sessionSelectionController = new SessionSelectionController({ + httpClient, + bridge, + sessionStore, + sendNotice: async (chatId, text) => { await sendText(chatId, text) }, + onServerMessage: handleServerMessage, + clearTransientState: clearTransientChatState, + clearProjectSelection: (chatId) => { projectSelectionController.clear(chatId) }, + isBusy: (chatId) => { + const runtime = getRuntimeState(chatId) + return runtime.state !== 'idle' || runtime.pendingPermissionCount > 0 + || (pendingPermissions.get(chatId)?.size ?? 0) > 0 + }, +}) + type ChatRuntimeState = { state: 'idle' | 'thinking' | 'streaming' | 'tool_executing' | 'permission_pending' verb?: string @@ -302,6 +319,7 @@ async function ensureSession(chatId: string): Promise { } async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { bridge.resetSession(chatId) clearTransientChatState(chatId) @@ -331,6 +349,7 @@ function formatProjectList(projects: RecentProject[]): string { } async function showProjectPicker(chatId: string): Promise { + sessionSelectionController.clear(chatId) try { const projects = await projectSelectionController.listProjects(chatId) if (projects.length === 0) { @@ -344,6 +363,7 @@ async function showProjectPicker(chatId: string): Promise { } function prepareNewSession(chatId: string): void { + sessionSelectionController.clear(chatId) bridge.resetSession(chatId) sessionStore.delete(chatId) clearTransientChatState(chatId) @@ -352,10 +372,12 @@ function prepareNewSession(chatId: string): void { async function handleServerMessage(chatId: string, msg: ServerMessage): Promise { const runtime = getRuntimeState(chatId) + if (syncImPermissionState(chatId, msg, runtime, pendingPermissions)) return switch (msg.type) { case 'connected': break + case 'status': runtime.state = msg.state runtime.verb = typeof msg.verb === 'string' ? msg.verb : undefined @@ -469,6 +491,8 @@ async function routeUserMessage(chatId: string, text: string, attachments: Attac if (!hasAttachments && handlePermissionCommand(chatId, trimmed)) return + if (!hasAttachments && await sessionSelectionController.handleInput(chatId, trimmed)) return + const projectOutcome = !hasAttachments ? await projectSelectionController.handleInput(chatId, trimmed) : null @@ -496,6 +520,7 @@ async function routeUserMessage(chatId: string, text: string, attachments: Attac await sendText(chatId, '⚠️ 无法发送 /clear,请先发送 /new 重新连接会话。') return } + getRuntimeState(chatId).state = 'thinking' await sendText(chatId, '🧹 已清空当前会话上下文。') return } @@ -518,7 +543,10 @@ async function routeUserMessage(chatId: string, text: string, attachments: Attac if (!ready) return const effectiveText = trimmed || (attachments.length > 0 ? '(用户发送了附件)' : '') if (!effectiveText && attachments.length === 0) return - if (!bridge.sendUserMessage(chatId, effectiveText, attachments.length ? attachments : undefined)) { + const sent = bridge.sendUserMessage(chatId, effectiveText, attachments.length ? attachments : undefined) + if (sent) { + getRuntimeState(chatId).state = 'thinking' + } else { await sendText(chatId, '⚠️ 消息发送失败,连接可能已断开。请发送 /new 重新开始。') } }) @@ -682,7 +710,9 @@ async function start(): Promise { process.once('SIGTERM', () => void shutdown()) } -start().catch((err) => { +if (import.meta.main || process.argv.includes('--dingtalk')) start().catch((err) => { console.error('[DingTalk] Fatal:', err instanceof Error ? err.message : err) process.exit(1) }) + +export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, routeUserMessage, handleRobotMessage, prepareNewSession } diff --git a/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts b/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts new file mode 100644 index 00000000..da7e2bce --- /dev/null +++ b/adapters/feishu/__tests__/legacy-session-selection-entrypoint.test.ts @@ -0,0 +1,322 @@ +import { afterAll, beforeAll, describe, expect, it, spyOn } from 'bun:test' +import * as fs from 'node:fs' +import * as os from 'node:os' +import * as path from 'node:path' +import { enqueue } from '../../common/chat-queue.js' +import { StreamingCard } from '../streaming-card.js' +import { FeishuMediaService } from '../media.js' +import { WechatMediaService } from '../../wechat/media.js' + +type Platform = 'feishu' | 'dingtalk' | 'wechat' | 'whatsapp' +const platforms: Platform[] = ['feishu', 'dingtalk', 'wechat', 'whatsapp'] +const chatFor = (platform: Platform) => platform === 'dingtalk' ? 'dingtalk:dm:fixture-user' : `${platform}-chat` +let feishu: typeof import('../index.js') +let dingtalk: typeof import('../../dingtalk/index.js') +let wechat: typeof import('../../wechat/index.js') +let whatsapp: typeof import('../../whatsapp/index.js') +let server: Bun.Server<{ sessionId: string }> +let temporaryRoot: string +let httpOrigin: string +let sequence = 0 +let newSessionCount = 0 +let preflightGate: Promise | undefined +let preflightStarted: (() => void) | undefined +const notices: string[] = [] +const sentPrompts: Array<{ sessionId: string; content: string }> = [] +const spies: Array<{ mockRestore(): void }> = [] +const originalFetch = globalThis.fetch +const environmentKeys = [ + 'HOME', 'CLAUDE_CONFIG_DIR', 'ADAPTER_SERVER_URL', 'ADAPTER_ALLOWED_PROJECT_ROOTS', + 'FEISHU_APP_ID', 'FEISHU_APP_SECRET', 'DINGTALK_CLIENT_ID', 'DINGTALK_CLIENT_SECRET', + 'WECHAT_ACCOUNT_ID', 'WECHAT_BOT_TOKEN', 'WECHAT_BASE_URL', 'WHATSAPP_AUTH_DIR', + 'CC_HAHA_LOCAL_ACCESS_TOKEN', +] +const savedEnvironment = new Map(environmentKeys.map((key) => [key, process.env[key]])) + +function entries() { + return [{ + id: 'history-active', title: 'Existing history', workDir: temporaryRoot, + projectRoot: temporaryRoot, projectPath: 'fixture', workDirExists: true, + modifiedAt: '2026-09-08T00:00:00Z', createdAt: '2026-09-01T00:00:00Z', messageCount: 4, + }] +} + +beforeAll(async () => { + temporaryRoot = fs.realpathSync(fs.mkdtempSync(path.join(os.tmpdir(), 'legacy-im-import-'))) + server = Bun.serve<{ sessionId: string }>({ + port: 0, + hostname: '127.0.0.1', + async fetch(request, currentServer) { + const url = new URL(request.url) + if (url.pathname.startsWith('/ws/')) { + if (currentServer.upgrade(request, { data: { sessionId: url.pathname.slice(4) } })) return + return new Response('upgrade required', { status: 400 }) + } + if (url.pathname === '/api/sessions/recent-projects') { + return Response.json({ projects: [{ + projectName: 'Fixture', realPath: temporaryRoot, projectPath: 'fixture', isGit: false, + repoName: null, branch: null, modifiedAt: '2026-09-08', sessionCount: 1, + }] }) + } + if (url.pathname === '/api/sessions' && request.method === 'POST') { + newSessionCount++ + return Response.json({ sessionId: `new-${newSessionCount}` }) + } + if (url.pathname === '/api/sessions') return Response.json({ sessions: entries(), total: 1 }) + if (url.pathname.endsWith('/git-info')) return Response.json({ repoName: 'Fixture', workDir: temporaryRoot, branch: 'main' }) + if (url.pathname.startsWith('/api/tasks/lists/')) return Response.json({ tasks: [{ id: 'task', subject: 'Fixture task', status: 'pending' }] }) + if (url.pathname.startsWith('/api/sessions/')) { + if (url.pathname.endsWith('/history-active') && preflightGate) { + preflightStarted?.() + await preflightGate + } + return Response.json({ workDir: temporaryRoot, permissionMode: 'default' }) + } + if (url.pathname === '/dingtalk' || url.pathname === '/ilink/bot/sendmessage') { + notices.push(await request.text()) + return Response.json({ ret: 0 }) + } + if (url.pathname === '/ilink/bot/getconfig') return Response.json({ ret: 0 }) + if (url.pathname === '/ilink/bot/sendtyping') return Response.json({ ret: 0 }) + return new Response(`Unexpected fixture request ${url.pathname}`, { status: 500 }) + }, + websocket: { + open(socket) { + socket.send(JSON.stringify({ type: 'connected' })) + socket.send(JSON.stringify({ + type: 'permission_requests_snapshot', toolRequestIds: [], + turnActive: socket.data.sessionId === 'history-active', + })) + }, + message(socket, raw) { + const message = JSON.parse(String(raw)) + if (message.type === 'user_message') sentPrompts.push({ sessionId: socket.data.sessionId, content: message.content }) + }, + }, + }) + httpOrigin = `http://127.0.0.1:${server.port}` + const env = { + HOME: temporaryRoot, + CLAUDE_CONFIG_DIR: path.join(temporaryRoot, '.claude'), + ADAPTER_SERVER_URL: httpOrigin.replace('http:', 'ws:'), + ADAPTER_ALLOWED_PROJECT_ROOTS: temporaryRoot, + FEISHU_APP_ID: 'fixture-app', FEISHU_APP_SECRET: 'fixture-secret', + DINGTALK_CLIENT_ID: 'fixture-client', DINGTALK_CLIENT_SECRET: 'fixture-secret', + WECHAT_ACCOUNT_ID: 'fixture-account', WECHAT_BOT_TOKEN: 'fixture-token', WECHAT_BASE_URL: httpOrigin, + WHATSAPP_AUTH_DIR: path.join(temporaryRoot, 'whatsapp-auth'), CC_HAHA_LOCAL_ACCESS_TOKEN: '', + } + Object.assign(process.env, env) + fs.mkdirSync(env.CLAUDE_CONFIG_DIR, { recursive: true }) + fs.mkdirSync(env.WHATSAPP_AUTH_DIR, { recursive: true }) + fs.writeFileSync(path.join(env.WHATSAPP_AUTH_DIR, 'creds.json'), JSON.stringify({ registered: true, me: { id: 'fixture' } })) + const platformConfig = { defaultWorkDir: temporaryRoot, allowedProjectRoots: [temporaryRoot], allowedUsers: ['fixture-user', ...platforms.map(chatFor)] } + fs.writeFileSync(path.join(env.CLAUDE_CONFIG_DIR, 'adapters.json'), JSON.stringify({ + serverUrl: env.ADAPTER_SERVER_URL, allowedProjectRoots: [temporaryRoot], + feishu: platformConfig, dingtalk: platformConfig, wechat: platformConfig, whatsapp: platformConfig, + })) + globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => { + const url = new URL(input instanceof Request ? input.url : String(input)) + if (url.origin === httpOrigin) return originalFetch(input, init) + if (url.href === 'https://api.dingtalk.com/v1.0/oauth2/accessToken') { + return Response.json({ accessToken: 'fixture-access-token', expireIn: 7200 }) + } + throw new Error(`External network is forbidden in this fixture: ${url.origin}${url.pathname}`) + }) as typeof fetch + feishu = await import('../index.js') + dingtalk = await import('../../dingtalk/index.js') + wechat = await import('../../wechat/index.js') + whatsapp = await import('../../whatsapp/index.js') + spies.push(spyOn(feishu.larkClient.im.message, 'create').mockImplementation(async (request: any) => { + notices.push(String(request.data.content)) + return { code: 0, data: { message_id: `fixture-${++sequence}` } } + })) + const attachment = { kind: 'image' as const, name: 'fixture.png', path: path.join(temporaryRoot, 'fixture.png'), buffer: Buffer.from('fixture'), size: 7, mimeType: 'image/png' } + spies.push(spyOn(FeishuMediaService.prototype, 'downloadResource').mockImplementation(async () => attachment)) + spies.push(spyOn(WechatMediaService.prototype, 'downloadCandidate').mockImplementation(async () => attachment)) + spies.push(spyOn(StreamingCard.prototype, 'ensureCreated').mockImplementation(async () => {})) + spies.push(spyOn(StreamingCard.prototype, 'abort').mockImplementation(async () => {})) + spies.push(spyOn(StreamingCard.prototype, 'finalize').mockImplementation(async () => {})) + whatsapp.useWhatsAppSocket({ sendMessage: async (_chatId: string, message: { text: string }) => { notices.push(message.text) } } as any) +}) + +afterAll(async () => { + for (const adapter of [feishu, dingtalk, wechat, whatsapp]) { + adapter?.bridge.destroy() + adapter?.dedup.destroy() + } + wechat?.typingController.destroy() + for (const spy of spies.reverse()) spy.mockRestore() + globalThis.fetch = originalFetch + server?.stop(true) + for (const [key, value] of savedEnvironment) { + if (value === undefined) delete process.env[key] + else process.env[key] = value + } + if (temporaryRoot) fs.rmSync(temporaryRoot, { recursive: true, force: true }) +}) + +function adapterFor(platform: Platform) { + return { feishu, dingtalk, wechat, whatsapp }[platform] +} + +async function send(platform: Platform, text: string, options: { unauthorized?: boolean; attachment?: boolean } = {}) { + const userId = options.unauthorized ? 'outsider' : 'fixture-user' + const chatId = options.unauthorized && platform === 'wechat' ? 'outsider' : chatFor(platform) + sequence++ + if (platform === 'feishu') { + await feishu.handleMessage({ + sender: { sender_id: { open_id: userId } }, + message: { message_id: `message-${sequence}`, chat_id: chatId, chat_type: 'p2p', message_type: options.attachment ? 'post' : 'text', content: JSON.stringify(options.attachment ? { zh_cn: { content: [[{ tag: 'text', text }, { tag: 'img', image_key: 'fixture-image' }]] } } : { text }) }, + }) + } else if (platform === 'dingtalk') { + if (options.attachment) { + await dingtalk.routeUserMessage(chatId, text, [{ type: 'image', data: 'Zml4dHVyZQ==', mimeType: 'image/png' }]) + } else await dingtalk.handleRobotMessage({ + conversationType: '1', senderStaffId: userId, senderId: chatId, + conversationId: chatId, msgtype: 'text', + text: { content: text }, sessionWebhook: `${httpOrigin}/dingtalk`, + }) + } else if (platform === 'wechat') { + await wechat.routeUserMessage({ from_user_id: chatId, message_id: sequence, item_list: [{ type: 1, text_item: { text } }, ...(options.attachment ? [{ type: 2, image_item: { media: { full_url: 'https://fixture.invalid/image.png' } } }] : [])] }) + } else { + await whatsapp.routeUserMessage(chatId, userId, 'Fixture User', text, options.attachment ? [{ type: 'image', data: 'Zml4dHVyZQ==', mimeType: 'image/png' }] : []) + } + // Flush the actual adapter's per-chat queue without assuming its handler awaits it. + await enqueue(platform === 'dingtalk' && options.unauthorized ? 'dingtalk:dm:outsider' : chatId, async () => {}) +} + +for (const platform of platforms) { + describe(`${platform} actual module session selection`, () => { + it('restores original history and honors the real initial active-turn snapshot', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + adapter.sessionStore.set(chatId, 'current', temporaryRoot) + const creationsBefore = newSessionCount + const sendSpy = spyOn(adapter.bridge, 'sendUserMessage') + spies.push(sendSpy) + await send(platform, '/sessions') + expect(adapter.sessionStore.get(chatId)?.sessionId).toBe('current') + await send(platform, '/resume 1') + await enqueue(chatId, async () => {}) + expect(adapter.sessionStore.get(chatId)?.sessionId).toBe('history-active') + expect(adapter.getRuntimeState(chatId).state).toBe('thinking') + expect(newSessionCount).toBe(creationsBefore) + expect(sendSpy).not.toHaveBeenCalled() + await send(platform, '/sessions') + await send(platform, '/resume 1') + expect(notices.at(-1)).toContain('正在运行') + await adapter.handleServerMessage(chatId, { type: 'status', state: 'idle' }) + sendSpy.mockRestore() + }) + + it('invalidates history selection through project picking, new sessions and project cards', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + await send(platform, '/sessions') + await send(platform, '/projects') + expect(await adapter.sessionSelectionController.handleInput(chatId, '1')).toBe(false) + await send(platform, '/new Fixture') + expect(adapter.sessionStore.get(chatId)?.sessionId).toStartWith('new-') + await send(platform, '/sessions') + if (platform === 'feishu') { + await feishu.handleCardAction({ context: { open_chat_id: chatId }, action: { value: { action: 'pick_project', realPath: temporaryRoot, projectName: 'Fixture' } } }) + } else { + expect(await adapter.createSessionForChat(chatId, temporaryRoot)).toBe(true) + } + expect(await adapter.sessionSelectionController.handleInput(chatId, '1')).toBe(false) + }) + + it('rejects unpaired users and leaves attachment text on the normal send path', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + const binding = adapter.sessionStore.get(chatId)?.sessionId + await send(platform, '/sessions', { unauthorized: true }) + expect(notices.at(-1)).toContain('未授权') + expect(adapter.sessionStore.get(chatId)?.sessionId).toBe(binding) + adapter.clearTransientChatState(chatId) + await send(platform, '/sessions') + const sendSpy = spyOn(adapter.bridge, 'sendUserMessage') + try { + await send(platform, '/sessions', { attachment: true }) + expect(sendSpy).toHaveBeenCalledTimes(1) + expect(sendSpy.mock.calls[0]?.[1]).toBe('/sessions') + expect(sendSpy.mock.calls[0]?.[2]?.[0]?.type).toBe('image') + } finally { sendSpy.mockRestore() } + adapter.clearTransientChatState(chatId) + }) + + it('keeps numeric approval ahead of history and accepts desktop permission resolution', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + await send(platform, '/sessions') + await adapter.handleServerMessage(chatId, { type: 'permission_requests_snapshot', toolRequestIds: ['approval'], turnActive: true }) + const binding = adapter.sessionStore.get(chatId)?.sessionId + await send(platform, '1') + expect(adapter.sessionStore.get(chatId)?.sessionId).toBe(binding) + expect(adapter.getRuntimeState(chatId).pendingPermissionCount).toBe(0) + await adapter.handleServerMessage(chatId, { type: 'permission_requests_snapshot', toolRequestIds: ['desktop-approval'], turnActive: true }) + await adapter.handleServerMessage(chatId, { type: 'permission_resolved', permissionType: 'tool', requestId: 'desktop-approval' }) + await adapter.handleServerMessage(chatId, { type: 'status', state: 'idle' }) + expect(adapter.getRuntimeState(chatId).pendingPermissionCount).toBe(0) + await send(platform, '/resume 1') + expect(adapter.sessionStore.get(chatId)?.sessionId).toBe('history-active') + adapter.clearTransientChatState(chatId) + }) + + it('keeps help, status and stop commands usable while history selection is open', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + await send(platform, '/sessions') + await send(platform, '/help') + expect(notices.at(-1)).toContain('/sessions') + await send(platform, '/status') + expect(notices.at(-1)).toContain('Fixture') + const stopSpy = spyOn(adapter.bridge, 'sendStopGeneration') + try { + await send(platform, '/stop') + expect(stopSpy).toHaveBeenCalledWith(chatId) + } finally { stopSpy.mockRestore() } + await send(platform, '/cancel') + expect(await adapter.sessionSelectionController.handleInput(chatId, '1')).toBe(false) + }) + + it('sets busy synchronously on successful normal and clear sends, while failed sends stay idle', async () => { + const adapter = adapterFor(platform) + const chatId = chatFor(platform) + for (const command of ['continue the previous task', '/clear']) { + adapter.clearTransientChatState(chatId) + await send(platform, command) + expect(adapter.getRuntimeState(chatId).state).toBe('thinking') + adapter.clearTransientChatState(chatId) + const sendSpy = spyOn(adapter.bridge, 'sendUserMessage').mockReturnValue(false) + try { + await send(platform, command) + expect(adapter.getRuntimeState(chatId).state).toBe('idle') + } finally { sendSpy.mockRestore() } + } + }) + }) +} + + +it('serializes Feishu project cards after an in-flight history selection', async () => { + const chatId = chatFor('feishu') + feishu.clearTransientChatState(chatId) + await send('feishu', '/sessions') + let release!: () => void + const started = new Promise((resolve) => { preflightStarted = resolve }) + preflightGate = new Promise((resolve) => { release = resolve }) + try { + const resume = send('feishu', '/resume 1') + await started + const card = feishu.handleCardAction({ context: { open_chat_id: chatId }, action: { value: { action: 'pick_project', realPath: temporaryRoot, projectName: 'Fixture' } } }) + release() + await Promise.all([resume, card]) + expect(feishu.sessionStore.get(chatId)?.sessionId).toStartWith('new-') + expect(await feishu.sessionSelectionController.handleInput(chatId, '1')).toBe(false) + } finally { + release() + preflightGate = undefined + preflightStarted = undefined + } +}) diff --git a/adapters/feishu/index.ts b/adapters/feishu/index.ts index c3378f4d..ba088149 100644 --- a/adapters/feishu/index.ts +++ b/adapters/feishu/index.ts @@ -26,6 +26,8 @@ import { type PermissionDecision, } from '../common/permission.js' import { SessionStore } from '../common/session-store.js' +import { syncImPermissionState } from '../common/permission-sync.js' +import { SessionSelectionController } from '../common/session-selection.js' import { type RecentProject } from '../common/http-client.js' import { createAdapterClient } from '../common/adapter-client.js' import { @@ -97,6 +99,21 @@ let botOpenId: string | null = null // WSClient reference for graceful shutdown let wsClient: InstanceType | null = null +const sessionSelectionController = new SessionSelectionController({ + httpClient, + bridge, + sessionStore, + sendNotice: async (chatId, text) => { await sendText(chatId, text) }, + onServerMessage: handleServerMessage, + clearTransientState: clearTransientChatState, + clearProjectSelection: (chatId) => { projectSelectionController.clear(chatId) }, + isBusy: (chatId) => { + const runtime = getRuntimeState(chatId) + return runtime.state !== 'idle' || runtime.pendingPermissionCount > 0 + || (pendingPermissions.get(chatId)?.size ?? 0) > 0 + }, +}) + type ChatRuntimeState = { state: 'idle' | 'thinking' | 'streaming' | 'tool_executing' | 'permission_pending' verb?: string @@ -643,6 +660,7 @@ async function ensureSession(chatId: string): Promise { } async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { // Always tear down any stale WS connection before creating a new session. // Without this, bridge.connectSession() below would short-circuit when an @@ -673,6 +691,7 @@ async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { const projects = await projectSelectionController.listProjects(chatId) if (projects.length === 0) { @@ -694,6 +713,7 @@ async function showProjectPicker(chatId: string): Promise { } function prepareNewSession(chatId: string): void { + sessionSelectionController.clear(chatId) bridge.resetSession(chatId) sessionStore.delete(chatId) // Abort any in-flight streaming card for the previous session @@ -712,11 +732,13 @@ function prepareNewSession(chatId: string): void { async function handleServerMessage(chatId: string, msg: ServerMessage): Promise { const runtime = getRuntimeState(chatId) + if (syncImPermissionState(chatId, msg, runtime, pendingPermissions)) return switch (msg.type) { case 'connected': break + case 'status': { runtime.state = msg.state runtime.verb = typeof msg.verb === 'string' ? msg.verb : undefined @@ -959,6 +981,8 @@ async function handleMessage(data: any): Promise { return } + if (!hasAttachments && await sessionSelectionController.handleInput(chatId, msgText)) return + const projectOutcome = !hasAttachments ? await projectSelectionController.handleInput(chatId, msgText) : null @@ -983,6 +1007,7 @@ async function handleMessage(data: any): Promise { } clearTransientChatState(chatId) const sent = bridge.sendUserMessage(chatId, '/clear') + if (sent) getRuntimeState(chatId).state = 'thinking' if (!sent) { await sendText(chatId, '⚠️ 无法发送 /clear,请先发送 /new 重新连接会话。') return @@ -1093,6 +1118,7 @@ async function handleMessage(data: any): Promise { }) const sent = bridge.sendUserMessage(chatId, effectiveText, attachments) + if (sent) getRuntimeState(chatId).state = 'thinking' if (!sent) { await sendText(chatId, '⚠️ 消息发送失败,连接可能已断开。请发送 /new 重新开始。') } @@ -1171,12 +1197,14 @@ async function handleCardAction(data: any): Promise { const projectName = event.action?.value?.projectName ?? realPath ?? '(unknown)' if (!realPath) return - projectSelectionController.clear(chatId) - // createSessionForChat handles its own error messaging on failure - const ok = await createSessionForChat(chatId, realPath) - if (ok) { - await sendText(chatId, `✅ 已新建会话:**${projectName}**`) - } + await enqueue(chatId, async () => { + projectSelectionController.clear(chatId) + // createSessionForChat handles its own error messaging on failure + const ok = await createSessionForChat(chatId, realPath) + if (ok) { + await sendText(chatId, `✅ 已新建会话:**${projectName}**`) + } + }) return { toast: { type: 'info', content: `📁 ${projectName}` } } } } @@ -1254,14 +1282,16 @@ async function start(): Promise { console.log('[Feishu] Bot is running! (WebSocket connected)') } -start().catch((err) => { +if (import.meta.main || process.argv.includes('--feishu')) start().catch((err) => { console.error('[Feishu] Failed to start:', err) process.exit(1) }) -process.on('SIGINT', () => { +if (import.meta.main || process.argv.includes('--feishu')) process.on('SIGINT', () => { console.log('[Feishu] Shutting down...') bridge.destroy() dedup.destroy() process.exit(0) }) + +export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, handleMessage, handleCardAction, larkClient, prepareNewSession } diff --git a/adapters/telegram/__tests__/commands.test.ts b/adapters/telegram/__tests__/commands.test.ts index 49e92e55..9621b35e 100644 --- a/adapters/telegram/__tests__/commands.test.ts +++ b/adapters/telegram/__tests__/commands.test.ts @@ -118,6 +118,7 @@ function createController(overrides?: Record) { title: 'Fix IM', createdAt: '2026-06-09T07:00:00.000Z', workDir: '/work/repo', + projectRoot: '/work/repo', projectPath: '/work/repo', workDirExists: true, modifiedAt: '2026-06-09T08:00:00.000Z', @@ -126,11 +127,15 @@ function createController(overrides?: Record) { ], total: 1, })), + sessionExists: mock(async () => true), }, defaultWorkDir: '/work/repo', isAllowedUser: mock(() => true), ensureExistingSession: mock(async () => ({ sessionId: 'active', workDir: '/work/repo' })), clearTransientChatState: mock((chatId: string) => bridgeEvents.push(`clear:${chatId}`)), + clearOtherSelections: mock(() => {}), + isBusy: mock(() => false), + getStoredSession: mock(() => ({ sessionId: 'active', workDir: '/work/repo', updatedAt: 1 })), setStoredSession: mock((chatId: string, sessionId: string, workDir: string) => { bridgeEvents.push(`store:${chatId}:${sessionId}:${workDir}`) }), @@ -140,6 +145,7 @@ function createController(overrides?: Record) { bridgeEvents.push(`connect:${chatId}:${sessionId}`) }), onBridgeServerMessage: mock((chatId: string) => bridgeEvents.push(`listen:${chatId}`)), + handleServerMessage: mock(async () => {}), waitForBridgeOpen: mock(async () => true), sendUserMessage: mock((chatId: string, content: string) => { sentUserMessages.push({ chatId, content }) @@ -239,7 +245,7 @@ describe('Telegram command controller helpers', () => { } as any registerTelegramExtendedCommands(bot, controller) - expect(commands).toEqual(['start', 'help', 'resume', 'provider', 'model', 'skills']) + expect(commands).toEqual(['start', 'help', 'provider', 'model', 'skills']) const handled = await tryHandleTelegramSelectionCallback( 'tgsel:model:pick:2', @@ -555,8 +561,7 @@ describe('Telegram command controller helpers', () => { }) expect(deps.httpClient.listSessions).toHaveBeenCalledWith({ - project: '/work/repo', - limit: 50, + limit: 100, offset: 0, }) expect(projectCallback.edits[0]).toContain('选择要恢复的会话') @@ -568,16 +573,115 @@ describe('Telegram command controller helpers', () => { index: 0, }) - expect(bridgeEvents).toEqual([ - 'reset:42', - 'clear:42', - 'store:42:session-123456789:/work/repo', - 'connect:42:session-123456789', - 'listen:42', - ]) + expect(bridgeEvents).toContain('connect:42:session-123456789') + expect(bridgeEvents.indexOf('store:42:session-123456789:/work/repo')) + .toBeGreaterThan(bridgeEvents.indexOf('connect:42:session-123456789')) + expect(deps.httpClient.sessionExists).toHaveBeenCalledWith('session-123456789') expect(sessionCallback.edits[0]).toContain('已恢复会话') }) + it('keeps the current binding when a selected historical session has been deleted', async () => { + const { controller, deps, bridgeEvents } = createController() + deps.httpClient.sessionExists = mock(async () => false) + await controller.handleResumeCommand(createCommandContext().ctx) + await controller.handleSelectionCallback(createCommandContext().ctx, { + kind: 'resume_project', action: 'pick', index: 0, + }) + const callback = createCommandContext() + await controller.handleSelectionCallback(callback.ctx, { + kind: 'resume_session', action: 'pick', index: 0, + }) + expect(bridgeEvents).toEqual([]) + expect(callback.edits[0]).toContain('不存在') + }) + + it('loads every history page and restores a worktree through the project resume menu', async () => { + const { controller, deps } = createController() + const histories = Array.from({ length: 105 }, (_, index) => ({ + id: `session-${String(index).padStart(3, '0')}`, + title: index === 104 ? 'Tree migration' : `History ${index}`, + createdAt: '2026-06-09T07:00:00.000Z', + modifiedAt: '2026-06-09T08:00:00.000Z', + messageCount: 5, + workDir: index === 104 ? '/work/repo-feature' : '/work/repo', + projectRoot: '/work/repo', + projectPath: index === 104 ? '-work-repo-feature' : '-work-repo', + workDirExists: true, + })) + histories.push({ ...histories[0], id: 'foreign', title: 'Other project', workDir: '/work/other', projectRoot: '/work/other' }) + deps.httpClient.listSessions.mockImplementation(async (query: { project?: string; limit: number; offset: number }) => { + const matches = query.project ? histories.filter((session) => session.workDir === query.project) : histories + return { sessions: matches.slice(query.offset, query.offset + query.limit), total: matches.length } + }) + + await controller.handleResumeCommand(createCommandContext().ctx) + await controller.handleSelectionCallback(createCommandContext().ctx, { + kind: 'resume_project', action: 'pick', index: 0, + }) + const lastPage = createCommandContext() + await controller.handleSelectionCallback(lastPage.ctx, { + kind: 'resume_session', action: 'page', index: 13, + }) + expect(lastPage.edits[0]).toContain('Tree migration') + expect(lastPage.edits[0]).not.toContain('Other project') + expect(deps.httpClient.listSessions.mock.calls).toEqual([ + [{ limit: 100, offset: 0 }], + [{ limit: 100, offset: 100 }], + ]) + + const callback = createCommandContext() + await controller.handleSelectionCallback(callback.ctx, { + kind: 'resume_session', action: 'pick', index: 104, + }) + expect(deps.setStoredSession).toHaveBeenCalledWith('42', 'session-104', '/work/repo-feature') + expect(deps.connectBridgeSession).toHaveBeenCalledWith('42', 'session-104') + expect(callback.edits[0]).toContain('已恢复会话') + }) + + it('refuses to switch an active turn or permission request from the resume menu', async () => { + const { controller, deps, bridgeEvents } = createController({ isBusy: () => true }) + await controller.handleResumeCommand(createCommandContext().ctx) + expect(deps.clearOtherSelections).toHaveBeenCalledWith('42') + await controller.handleSelectionCallback(createCommandContext().ctx, { + kind: 'resume_project', action: 'pick', index: 0, + }) + const callback = createCommandContext() + await controller.handleSelectionCallback(callback.ctx, { + kind: 'resume_session', action: 'pick', index: 0, + }) + expect(bridgeEvents).toEqual([]) + expect(callback.edits[0]).toContain('/stop') + expect(deps.httpClient.sessionExists).not.toHaveBeenCalled() + }) + + it('clears an older menu before a new resume list fails to load', async () => { + const { controller, deps } = createController() + await controller.handleResumeCommand(createCommandContext().ctx) + deps.httpClient.listRecentProjects.mockRejectedValueOnce(new Error('offline')) + await controller.handleResumeCommand(createCommandContext().ctx) + const callback = createCommandContext() + await controller.handleSelectionCallback(callback.ctx, { + kind: 'resume_project', action: 'pick', index: 0, + }) + expect(callback.answers[0]).toContain('选择已过期') + expect(deps.httpClient.listSessions).not.toHaveBeenCalled() + }) + + it('reports a preflight failure without disconnecting the current session', async () => { + const { controller, deps, bridgeEvents } = createController() + deps.httpClient.sessionExists.mockRejectedValueOnce(new Error('server unavailable')) + await controller.handleResumeCommand(createCommandContext().ctx) + await controller.handleSelectionCallback(createCommandContext().ctx, { + kind: 'resume_project', action: 'pick', index: 0, + }) + const callback = createCommandContext() + await controller.handleSelectionCallback(callback.ctx, { + kind: 'resume_session', action: 'pick', index: 0, + }) + expect(bridgeEvents).toEqual([]) + expect(callback.edits[0]).toContain('server unavailable') + }) + it('handles selection callback edge cases and resume timeout cleanup', async () => { const unauthorized = createController({ isAllowedUser: mock(() => false) }) const denied = createCommandContext() @@ -629,9 +733,10 @@ describe('Telegram command controller helpers', () => { index: 0, }) - expect(deps.deleteStoredSession).toHaveBeenCalledWith('42') - expect(bridgeEvents).toContain('delete:42') - expect(timeout.edits[0]).toContain('连接服务器超时') + expect(deps.deleteStoredSession).not.toHaveBeenCalled() + expect(deps.setStoredSession).not.toHaveBeenCalled() + expect(bridgeEvents).toContain('connect:42:active') + expect(timeout.edits[0]).toContain('超时') }) it('handles selection callback failures for provider, model, and session lists', async () => { @@ -698,7 +803,10 @@ describe('Telegram command controller helpers', () => { defaultWorkDir: '/work/repo', bridge: { resetSession: (chatId) => events.push(`reset:${chatId}`), - connectSession: (chatId, sessionId) => events.push(`connect:${chatId}:${sessionId}`), + connectSession: (chatId, sessionId) => { + events.push(`connect:${chatId}:${sessionId}`) + return true + }, onServerMessage: (chatId, handler) => { events.push(`listen:${chatId}`) void handler({ type: 'connected' }) @@ -714,16 +822,20 @@ describe('Telegram command controller helpers', () => { }, }, sessionStore: { + get: () => ({ sessionId: 'active', workDir: '/work/repo', updatedAt: 1 }), set: (chatId, sessionId, workDir) => events.push(`store:${chatId}:${sessionId}:${workDir}`), delete: (chatId) => events.push(`delete:${chatId}`), }, isAllowedUser: () => allowPermissionUser, ensureExistingSession: mock(async () => ({ sessionId: 'active', workDir: '/work/repo' })), clearTransientChatState: (chatId) => events.push(`clear:${chatId}`), + clearOtherSelections: () => {}, + isBusy: () => false, handleServerMessage: (chatId, msg) => { events.push(`message:${chatId}:${(msg as any).type}`) }, setRuntimeModel: (chatId, modelId) => events.push(`model:${chatId}:${modelId}`), + setRuntimeBusy: (chatId) => events.push(`busy:${chatId}`), }) await controller.setModelFromCommand('42', 'model-x') @@ -750,6 +862,13 @@ describe('Telegram command controller helpers', () => { expect(events).toContain('decrement:42') expect(permissionCtx.edits[0]).toContain('已允许') + await controller.handleSkillsCommand(createCommandContext().ctx) + await controller.handleSelectionCallback(createCommandContext().ctx, { + kind: 'skill', action: 'pick', index: 0, + }) + expect(events).toContain('send:42:/skill-a') + expect(events).toContain('busy:42') + const missingIdentityCtx = createCommandContext() delete (missingIdentityCtx.ctx as any).from await expect(controller.handlePermissionCallback(missingIdentityCtx.ctx, { diff --git a/adapters/telegram/__tests__/entrypoint-session-routing.test.ts b/adapters/telegram/__tests__/entrypoint-session-routing.test.ts new file mode 100644 index 00000000..7db606b6 --- /dev/null +++ b/adapters/telegram/__tests__/entrypoint-session-routing.test.ts @@ -0,0 +1,302 @@ +import { afterAll, beforeAll, describe, expect, it, spyOn } from 'bun:test' +import { mkdtempSync, mkdirSync, realpathSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import type { ServerWebSocket } from 'bun' +import { SessionStore } from '../../common/session-store.js' +import { AttachmentStore } from '../../common/attachment/attachment-store.js' + +// Import the actual entrypoint with isolated configuration. Telegram API calls +// terminate in grammY's documented transformer; HTTP and WS use loopback only. +describe('Telegram entrypoint session routing', () => { + const envKeys = ['CLAUDE_CONFIG_DIR', 'TELEGRAM_BOT_TOKEN', 'ADAPTER_SERVER_URL', 'ADAPTER_ALLOWED_PROJECT_ROOTS', 'ADAPTER_DEFAULT_PROJECT_DIR', 'CLAUDE_ADAPTER_DEFAULT_WORK_DIR', 'CC_HAHA_LOCAL_ACCESS_TOKEN'] + const previousEnv = new Map() + let directory: string + let project: string + let worktree: string + let entry: typeof import('../index.js') + let server: ReturnType> + let store: SessionStore + let nextId = 100 + const apiCalls: Array<{ method: string; payload: any }> = [] + const requests: string[] = [] + const messages: Array<{ sessionId: string; message: any }> = [] + const sockets = new Map>>() + const sessionPaths = new Map() + + async function eventually(assertion: () => void): Promise { + const deadline = Date.now() + 2500 + while (true) { + try { assertion(); return } catch (error) { + if (Date.now() >= deadline) throw error + await new Promise((resolve) => setTimeout(resolve, 10)) + } + } + } + + function texts(chatId: number): string[] { + return apiCalls.filter((call) => call.payload.chat_id === chatId && call.payload.text).map((call) => call.payload.text) + } + + async function text(chatId: number, value: string, options: { userId?: number; messageId?: number; photo?: boolean } = {}): Promise { + const messageId = options.messageId ?? nextId++ + const message = { + message_id: messageId, + date: 1, + chat: { id: chatId, type: 'private' }, + from: { id: options.userId ?? 7, is_bot: false, first_name: 'Fixture' }, + ...(options.photo ? { + photo: [{ file_id: 'fixture-photo', file_unique_id: 'unique', width: 1, height: 1 }], + caption: value, + } : { + text: value, + ...(value.startsWith('/') ? { entities: [{ type: 'bot_command', offset: 0, length: value.split(' ')[0].length }] } : {}), + }), + } + await entry.bot.handleUpdate({ update_id: messageId, message } as any) + } + + async function callback(chatId: number, data: string, id = `callback-${nextId++}`): Promise { + await entry.bot.handleUpdate({ + update_id: nextId++, + callback_query: { + id, data, chat_instance: 'fixture', + from: { id: 7, is_bot: false, first_name: 'Fixture' }, + message: { message_id: 1, date: 1, chat: { id: chatId, type: 'private' }, text: 'fixture menu' }, + }, + } as any) + } + + function broadcast(sessionId: string, message: unknown): void { + for (const socket of sockets.get(sessionId) ?? []) socket.send(JSON.stringify(message)) + } + + beforeAll(async () => { + for (const key of envKeys) previousEnv.set(key, process.env[key]) + directory = realpathSync(mkdtempSync(join(tmpdir(), 'telegram-entry-'))) + project = join(directory, 'repo') + worktree = join(directory, 'repo-feature') + mkdirSync(project) + mkdirSync(worktree) + for (const id of ['old', 'history', 'running', 'stream']) sessionPaths.set(id, id === 'history' ? worktree : project) + server = Bun.serve<{ sessionId: string }>({ + hostname: '127.0.0.1', port: 0, + async fetch(request, server) { + const url = new URL(request.url) + requests.push(`${request.method} ${url.pathname}`) + if (url.pathname.startsWith('/ws/')) { + if (server.upgrade(request, { data: { sessionId: url.pathname.split('/')[2] } })) return + return new Response('upgrade failed', { status: 400 }) + } + if (url.pathname === '/api/sessions/recent-projects') return Response.json({ projects: [{ projectName: 'repo', realPath: project, projectPath: '-fixture-repo', branch: 'main', sessionCount: 3 }] }) + if (url.pathname === '/api/sessions' && request.method === 'POST') { + const body = await request.json() as { workDir: string } + const sessionId = `created-${nextId++}` + sessionPaths.set(sessionId, body.workDir) + return Response.json({ sessionId }) + } + if (url.pathname === '/api/sessions') return Response.json({ + sessions: ['history', 'running', 'old'].map((id, index) => ({ id, title: `${id} title`, createdAt: '2026-01-01', modifiedAt: `2026-06-0${3 - index}`, workDir: sessionPaths.get(id), projectRoot: project, projectPath: '-fixture-repo', workDirExists: true, messageCount: 4 })), total: 3, + }) + const sessionId = url.pathname.split('/')[3] + if (sessionPaths.has(sessionId)) return Response.json({ workDir: sessionPaths.get(sessionId), repoName: 'repo', branch: 'main' }) + if (url.pathname === '/api/skills') return Response.json({ skills: [{ name: 'fixture', displayName: 'Fixture', description: 'Fixture skill', source: 'plugin', userInvocable: true }] }) + if (url.pathname === '/api/models/current') return Response.json({ model: { id: 'fixture-model' } }) + if (url.pathname === '/api/tasks') return Response.json({ tasks: [] }) + return new Response('Unexpected fixture endpoint', { status: 404 }) + }, + websocket: { + open(socket) { + const peers = sockets.get(socket.data.sessionId) ?? new Set() + peers.add(socket) + sockets.set(socket.data.sessionId, peers) + socket.send(JSON.stringify({ type: 'connected' })) + socket.send(JSON.stringify({ type: 'permission_requests_snapshot', turnActive: socket.data.sessionId === 'running', toolRequestIds: [], computerUseRequestIds: [] })) + }, + message(socket, raw) { + const message = JSON.parse(String(raw)) + messages.push({ sessionId: socket.data.sessionId, message }) + if (message.content === '/clear') socket.send(JSON.stringify({ type: 'message_complete' })) + }, + close(socket) { sockets.get(socket.data.sessionId)?.delete(socket) }, + }, + }) + process.env.CLAUDE_CONFIG_DIR = directory + process.env.TELEGRAM_BOT_TOKEN = '12345:fixture-token' + process.env.ADAPTER_SERVER_URL = `ws://127.0.0.1:${server.port}` + process.env.ADAPTER_ALLOWED_PROJECT_ROOTS = directory + process.env.ADAPTER_DEFAULT_PROJECT_DIR = project + process.env.CLAUDE_ADAPTER_DEFAULT_WORK_DIR = project + process.env.CC_HAHA_LOCAL_ACCESS_TOKEN = 'fixture-local-token' + writeFileSync(join(directory, 'adapters.json'), JSON.stringify({ telegram: { allowedUsers: [7], defaultWorkDir: project, allowedProjectRoots: [directory] } })) + store = new SessionStore(join(directory, 'adapter-sessions.json')) + entry = await import('../index.js') + entry.bot.botInfo = { id: 12345, is_bot: true, first_name: 'Fixture', username: 'fixture_bot', can_join_groups: false, can_read_all_group_messages: false, supports_inline_queries: false, can_manage_bots: false, can_connect_to_business: false, has_main_web_app: false, has_topics_enabled: false, allows_users_to_create_topics: false } + entry.bot.api.config.use(async (_previous, method, payload) => { + apiCalls.push({ method, payload }) + if (method === 'getFile') throw new Error('Fixture download rejected before network') + return { ok: true, result: ['answerCallbackQuery', 'deleteMessage'].includes(method) ? true : { message_id: nextId++, date: 1, chat: { id: (payload as any).chat_id, type: 'private' }, text: (payload as any).text } } as any + }) + }) + + afterAll(async () => { + entry?.stopTelegramAdapter() + await server?.stop(true) + for (const [key, value] of previousEnv) { + if (value === undefined) delete process.env[key] + else process.env[key] = value + } + if (directory) rmSync(directory, { recursive: true, force: true }) + }) + + it('runs registered history commands after authorization and deduplication', async () => { + const before = requests.length + await text(701, '/sessions', { userId: 99 }) + expect(texts(701).at(-1)).toContain('未授权') + expect(requests.length).toBe(before) + store.set('702', 'old', project) + await text(702, '/sessions', { messageId: 10001 }) + const afterFirst = requests.length + await text(702, '/sessions', { messageId: 10001 }) + expect(requests.length).toBe(afterFirst) + await text(702, '/resume 1') + expect(store.get('702')?.sessionId).toBe('history') + expect(store.get('702')?.workDir).toBe(worktree) + await text(702, 'Continue this history') + await eventually(() => expect(messages.some((item) => item.sessionId === 'history' && item.message.content === 'Continue this history')).toBe(true)) + await text(702, '/sessions') + await text(702, '/resume 3') + expect(store.get('702')?.sessionId).toBe('history') + expect(texts(702).at(-1)).toContain('/stop') + broadcast('history', { type: 'message_complete' }) + }) + + it('reads active-turn snapshots when resuming history before any status event', async () => { + store.set('703', 'old', project) + await text(703, '/sessions') + await text(703, '/resume 2') + expect(store.get('703')?.sessionId).toBe('running') + await text(703, '/sessions') + await text(703, '1') + expect(store.get('703')?.sessionId).toBe('running') + expect(texts(703).at(-1)).toContain('/stop') + }) + + it('serializes menu and text input and clears conflicting selection states', async () => { + store.set('704', 'old', project) + await Promise.all([text(704, '/sessions'), text(704, '/resume')]) + expect(texts(704).at(-1)).toContain('选择要恢复的项目') + await text(704, '/resume 1') + expect(texts(704).at(-1)).toContain('已过期') + await callback(704, 'tgsel:resume_project:pick:0', 'dedup-callback') + const afterFirst = requests.length + await callback(704, 'tgsel:resume_project:pick:0', 'dedup-callback') + expect(requests.length).toBe(afterFirst) + await callback(704, 'tgsel:resume_session:pick:0') + expect(store.get('704')?.sessionId).toBe('history') + await text(704, '/sessions') + await text(704, '/projects') + await text(704, 'repo') + expect(store.get('704')?.sessionId).toStartWith('created-') + expect(store.get('704')?.workDir).toBe(project) + await text(704, '/projects') + await text(704, worktree) + expect(store.get('704')?.workDir).toBe(worktree) + await text(704, '/new') + expect(store.get('704')?.workDir).toBe(project) + }) + + it('keeps failed-download captions as conversation content and routes permission callbacks', async () => { + store.set('705', 'old', project) + await text(705, '/sessions') + const logError = spyOn(console, 'error').mockImplementation(() => {}) + try { + await text(705, '/resume 1', { photo: true }) + expect(logError).toHaveBeenCalled() + } finally { + logError.mockRestore() + } + expect(store.get('705')?.sessionId).toBe('old') + await eventually(() => expect(messages.some((item) => item.sessionId === 'old' && item.message.content === '/resume 1')).toBe(true)) + broadcast('old', { type: 'permission_request', requestId: 'fixture-request', toolName: 'Bash', input: { command: 'echo fixture' } }) + await eventually(() => expect(texts(705).some((value) => value.includes('fixture-request'))).toBe(true)) + await callback(705, 'permit:fixture-request:yes') + await eventually(() => expect(messages.some((item) => item.message.type === 'permission_response' && item.message.requestId === 'fixture-request')).toBe(true)) + broadcast('old', { type: 'message_complete' }) + await text(705, '/clear') + await eventually(() => expect(messages.some((item) => item.message.content === '/clear')).toBe(true)) + }) + + it('updates model and busy state through existing Skill and model menu wiring', async () => { + store.set('706', 'old', project) + await text(706, '/model fixture-model') + await eventually(() => expect(texts(706).some((value) => value.includes('已切换模型'))).toBe(true)) + await text(706, '/skills') + await eventually(() => expect(texts(706).some((value) => value.includes('当前项目可用 Skills'))).toBe(true)) + await callback(706, 'tgsel:skill:pick:0') + await text(706, '/sessions') + await text(706, '/resume 1') + expect(store.get('706')?.sessionId).toBe('old') + expect(texts(706).at(-1)).toContain('/stop') + }) + + it('accepts desktop permission resolution and restores history after completion', async () => { + store.set('707', 'old', project) + await text(707, 'Need approval') + broadcast('old', { type: 'permission_request', requestId: 'desktop-request', toolName: 'Bash', input: { command: 'echo fixture' } }) + await eventually(() => expect(texts(707).some((value) => value.includes('desktop-request'))).toBe(true)) + broadcast('old', { type: 'permission_resolved', permissionType: 'tool', requestId: 'desktop-request' }) + broadcast('old', { type: 'message_complete' }) + await text(707, '/status') + await eventually(() => expect(texts(707).some((value) => value.includes('old'))).toBe(true)) + await text(707, '/stop') + await eventually(() => expect(messages.some((item) => item.sessionId === 'old' && item.message.type === 'stop_generation')).toBe(true)) + await text(707, '/sessions') + await text(707, '/resume 1') + expect(store.get('707')?.sessionId).toBe('history') + expect(texts(707).at(-1)).toContain('已恢复会话') + }) + + it('delivers a resumed answer and accepts text approval through the same entrypoint', async () => { + store.set('708', 'stream', project) + await text(708, 'Show the result') + await eventually(() => expect(messages.some((item) => item.sessionId === 'stream' && item.message.content === 'Show the result')).toBe(true)) + broadcast('stream', { type: 'status', state: 'thinking', verb: 'Thinking' }) + broadcast('stream', { type: 'thinking', text: 'Checking the previous context' }) + broadcast('stream', { type: 'content_start', blockType: 'text' }) + broadcast('stream', { type: 'content_delta', text: 'Result from the restored session.' }) + broadcast('stream', { type: 'content_start', blockType: 'tool_use' }) + broadcast('stream', { type: 'tool_use_complete' }) + broadcast('stream', { type: 'tool_result' }) + broadcast('stream', { type: 'system_notification', subtype: 'init', data: { model: 'fixture-model' } }) + broadcast('stream', { type: 'permission_request', requestId: 'text-request', toolName: 'Bash', input: { command: 'echo fixture' } }) + await eventually(() => expect(texts(708).some((value) => value.includes('text-request'))).toBe(true)) + await text(708, '/allow text-request') + await eventually(() => expect(messages.some((item) => item.sessionId === 'stream' && item.message.type === 'permission_response' && item.message.requestId === 'text-request')).toBe(true)) + broadcast('stream', { type: 'message_complete' }) + broadcast('stream', { type: 'error', message: 'Fixture turn failed' }) + await eventually(() => expect(texts(708).some((value) => value.includes('Fixture turn failed'))).toBe(true)) + expect(texts(708).some((value) => value.includes('Result from the restored session.'))).toBe(true) + expect(texts(708).some((value) => value.includes('Checking the previous context'))).toBe(true) + }) + + it('starts the registered bot and publishes its menu without external access', async () => { + const gc = spyOn(AttachmentStore.prototype, 'gc').mockResolvedValue({ removed: 0, bytes: 0 }) + const start = spyOn(entry.bot, 'start').mockImplementation(async (options) => { await options?.onStart?.(entry.bot.botInfo) }) + const previousListeners = process.listeners('SIGINT') + try { + entry.startTelegramAdapter() + await eventually(() => expect(apiCalls.some((call) => call.method === 'setMyCommands')).toBe(true)) + expect(gc).toHaveBeenCalledTimes(1) + expect(start).toHaveBeenCalledTimes(1) + const commands = apiCalls.find((call) => call.method === 'setMyCommands')!.payload.commands + expect(commands.some((command: { command: string }) => command.command === 'sessions')).toBe(true) + } finally { + for (const listener of process.listeners('SIGINT')) { + if (!previousListeners.includes(listener)) process.removeListener('SIGINT', listener) + } + start.mockRestore() + gc.mockRestore() + } + }) +}) diff --git a/adapters/telegram/__tests__/session-input.test.ts b/adapters/telegram/__tests__/session-input.test.ts new file mode 100644 index 00000000..6f08adca --- /dev/null +++ b/adapters/telegram/__tests__/session-input.test.ts @@ -0,0 +1,63 @@ +import { describe, expect, it, mock } from 'bun:test' +import { + registerTelegramSessionCommands, + tryHandleTelegramSessionInput, +} from '../commands.js' + +function createSessionRoutes() { + return { + startNewSession: mock(async () => {}), + showProjectPicker: mock(async () => {}), + showResumeProjectPicker: mock(async () => {}), + handleSessionInput: mock(async () => false), + } +} + +describe('Telegram session input routing', () => { + it('routes registered session commands through the normal message pipeline', async () => { + const handlers = new Map unknown>() + const routeInput = mock(async () => {}) + registerTelegramSessionCommands({ + command: (name, handler) => handlers.set(name, handler), + }, routeInput) + + const ctx = { match: '2' } as any + await handlers.get('resume')!(ctx) + expect(routeInput).toHaveBeenCalledWith(ctx, '/resume 2') + await handlers.get('sessions')!({ match: '/work/my project' } as any) + expect(routeInput).toHaveBeenLastCalledWith(expect.anything(), '/sessions /work/my project') + expect([...handlers.keys()]).toEqual(['new', 'projects', 'sessions', 'resume']) + }) + + it('keeps the existing resume menu and sends numbered selection to shared history', async () => { + const routes = createSessionRoutes() + expect(await tryHandleTelegramSessionInput('42', '/resume', false, routes)).toBe(true) + expect(routes.showResumeProjectPicker).toHaveBeenCalledWith('42') + expect(routes.handleSessionInput).not.toHaveBeenCalled() + + routes.handleSessionInput.mockResolvedValueOnce(true) + expect(await tryHandleTelegramSessionInput('42', '/resume 2', false, routes)).toBe(true) + expect(routes.handleSessionInput).toHaveBeenCalledWith('42', '/resume 2') + + routes.handleSessionInput.mockResolvedValueOnce(true) + expect(await tryHandleTelegramSessionInput('42', '2', false, routes)).toBe(true) + expect(routes.handleSessionInput).toHaveBeenLastCalledWith('42', '2') + }) + + it('leaves attachment captions out of all session commands and pickers', async () => { + const routes = createSessionRoutes() + for (const text of ['/new /tmp/repo', '/projects', '/sessions', '/resume', '/resume 1', '1']) { + expect(await tryHandleTelegramSessionInput('42', text, true, routes)).toBe(false) + } + for (const route of Object.values(routes)) expect(route).not.toHaveBeenCalled() + }) + + it('handles new and projects inside the same serialized message path', async () => { + const routes = createSessionRoutes() + expect(await tryHandleTelegramSessionInput('42', '/new /work/my project', false, routes)).toBe(true) + expect(routes.startNewSession).toHaveBeenCalledWith('42', '/work/my project') + expect(await tryHandleTelegramSessionInput('42', '/projects', false, routes)).toBe(true) + expect(routes.showProjectPicker).toHaveBeenCalledWith('42') + expect(await tryHandleTelegramSessionInput('42', 'ordinary text', false, routes)).toBe(false) + }) +}) diff --git a/adapters/telegram/commands.ts b/adapters/telegram/commands.ts index 021ea185..077d4e56 100644 --- a/adapters/telegram/commands.ts +++ b/adapters/telegram/commands.ts @@ -1,4 +1,7 @@ import { formatImHelp } from '../common/format.js' +import { listProjectSessionHistory, restoreSelectedSession } from '../common/session-selection.js' +import type { SessionEntry } from '../common/session-store.js' +import type { ServerMessage } from '../common/ws-bridge.js' import { formatPermissionDecisionStatus, type PermissionDecision, @@ -69,11 +72,15 @@ export type TelegramCommandControllerDeps = { isAllowedUser: (userId: number) => boolean ensureExistingSession: (chatId: string) => Promise<{ sessionId: string; workDir: string } | null> clearTransientChatState: (chatId: string) => void + clearOtherSelections: (chatId: string) => void + isBusy: (chatId: string) => boolean + getStoredSession: (chatId: string) => SessionEntry | null setStoredSession: (chatId: string, sessionId: string, workDir: string) => void deleteStoredSession: (chatId: string) => void resetBridgeSession: (chatId: string) => void - connectBridgeSession: (chatId: string, sessionId: string) => void - onBridgeServerMessage: (chatId: string) => void + connectBridgeSession: (chatId: string, sessionId: string) => boolean + onBridgeServerMessage: (chatId: string, handler: (msg: ServerMessage) => void | Promise) => void + handleServerMessage: (chatId: string, msg: ServerMessage) => void | Promise waitForBridgeOpen: (chatId: string) => Promise sendUserMessage: (chatId: string, content: string) => boolean setRuntimeModel: RuntimeModelSetter @@ -216,7 +223,7 @@ export type TelegramRuntimeCommandControllerDeps = { defaultWorkDir: string bridge: { resetSession: (chatId: string) => void - connectSession: (chatId: string, sessionId: string) => void + connectSession: (chatId: string, sessionId: string) => boolean onServerMessage: (chatId: string, handler: (msg: unknown) => void | Promise) => void waitForOpen: (chatId: string) => Promise sendUserMessage: (chatId: string, content: string) => boolean @@ -228,14 +235,18 @@ export type TelegramRuntimeCommandControllerDeps = { ) => boolean } sessionStore: { + get: (chatId: string) => SessionEntry | null set: (chatId: string, sessionId: string, workDir: string) => void delete: (chatId: string) => void } isAllowedUser: (userId: number) => boolean ensureExistingSession: (chatId: string) => Promise<{ sessionId: string; workDir: string } | null> clearTransientChatState: (chatId: string) => void + clearOtherSelections: (chatId: string) => void + isBusy: (chatId: string) => boolean handleServerMessage: (chatId: string, msg: unknown) => void | Promise setRuntimeModel: (chatId: string, modelId: string) => void + setRuntimeBusy: (chatId: string) => void } export function createTelegramRuntimeCommandController( @@ -248,16 +259,24 @@ export function createTelegramRuntimeCommandController( isAllowedUser: deps.isAllowedUser, ensureExistingSession: deps.ensureExistingSession, clearTransientChatState: deps.clearTransientChatState, + clearOtherSelections: deps.clearOtherSelections, + isBusy: deps.isBusy, + getStoredSession: (chatId) => deps.sessionStore.get(chatId), setStoredSession: (chatId, sessionId, workDir) => deps.sessionStore.set(chatId, sessionId, workDir), deleteStoredSession: (chatId) => deps.sessionStore.delete(chatId), resetBridgeSession: (chatId) => deps.bridge.resetSession(chatId), connectBridgeSession: (chatId, sessionId) => deps.bridge.connectSession(chatId, sessionId), - onBridgeServerMessage: (chatId) => deps.bridge.onServerMessage( + onBridgeServerMessage: (chatId, handler) => deps.bridge.onServerMessage( chatId, - (msg) => deps.handleServerMessage(chatId, msg), + (msg) => handler(msg as ServerMessage), ), + handleServerMessage: deps.handleServerMessage, waitForBridgeOpen: (chatId) => deps.bridge.waitForOpen(chatId), - sendUserMessage: (chatId, content) => deps.bridge.sendUserMessage(chatId, content), + sendUserMessage: (chatId, content) => { + const sent = deps.bridge.sendUserMessage(chatId, content) + if (sent) deps.setRuntimeBusy(chatId) + return sent + }, setRuntimeModel: deps.setRuntimeModel, }) return { @@ -282,12 +301,53 @@ export function registerTelegramExtendedCommands( ): void { bot.command('start', (ctx) => void controller.sendHelp(ctx)) bot.command('help', (ctx) => void controller.sendHelp(ctx)) - bot.command('resume', (ctx) => void controller.handleResumeCommand(ctx)) bot.command('provider', (ctx) => void controller.handleProviderCommand(ctx)) bot.command('model', (ctx) => void controller.handleModelCommand(ctx)) bot.command('skills', (ctx) => void controller.handleSkillsCommand(ctx)) } +/** Session commands use the same authorization, deduplication and queue as messages. */ +export function registerTelegramSessionCommands( + bot: TelegramCommandRegistrar, + routeInput: (ctx: TelegramCommandContext, text: string) => Promise, +): void { + for (const command of ['new', 'projects', 'sessions', 'resume']) { + bot.command(command, (ctx) => { + const query = getCommandMatchText(ctx) + return routeInput(ctx, `/${command}${query ? ` ${query}` : ''}`) + }) + } +} + +export async function tryHandleTelegramSessionInput( + chatId: string, + text: string, + hasAttachments: boolean, + deps: { + startNewSession: (chatId: string, query?: string) => Promise + showProjectPicker: (chatId: string) => Promise + showResumeProjectPicker: (chatId: string) => Promise + handleSessionInput: (chatId: string, text: string) => Promise + }, +): Promise { + if (hasAttachments) return false + const trimmed = text.trim() + const newCommand = /^\/new(?:\s+([\s\S]+))?$/.exec(trimmed) + if (newCommand) { + await deps.startNewSession(chatId, newCommand[1]) + return true + } + if (trimmed === '/projects') { + await deps.showProjectPicker(chatId) + return true + } + if (trimmed === '/resume') { + await deps.showResumeProjectPicker(chatId) + return true + } + return await deps.handleSessionInput(chatId, text) +} + export async function tryHandleTelegramSelectionCallback( data: string, ctx: TelegramCommandContext, @@ -529,6 +589,8 @@ export function createTelegramCommandController(deps: TelegramCommandControllerD } const showResumeProjectPicker = async (chatId: string): Promise => { + deps.clearOtherSelections(chatId) + pendingSelections.delete(chatId) try { const projects = await deps.httpClient.listRecentProjects() if (projects.length === 0) { @@ -563,12 +625,7 @@ export function createTelegramCommandController(deps: TelegramCommandControllerD const chatId = getCallbackChatId(ctx) if (!chatId) return try { - const { sessions } = await deps.httpClient.listSessions({ - project: project.value, - limit: 50, - offset: 0, - }) - const resumableSessions = sessions.filter((session) => session.workDir) + const resumableSessions = await listProjectSessionHistory(deps.httpClient, project.value) if (resumableSessions.length === 0) { pendingSelections.delete(chatId) await ctx.editMessageText(`没有可恢复会话:${project.label}`) @@ -598,25 +655,25 @@ export function createTelegramCommandController(deps: TelegramCommandControllerD return } - deps.resetBridgeSession(chatId) - deps.clearTransientChatState(chatId) - deps.setStoredSession(chatId, item.value, workDir) - deps.connectBridgeSession(chatId, item.value) - deps.onBridgeServerMessage(chatId) - const opened = await deps.waitForBridgeOpen(chatId) - if (!opened) { - deps.deleteStoredSession(chatId) - await ctx.editMessageText('⚠️ 恢复会话时连接服务器超时,请重试。') - return - } - - pendingSelections.delete(chatId) - await ctx.editMessageText([ - `✅ 已恢复会话:${item.label}`, - workDir, - '', - '可以继续发送消息。', - ].join('\n')) + const result = await restoreSelectedSession({ + httpClient: deps.httpClient, + bridge: { + resetSession: deps.resetBridgeSession, + connectSession: deps.connectBridgeSession, + onServerMessage: deps.onBridgeServerMessage, + waitForOpen: deps.waitForBridgeOpen, + }, + sessionStore: { + get: deps.getStoredSession, + set: deps.setStoredSession, + delete: deps.deleteStoredSession, + }, + onServerMessage: deps.handleServerMessage, + clearTransientState: deps.clearTransientChatState, + isBusy: deps.isBusy, + }, chatId, { id: item.value, workDir, title: item.label }) + if (result.ok) pendingSelections.delete(chatId) + await ctx.editMessageText(result.message) } const handleSelectionCallback = async ( @@ -699,7 +756,7 @@ export function buildTelegramHelpText(): string { formatImHelp(), '', 'Telegram 扩展命令:', - '/resume — 恢复历史会话', + '/resume — 按钮选择项目与历史会话', '/provider — 切换 Provider', '/model [model] — 查看或切换模型', '/skills — 查看当前项目可用 Skills', diff --git a/adapters/telegram/index.ts b/adapters/telegram/index.ts index f1f9b4c6..480ab24d 100644 --- a/adapters/telegram/index.ts +++ b/adapters/telegram/index.ts @@ -29,6 +29,8 @@ import { import { SessionStore } from '../common/session-store.js' import { createAdapterClient } from '../common/adapter-client.js' import { restoreStoredSessionBinding } 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' import { TelegramMediaService } from './media.js' import { AttachmentStore } from '../common/attachment/attachment-store.js' @@ -38,7 +40,7 @@ import { ImageBlockWatcher } from '../common/attachment/image-block-watcher.js' import type { PendingUpload } from '../common/attachment/attachment-types.js' import { sendSafeOutboundImage } from '../common/attachment/outbound-image.js' import { syncTelegramBotCommands } from './menu.js' -import { createTelegramRuntimeCommandController, registerAuthorizedTelegramCommand, registerTelegramExtendedCommands, shouldProcessTelegramMessage, tryHandleTelegramSelectionCallback } from './commands.js' +import { createTelegramRuntimeCommandController, registerAuthorizedTelegramCommand, registerTelegramExtendedCommands, registerTelegramSessionCommands, shouldProcessTelegramMessage, tryHandleTelegramSelectionCallback, tryHandleTelegramSessionInput } from './commands.js' // ---------- init ---------- @@ -48,7 +50,7 @@ if (!config.telegram.botToken) { process.exit(1) } -const bot = new Bot(config.telegram.botToken) +export const bot = new Bot(config.telegram.botToken) const bridge = new WsBridge(config.serverUrl, 'tg') const streamDelivery = new TelegramStreamDelivery(bot.api) const dedup = new MessageDedup() @@ -56,9 +58,6 @@ const sessionStore = new SessionStore() const { httpClient, defaultWorkDir } = createAdapterClient(config, config.telegram) const attachmentStore = new AttachmentStore() const media = new TelegramMediaService(bot, attachmentStore) -attachmentStore.gc().catch((err) => { - console.warn('[Telegram] AttachmentStore.gc failed:', err instanceof Error ? err.message : err) -}) const accumulatedThinkingText = new Map() // Track chats waiting for project selection @@ -84,7 +83,31 @@ type ChatRuntimeState = { pendingPermissionCount: number } -const commandController = createTelegramRuntimeCommandController({ botApi: bot.api, httpClient, defaultWorkDir, bridge, sessionStore, ensureExistingSession, clearTransientChatState, isAllowedUser: (userId) => isAllowedUser('telegram', userId), handleServerMessage: (chatId, msg) => handleServerMessage(chatId, msg as ServerMessage), setRuntimeModel: (chatId, modelId) => { getRuntimeState(chatId).model = modelId } }) +const isChatBusy = (chatId: string) => getRuntimeState(chatId).state !== 'idle' || Boolean(pendingPermissions.get(chatId)?.size) +const commandController = createTelegramRuntimeCommandController({ + botApi: bot.api, httpClient, defaultWorkDir, bridge, sessionStore, + ensureExistingSession, clearTransientChatState, + clearOtherSelections: (chatId) => { + pendingProjectSelection.delete(chatId) + sessionSelection.clear(chatId) + }, + isBusy: isChatBusy, + isAllowedUser: (userId) => isAllowedUser('telegram', userId), + handleServerMessage: (chatId, msg) => handleServerMessage(chatId, msg as ServerMessage), + setRuntimeModel: (chatId, modelId) => { getRuntimeState(chatId).model = modelId }, + setRuntimeBusy: (chatId) => { getRuntimeState(chatId).state = 'thinking' }, +}) +const sessionSelection = new SessionSelectionController({ + httpClient, bridge, sessionStore, + sendNotice: async (chatId, text) => { await bot.api.sendMessage(Number(chatId), text) }, + onServerMessage: handleServerMessage, + clearTransientState: clearTransientChatState, + clearProjectSelection: (chatId) => { + pendingProjectSelection.delete(chatId) + commandController.clearPendingSelections(chatId) + }, + isBusy: isChatBusy, +}) // ---------- helpers ---------- @@ -231,6 +254,9 @@ async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelection.clear(chatId) + commandController.clearPendingSelections(chatId) + pendingProjectSelection.delete(chatId) const numericChatId = Number(chatId) try { const projects = await httpClient.listRecentProjects() @@ -276,6 +302,7 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< const numericChatId = Number(chatId) const runtime = getRuntimeState(chatId) + if (syncImPermissionState(chatId, msg, runtime, pendingPermissions)) return switch (msg.type) { case 'connected': break @@ -402,6 +429,7 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< // ---------- bot handlers ---------- registerTelegramExtendedCommands(bot, commandController) +registerTelegramSessionCommands(bot, (ctx, text) => routeUserMessage(ctx as Context, text, [])) /** Reset session state and start a new session for chatId. * If `query` is provided, match a project by index or name; @@ -413,6 +441,7 @@ async function startNewSession(chatId: string, query?: string): Promise { sessionStore.delete(chatId) streamDelivery.clear(chatId) pendingProjectSelection.delete(chatId) + sessionSelection.clear(chatId) commandController.clearPendingSelections(chatId) pendingPermissions.delete(chatId) runtimeStates.delete(chatId) @@ -454,19 +483,8 @@ async function startNewSession(chatId: string, query?: string): Promise { const isAuthorizedTelegramUser = (userId: number) => isAllowedUser('telegram', userId) -registerAuthorizedTelegramCommand(bot, 'new', isAuthorizedTelegramUser, async (ctx) => { - const chatId = String(ctx.chat.id) - const query = typeof ctx.match === 'string' ? ctx.match.trim() : undefined - await startNewSession(chatId, query || undefined) -}) - -registerAuthorizedTelegramCommand(bot, 'projects', isAuthorizedTelegramUser, async (ctx) => { - const chatId = String(ctx.chat.id) - await showProjectPicker(chatId) -}) - registerAuthorizedTelegramCommand(bot, 'stop', isAuthorizedTelegramUser, (ctx) => { - const chatId = String(ctx.chat.id) + const chatId = String(ctx.chat!.id) void (async () => { const stored = await ensureExistingSession(chatId) if (!stored) { @@ -479,12 +497,12 @@ registerAuthorizedTelegramCommand(bot, 'stop', isAuthorizedTelegramUser, (ctx) = }) registerAuthorizedTelegramCommand(bot, 'status', isAuthorizedTelegramUser, async (ctx) => { - const chatId = String(ctx.chat.id) + const chatId = String(ctx.chat!.id) await ctx.reply(await buildStatusText(chatId)) }) registerAuthorizedTelegramCommand(bot, 'clear', isAuthorizedTelegramUser, (ctx) => { - const chatId = String(ctx.chat.id) + const chatId = String(ctx.chat!.id) void (async () => { const stored = await ensureExistingSession(chatId) if (!stored) { @@ -497,6 +515,7 @@ registerAuthorizedTelegramCommand(bot, 'clear', isAuthorizedTelegramUser, (ctx) await ctx.reply('⚠️ 无法发送 /clear,请先发送 /new 重新连接会话。') return } + getRuntimeState(chatId).state = 'thinking' await ctx.reply('🧹 已清空当前会话上下文。') })() }) @@ -532,8 +551,12 @@ async function routeUserMessage( return } - enqueue(chatId, async () => { - const permissionDecision = attachments.length === 0 + // Captions remain conversation content even when an attachment download fails. + const hasAttachments = attachments.length > 0 || Boolean( + ctx.message?.photo || ctx.message?.document || ctx.message?.video || ctx.message?.audio || ctx.message?.voice, + ) + await enqueue(chatId, async () => { + const permissionDecision = !hasAttachments ? parsePermissionCommand(text, pendingPermissions.get(chatId)) : null if (permissionDecision) { @@ -541,7 +564,14 @@ async function routeUserMessage( return } - if (pendingProjectSelection.has(chatId)) { + if (await tryHandleTelegramSessionInput(chatId, text, hasAttachments, { + startNewSession, + showProjectPicker, + showResumeProjectPicker: commandController.showResumeProjectPicker, + handleSessionInput: (id, input) => sessionSelection.handleInput(id, input), + })) return + + if (!hasAttachments && pendingProjectSelection.has(chatId)) { if (text.trim()) await startNewSession(chatId, text.trim()) return } @@ -553,6 +583,8 @@ async function routeUserMessage( const sent = bridge.sendUserMessage(chatId, effective, attachments.length ? attachments : undefined) if (!sent) { await bot.api.sendMessage(Number(chatId), '⚠️ 消息发送失败,连接可能已断开。请发送 /new 重新开始。') + } else { + getRuntimeState(chatId).state = 'thinking' } }) } @@ -645,33 +677,43 @@ bot.on( ) bot.on('callback_query:data', async (ctx) => { + if (!ctx.from || ctx.chat?.type !== 'private') return + if (!dedup.tryRecord(`telegram:callback:${ctx.callbackQuery.id}`)) return const data = ctx.callbackQuery.data - if (await tryHandleTelegramSelectionCallback(data, ctx, commandController)) return + await enqueue(String(ctx.chat.id), async () => { + if (await tryHandleTelegramSelectionCallback(data, ctx, commandController)) return - if (!data.startsWith('permit:')) return + if (!data.startsWith('permit:')) return - const decision = parsePermitCallbackData(data) - if (!decision) return - await commandController.handlePermissionCallback(ctx, decision, pendingPermissions, (chatId) => getRuntimeState(chatId).pendingPermissionCount = Math.max(0, getRuntimeState(chatId).pendingPermissionCount - 1)) + const decision = parsePermitCallbackData(data) + if (!decision) return + await commandController.handlePermissionCallback(ctx, decision, pendingPermissions, (chatId) => getRuntimeState(chatId).pendingPermissionCount = Math.max(0, getRuntimeState(chatId).pendingPermissionCount - 1)) + }) }) // ---------- start ---------- -console.log('[Telegram] Starting bot...') -console.log(`[Telegram] Server: ${config.serverUrl}`) -console.log(`[Telegram] Allowed users: ${config.telegram.allowedUsers.length === 0 ? 'paired users only' : config.telegram.allowedUsers.join(', ')}`) - -void syncTelegramBotCommands(bot.api).then(() => console.log('[Telegram] Command menu synced')).catch((err) => console.warn('[Telegram] Command menu sync failed:', err instanceof Error ? err.message : err)) - -bot.start({ - onStart: () => console.log('[Telegram] Bot is running!'), -}) - -// Graceful shutdown -process.on('SIGINT', () => { - console.log('[Telegram] Shutting down...') - bot.stop() +export function stopTelegramAdapter(): void { + if (bot.isRunning()) void bot.stop() bridge.destroy() dedup.destroy() - process.exit(0) -}) +} + +export function startTelegramAdapter(): void { + console.log('[Telegram] Starting bot...') + console.log(`[Telegram] Server: ${config.serverUrl}`) + console.log(`[Telegram] Allowed users: ${config.telegram.allowedUsers.length === 0 ? 'paired users only' : config.telegram.allowedUsers.join(', ')}`) + void attachmentStore.gc().catch((err) => { + console.warn('[Telegram] AttachmentStore.gc failed:', err instanceof Error ? err.message : err) + }) + void syncTelegramBotCommands(bot.api).then(() => console.log('[Telegram] Command menu synced')).catch((err) => console.warn('[Telegram] Command menu sync failed:', err instanceof Error ? err.message : err)) + void bot.start({ onStart: () => console.log('[Telegram] Bot is running!') }) + process.once('SIGINT', () => { + console.log('[Telegram] Shutting down...') + stopTelegramAdapter() + process.exit(0) + }) +} + +// Desktop's shared sidecar imports this module with its explicit adapter flag. +if (import.meta.main || process.argv.includes('--telegram')) startTelegramAdapter() diff --git a/adapters/telegram/menu.ts b/adapters/telegram/menu.ts index f590c90a..ae63fcc6 100644 --- a/adapters/telegram/menu.ts +++ b/adapters/telegram/menu.ts @@ -47,6 +47,7 @@ export const TELEGRAM_BOT_COMMANDS: TelegramBotCommand[] = [ { command: 'help', description: '查看帮助' }, { command: 'new', description: '新建会话或切换项目' }, { command: 'projects', description: '查看最近项目' }, + { command: 'sessions', description: '查看项目历史会话' }, { command: 'resume', description: '恢复历史会话' }, { command: 'status', description: '查看当前状态' }, { command: 'clear', description: '清空当前上下文' }, diff --git a/adapters/wechat/index.ts b/adapters/wechat/index.ts index 15db8417..be7f219e 100644 --- a/adapters/wechat/index.ts +++ b/adapters/wechat/index.ts @@ -16,6 +16,8 @@ import { parsePermissionCommand, } from '../common/permission.js' 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 { isAllowedUser, tryPair } from '../common/pairing.js' @@ -66,6 +68,21 @@ attachmentStore.gc().catch((err) => { console.warn('[WeChat] AttachmentStore.gc failed:', err instanceof Error ? err.message : err) }) +const sessionSelectionController = new SessionSelectionController({ + httpClient, + bridge, + sessionStore, + sendNotice: async (chatId, text) => { await sendText(chatId, text) }, + onServerMessage: handleServerMessage, + clearTransientState: clearTransientChatState, + clearProjectSelection: (chatId) => { pendingProjectSelection.delete(chatId) }, + isBusy: (chatId) => { + const runtime = getRuntimeState(chatId) + return runtime.state !== 'idle' || runtime.pendingPermissionCount > 0 + || (pendingPermissions.get(chatId)?.size ?? 0) > 0 + }, +}) + type ChatRuntimeState = { state: 'idle' | 'thinking' | 'streaming' | 'tool_executing' | 'permission_pending' verb?: string @@ -277,6 +294,7 @@ async function ensureSession(chatId: string): Promise { } async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { bridge.resetSession(chatId) clearTransientChatState(chatId) @@ -297,6 +315,7 @@ async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { const projects = await httpClient.listRecentProjects() if (projects.length === 0) { @@ -315,6 +334,7 @@ async function showProjectPicker(chatId: string): Promise { } async function startNewSession(chatId: string, query?: string): Promise { + sessionSelectionController.clear(chatId) bridge.resetSession(chatId) sessionStore.delete(chatId) clearTransientChatState(chatId) @@ -351,10 +371,12 @@ async function startNewSession(chatId: string, query?: string): Promise { async function handleServerMessage(chatId: string, msg: ServerMessage): Promise { const runtime = getRuntimeState(chatId) + if (syncImPermissionState(chatId, msg, runtime, pendingPermissions)) return switch (msg.type) { case 'connected': break + case 'status': runtime.state = msg.state runtime.verb = typeof msg.verb === 'string' ? msg.verb : undefined @@ -491,6 +513,7 @@ async function routeUserMessage(message: WechatMessage): Promise { } clearTransientChatState(chatId) const sent = bridge.sendUserMessage(chatId, '/clear') + if (sent) getRuntimeState(chatId).state = 'thinking' await sendText(chatId, sent ? '已清空当前会话上下文。' : '无法发送 /clear,请先发送 /new 重新连接会话。') return } @@ -511,6 +534,8 @@ async function routeUserMessage(message: WechatMessage): Promise { await sendText(chatId, sent ? `${formatPermissionDecisionStatus(permissionDecision)}。` : '权限响应发送失败,请检查会话状态。') return } + if (!hasAttachments && await sessionSelectionController.handleInput(chatId, text)) return + if (!hasAttachments && pendingProjectSelection.has(chatId)) { await startNewSession(chatId, text) return @@ -523,6 +548,7 @@ async function routeUserMessage(message: WechatMessage): Promise { if (!effectiveText && attachments.length === 0) return typingController.start(chatId) const sent = bridge.sendUserMessage(chatId, effectiveText, attachments.length ? attachments : undefined) + if (sent) getRuntimeState(chatId).state = 'thinking' if (!sent) await sendText(chatId, '消息发送失败,连接可能已断开。请发送 /new 重新开始。') }) } @@ -612,9 +638,9 @@ function redactChatId(chatId: string): string { console.log('[WeChat] Starting adapter...') console.log(`[WeChat] Account: ${accountId}`) -void pollLoop() +if (import.meta.main || process.argv.includes('--wechat')) void pollLoop() -process.on('SIGINT', () => { +if (import.meta.main || process.argv.includes('--wechat')) process.on('SIGINT', () => { console.log('[WeChat] Shutting down...') stopped = true typingController.destroy() @@ -622,3 +648,5 @@ process.on('SIGINT', () => { dedup.destroy() process.exit(0) }) + +export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, routeUserMessage, startNewSession, typingController } diff --git a/adapters/whatsapp/index.ts b/adapters/whatsapp/index.ts index d1d07065..84ce03ae 100644 --- a/adapters/whatsapp/index.ts +++ b/adapters/whatsapp/index.ts @@ -26,6 +26,8 @@ import { type PermissionDecision, } from '../common/permission.js' 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 { isAllowedUser, tryPair } from '../common/pairing.js' @@ -82,6 +84,21 @@ const runtimeStates = new Map() const pendingPermissions = new Map>() const imageWatchers = new Map() +const sessionSelectionController = new SessionSelectionController({ + httpClient, + bridge, + sessionStore, + sendNotice: async (chatId, text) => { await sendWhatsAppText(chatId, text) }, + onServerMessage: handleServerMessage, + clearTransientState: clearTransientChatState, + clearProjectSelection: (chatId) => { pendingProjectSelection.delete(chatId) }, + isBusy: (chatId) => { + const runtime = getRuntimeState(chatId) + return runtime.state !== 'idle' || runtime.pendingPermissionCount > 0 + || (pendingPermissions.get(chatId)?.size ?? 0) > 0 + }, +}) + type ChatRuntimeState = { state: 'idle' | 'thinking' | 'streaming' | 'tool_executing' | 'permission_pending' verb?: string @@ -202,6 +219,7 @@ async function ensureSession(chatId: string): Promise { } async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { bridge.resetSession(chatId) const sessionId = await httpClient.createSession(workDir) @@ -221,6 +239,7 @@ async function createSessionForChat(chatId: string, workDir: string): Promise { + sessionSelectionController.clear(chatId) try { const projects = await httpClient.listRecentProjects() if (projects.length === 0) { @@ -270,11 +289,13 @@ async function flushAccumulatedText(chatId: string): Promise { async function handleServerMessage(chatId: string, msg: ServerMessage): Promise { const runtime = getRuntimeState(chatId) + if (syncImPermissionState(chatId, msg, runtime, pendingPermissions)) return switch (msg.type) { case 'connected': break + case 'status': runtime.state = msg.state runtime.verb = typeof msg.verb === 'string' ? msg.verb : undefined @@ -364,6 +385,7 @@ async function handleServerMessage(chatId: string, msg: ServerMessage): Promise< } async function startNewSession(chatId: string, query?: string): Promise { + sessionSelectionController.clear(chatId) bridge.resetSession(chatId) sessionStore.delete(chatId) accumulatedText.delete(chatId) @@ -460,6 +482,7 @@ async function routeUserMessage( } clearTransientChatState(chatId) const sent = bridge.sendUserMessage(chatId, '/clear') + if (sent) getRuntimeState(chatId).state = 'thinking' await sendWhatsAppText(chatId, sent ? '已清空当前会话上下文。' : '无法发送 /clear,请先发送 /new 重新连接会话。') return } @@ -472,7 +495,9 @@ async function routeUserMessage( return } - if (pendingProjectSelection.has(chatId)) { + if (attachments.length === 0 && await sessionSelectionController.handleInput(chatId, command)) return + + if (attachments.length === 0 && pendingProjectSelection.has(chatId)) { if (text.trim()) await startNewSession(chatId, text.trim()) return } @@ -482,6 +507,7 @@ async function routeUserMessage( const effective = text || (attachments.length > 0 ? '(用户发送了附件)' : '') if (!effective && attachments.length === 0) return const sent = bridge.sendUserMessage(chatId, effective, attachments.length ? attachments : undefined) + if (sent) getRuntimeState(chatId).state = 'thinking' if (!sent) { await sendWhatsAppText(chatId, '消息发送失败,连接可能已断开。请发送 /new 重新开始。') } @@ -567,13 +593,17 @@ async function handleIncomingMessage(message: proto.IWebMessageInfo): Promise { if (reconnectTimer) { clearTimeout(reconnectTimer) reconnectTimer = null } - sock = await createWhatsAppSocket({ authDir }) - media = new WhatsAppMediaService(sock, attachmentStore) + useWhatsAppSocket(await createWhatsAppSocket({ authDir })) sock.ev.on('messages.upsert', ({ type, messages }) => { if (type !== 'notify') return @@ -617,9 +647,9 @@ console.log(`[WhatsApp] Server: ${config.serverUrl}`) console.log(`[WhatsApp] Auth dir: ${authDir}`) console.log(`[WhatsApp] Allowed users: ${config.whatsapp.allowedUsers.length === 0 ? 'paired users only' : config.whatsapp.allowedUsers.join(', ')}`) -await startSocket() +if (import.meta.main || process.argv.includes('--whatsapp')) await startSocket() -process.on('SIGINT', () => { +if (import.meta.main || process.argv.includes('--whatsapp')) process.on('SIGINT', () => { console.log('[WhatsApp] Shutting down...') shuttingDown = true if (reconnectTimer) clearTimeout(reconnectTimer) @@ -628,3 +658,5 @@ process.on('SIGINT', () => { dedup.destroy() process.exit(0) }) + +export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, routeUserMessage, startNewSession } diff --git a/docs/en/im/index.md b/docs/en/im/index.md index d6ac83c6..8b798ee2 100644 --- a/docs/en/im/index.md +++ b/docs/en/im/index.md @@ -16,6 +16,7 @@ The chat partner is a bot or account you bound yourself. Messages reach your loc ## What you get - **The same session, continued.** Messages sent from your phone enter the Claude Code session on your computer, where file edits, commands, and reads really happen. +- **Resume history.** `/sessions` lists history in the current project; `/sessions ` opens another project. Use `/resume ` to continue, including sessions created on Desktop. - **Project switching.** `/projects` lists recent projects and switches to the one you pick; `/new` starts a fresh session. - **Permission approval.** When Claude wants to write a file or run a risky command, the request is pushed to the chat. Feishu and DingTalk send interactive cards, Telegram sends buttons, and every other platform expects a text reply. - **Status and stop.** `/status` reports the current project, model, and run state; `/stop` interrupts the current turn. @@ -63,6 +64,16 @@ Paired accounts appear under **Paired Users**, where **Unbind** revokes one of t Later messages in the same chat reuse that session, and the mapping survives a Desktop restart. `/new` changes the directory; `/clear` empties the context while keeping the project binding. +## Find an old session from your phone + +Send `/sessions` to list history for your current project. Without a session binding, it starts with a project picker. Use `/sessions projects` to choose a different project, or `/sessions ` to open one directly. Ambiguous names show a list to choose from. + +Each page shows eight sessions with their title, update time, message count, and a marker for the current session. Reply with a number shown on that page or send `/resume ` to continue the selected conversation. Use `/sessions next` and `/sessions prev` to turn pages. Accessible worktree sessions are grouped under their project and resume in their original working directory. + +Browsing preserves the current binding. `/cancel` exits selection; ordinary chat text also leaves the history picker and continues your current conversation. Lists expire after 15 minutes. Finish pending approvals or send `/stop` and wait for the current turn to stop before switching. If a session was deleted, its directory is unavailable, or connecting fails, the original binding is retained. + +Telegram also keeps its `/resume` project and session button menu. The text commands above work on all platforms. + ## Common commands Entry points differ slightly per platform — Feishu can expose commands as a bot menu — but these work everywhere: @@ -71,6 +82,9 @@ Entry points differ slightly per platform — Feishu can expose commands as a bo - `/status` — current project, model, and run state - `/projects` — list recent projects and switch - `/new` — start a new session, optionally with a project number or path +- `/sessions [project]` — list old sessions; `/sessions projects` chooses another project +- `/resume ` — continue a session from the list +- `/cancel` — exit selection and keep the current session - `/clear` — clear context, keep the project binding - `/stop` — stop the current generation diff --git a/docs/im/index.md b/docs/im/index.md index 5896c47a..39803578 100644 --- a/docs/im/index.md +++ b/docs/im/index.md @@ -16,6 +16,7 @@ order: 0 ## 接进来之后能做什么 - **接着聊同一条会话**:手机上发的消息进的是本机的 Claude Code 会话,改文件、跑命令、读代码都在你电脑上真实发生。 +- **恢复旧会话**:发 `/sessions` 查看当前项目的历史,或 `/sessions <项目名或绝对路径>` 找到其他项目,再用 `/resume <编号>` 接着聊。桌面端创建的会话也能选。 - **换项目**:发 `/projects` 列出最近用过的项目,回复编号或路径就切过去;`/new` 直接开一条新会话。 - **批权限**:Claude 要写文件或执行高风险命令时,会把请求推到 IM 里。飞书和钉钉是可点的卡片,Telegram 是按钮,其余平台回复一条文本命令。 - **看状态、叫停**:`/status` 看当前项目、模型和运行状态,`/stop` 中断正在跑的这一轮。 @@ -63,6 +64,16 @@ order: 0 同一个 IM 聊天窗口后续的消息会复用同一条会话,桌面端重启后也能接回去。想换目录发 `/new`,想清空上下文但保留项目发 `/clear`。 +## 在手机上找回旧会话 + +发 `/sessions` 查看当前绑定项目的历史。如果还没有绑定会话,会先显示项目列表;已经绑定时,也可以发 `/sessions projects` 选择其他项目。发 `/sessions <项目名或绝对路径>` 可以直接查看指定项目,重名项目会列出来让你选。 + +会话列表显示标题、更新时间、消息数和当前会话标记,每页 8 条。回复当前页的编号,或发 `/resume <编号>`,之后的消息就会继续那条旧会话。用 `/sessions next` 和 `/sessions prev` 翻页。列表包含同一项目下可访问的 worktree 会话,选中后仍在原工作目录继续。 + +查看列表不会打断当前会话;发 `/cancel` 退出选择,直接发普通聊天文字也会退出历史选择并继续当前对话。列表 15 分钟后过期,需要重新查询。会话正在生成或等待审批时,请先处理审批或 `/stop`,等停止后再切换。旧会话已删除、目录已移除或无法连接时,原来的绑定会保留。 + +Telegram 还保留 `/resume` 的项目、会话按钮菜单;各平台都支持上面的文本命令。 + ## 允许访问的项目目录决定它能碰哪些项目 「默认项目」只决定新会话开在哪,**不是**访问边界。真正的边界是「允许访问的项目目录」:机器人只能列出、打开并在这些目录内的项目里开会话,`/projects`、`/sessions` 和按名字、绝对路径选项目都受它约束。 @@ -98,6 +109,9 @@ order: 0 - `/status` — 当前项目、模型、运行状态 - `/projects` — 列出最近项目并切换 - `/new` — 开一条新会话,可带项目编号或路径 +- `/sessions [项目]` — 查看旧会话;`/sessions projects` 选择其他项目 +- `/resume <编号>` — 继续列表中的旧会话 +- `/cancel` — 退出选择,保留当前会话 - `/clear` — 清空上下文,保留项目绑定 - `/stop` — 停止本轮生成