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
This commit is contained in:
程序员阿江(Relakkes)
2026-09-08 13:56:16 +08:00
parent 82447fcf23
commit 59c7857beb
24 changed files with 2023 additions and 124 deletions
+99 -1
View File
@@ -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<string>()
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')
})
})
+48 -5
View File
@@ -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) {
@@ -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<string, Set<string>>()
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 })
})
})
@@ -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> = {}): 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([])
})
})
@@ -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<void>((resolve) => { release = resolve })
const firstStarted = new Promise<void>((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<void>((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<void>
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[] = []
+22
View File
@@ -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<string, Set<string>>()
private readonly pendingProjectSelection = new Set<string>()
private readonly imageWatchers = new Map<string, ImageBlockWatcher>()
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<void> {
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<void> {
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<void> {
this.sessionSelection.clear(chatId)
this.bridge.resetSession(chatId)
this.sessionStore.delete(chatId)
this.clearTransientChatState(chatId)
+6 -1
View File
@@ -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 = ['当前会话状态:']
+6 -6
View File
@@ -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)
}
+46
View File
@@ -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<string, Set<string>>,
): 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<string>(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
}
+262
View File
@@ -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<AdapterHttpClient, 'sessionExists'>
bridge: Pick<WsBridge, 'resetSession' | 'connectSession' | 'onServerMessage' | 'waitForOpen'>
sessionStore: Pick<SessionStore, 'get' | 'set' | 'delete'>
onServerMessage: (chatId: string, message: ServerMessage) => void | Promise<void>
clearTransientState: (chatId: string) => void
isBusy?: (chatId: string) => boolean
}
export type SessionSelectionDeps = SessionRestoreDeps & {
httpClient: Pick<AdapterHttpClient, 'sessionExists' | 'listRecentProjects' | 'matchProject' | 'listSessions'>
sendNotice: (chatId: string, text: string) => Promise<void>
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<SessionListItem, 'id' | 'workDir' | 'title'>,
): 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<AdapterHttpClient, 'listSessions'>,
project: string,
all?: SessionListItem[],
): Promise<SessionListItem[]> {
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<AdapterHttpClient, 'listSessions'>): Promise<SessionListItem[]> {
const sessions = new Map<string, SessionListItem>()
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<string, Picker>()
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<boolean> {
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<Picker | null> {
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<void> {
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<void> {
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<void> {
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<void> {
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<void> {
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'))
}
}
+9 -3
View File
@@ -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)
})
+32 -2
View File
@@ -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<boolean> {
}
async function createSessionForChat(chatId: string, workDir: string): Promise<boolean> {
sessionSelectionController.clear(chatId)
try {
bridge.resetSession(chatId)
clearTransientChatState(chatId)
@@ -331,6 +349,7 @@ function formatProjectList(projects: RecentProject[]): string {
}
async function showProjectPicker(chatId: string): Promise<void> {
sessionSelectionController.clear(chatId)
try {
const projects = await projectSelectionController.listProjects(chatId)
if (projects.length === 0) {
@@ -344,6 +363,7 @@ async function showProjectPicker(chatId: string): Promise<void> {
}
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<void> {
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<void> {
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 }
@@ -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<void> | 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<void>((resolve) => { preflightStarted = resolve })
preflightGate = new Promise<void>((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
}
})
+38 -8
View File
@@ -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<typeof Lark.WSClient> | 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<boolean> {
}
async function createSessionForChat(chatId: string, workDir: string): Promise<boolean> {
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<bo
}
async function showProjectPicker(chatId: string): Promise<void> {
sessionSelectionController.clear(chatId)
try {
const projects = await projectSelectionController.listProjects(chatId)
if (projects.length === 0) {
@@ -694,6 +713,7 @@ async function showProjectPicker(chatId: string): Promise<void> {
}
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<void> {
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<void> {
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<void> {
}
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<void> {
})
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<any> {
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<void> {
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 }
+133 -14
View File
@@ -118,6 +118,7 @@ function createController(overrides?: Record<string, unknown>) {
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<string, unknown>) {
],
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<string, unknown>) {
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, {
@@ -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<string, string | undefined>()
let directory: string
let project: string
let worktree: string
let entry: typeof import('../index.js')
let server: ReturnType<typeof Bun.serve<{ sessionId: string }>>
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<string, Set<ServerWebSocket<{ sessionId: string }>>>()
const sessionPaths = new Map<string, string>()
async function eventually(assertion: () => void): Promise<void> {
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<void> {
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<void> {
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()
}
})
})
@@ -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<string, (ctx: any) => 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)
})
})
+90 -33
View File
@@ -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>) => void
handleServerMessage: (chatId: string, msg: ServerMessage) => void | Promise<void>
waitForBridgeOpen: (chatId: string) => Promise<boolean>
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>) => void
waitForOpen: (chatId: string) => Promise<boolean>
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<void>
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>,
): 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<void>
showProjectPicker: (chatId: string) => Promise<void>
showResumeProjectPicker: (chatId: string) => Promise<void>
handleSessionInput: (chatId: string, text: string) => Promise<boolean>
},
): Promise<boolean> {
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<void> => {
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',
+86 -44
View File
@@ -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<string, string>()
// 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<bo
}
async function showProjectPicker(chatId: string): Promise<void> {
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<void> {
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<void> {
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()
+1
View File
@@ -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: '清空当前上下文' },
+30 -2
View File
@@ -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<boolean> {
}
async function createSessionForChat(chatId: string, workDir: string): Promise<boolean> {
sessionSelectionController.clear(chatId)
try {
bridge.resetSession(chatId)
clearTransientChatState(chatId)
@@ -297,6 +315,7 @@ async function createSessionForChat(chatId: string, workDir: string): Promise<bo
}
async function showProjectPicker(chatId: string): Promise<void> {
sessionSelectionController.clear(chatId)
try {
const projects = await httpClient.listRecentProjects()
if (projects.length === 0) {
@@ -315,6 +334,7 @@ async function showProjectPicker(chatId: string): Promise<void> {
}
async function startNewSession(chatId: string, query?: string): Promise<void> {
sessionSelectionController.clear(chatId)
bridge.resetSession(chatId)
sessionStore.delete(chatId)
clearTransientChatState(chatId)
@@ -351,10 +371,12 @@ async function startNewSession(chatId: string, query?: string): Promise<void> {
async function handleServerMessage(chatId: string, msg: ServerMessage): Promise<void> {
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<void> {
}
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<void> {
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<void> {
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 }
+37 -5
View File
@@ -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<string, ChatRuntimeState>()
const pendingPermissions = new Map<string, Set<string>>()
const imageWatchers = new Map<string, ImageBlockWatcher>()
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<boolean> {
}
async function createSessionForChat(chatId: string, workDir: string): Promise<boolean> {
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<bo
}
async function showProjectPicker(chatId: string): Promise<void> {
sessionSelectionController.clear(chatId)
try {
const projects = await httpClient.listRecentProjects()
if (projects.length === 0) {
@@ -270,11 +289,13 @@ async function flushAccumulatedText(chatId: string): Promise<void> {
async function handleServerMessage(chatId: string, msg: ServerMessage): Promise<void> {
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<void> {
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<vo
)
}
export function useWhatsAppSocket(socket: WhatsAppSocket): void {
sock = socket
media = new WhatsAppMediaService(sock, attachmentStore)
}
async function startSocket(): Promise<void> {
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 }