fix(adapters): preserve IM session binding on transient reconnect failures

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