feat(desktop): add safe data directory migration (#1433)

This commit is contained in:
程序员阿江-Relakkes
2026-10-04 00:39:20 +08:00
committed by GitHub
parent c144148d20
commit 795ff9d4df
101 changed files with 5024 additions and 226 deletions
@@ -0,0 +1,41 @@
import { expect, test } from 'bun:test'
import { PassThrough } from 'node:stream'
import { installAdapterMigrationControl } from '../migration-control.js'
test('authenticated inherited pipe drains writes before acknowledgement and clean exit', async () => {
const input = new PassThrough()
const output = new PassThrough()
let release!: () => void
let didExit!: (code: number) => void
const exited = new Promise<number>(resolve => { didExit = resolve })
const write = new Promise<void>(resolve => { release = resolve })
let messages = ''
let drains = 0
output.on('data', chunk => { messages += chunk })
installAdapterMigrationControl({ token: 'isolated-secret', input, output, quiesce: async () => { drains++; await write }, exit: didExit })
input.write(JSON.stringify({ type: 'migration_quiesce', requestId: 'wrong', token: 'wrong' }) + '\n')
expect(drains).toBe(0)
input.write(JSON.stringify({ type: 'migration_quiesce', requestId: 'right', token: 'isolated-secret' }) + '\n')
await Promise.resolve()
expect(messages).toBe('')
release()
expect(await exited).toBe(0)
expect(JSON.parse(messages)).toEqual({ type: 'migration_quiesced', requestId: 'right' })
input.destroy()
output.destroy()
})
test('failed credential drain cannot emit a successful migration acknowledgement', async () => {
const input = new PassThrough()
const output = new PassThrough()
let didExit!: (code: number) => void
const exited = new Promise<number>(resolve => { didExit = resolve })
let messages = ''
output.on('data', chunk => { messages += chunk })
installAdapterMigrationControl({ token: 'isolated-secret', input, output, quiesce: async () => { throw new Error('Disk full') }, exit: didExit })
input.write(JSON.stringify({ type: 'migration_quiesce', requestId: 'failure', token: 'isolated-secret' }) + '\n')
expect(await exited).toBe(1)
expect(JSON.parse(messages).type).toBe('migration_quiesce_failed')
input.destroy()
output.destroy()
})
@@ -0,0 +1,50 @@
import { describe, expect, it } from 'bun:test'
import { AdapterMigrationLifecycle } from '../migration-lifecycle.js'
describe('adapter migration lifecycle', () => {
it('closes ingress before teardown and waits for in-flight credential writes', async () => {
const lifecycle = new AdapterMigrationLifecycle()
let release!: () => void
const calls: string[] = []
lifecycle.track(new Promise<void>(resolve => { release = resolve }))
lifecycle.registerShutdown(() => { expect(lifecycle.isQuiescing).toBe(true); calls.push('transport') })
let completed = false
const drain = lifecycle.quiesce().then(() => { completed = true })
await Promise.resolve()
expect(completed).toBe(false)
release()
await drain
expect(calls).toEqual(['transport'])
})
it('does not acknowledge a failed pending write', async () => {
const lifecycle = new AdapterMigrationLifecycle()
let fail!: (error: Error) => void
const pending = lifecycle.track(new Promise<void>((_resolve, reject) => { fail = reject }))
void pending.catch(() => {})
const drain = lifecycle.quiesce()
fail(new Error('Cannot save credentials'))
await expect(drain).rejects.toThrow('Cannot save credentials')
})
it('drains writes even when transport cleanup fails and closes each transport once', async () => {
const lifecycle = new AdapterMigrationLifecycle()
let release!: () => void
let closed = 0
lifecycle.track(new Promise<void>(resolve => { release = resolve }))
lifecycle.registerShutdown(() => {
closed++
throw new Error('Transport close failed')
})
const stopping = lifecycle.quiesce()
expect(lifecycle.quiesce()).toBe(stopping)
let settled = false
void stopping.catch(() => { settled = true })
await Promise.resolve()
await Promise.resolve()
expect(settled).toBe(false)
release()
await expect(stopping).rejects.toThrow('Transport close failed')
expect(closed).toBe(1)
})
})
@@ -16,6 +16,21 @@ afterEach(async () => {
})
describe('AttachmentStore', () => {
it('keeps default IM downloads inside the active storage directory', async () => {
const previous = process.env.CLAUDE_CONFIG_DIR
process.env.CLAUDE_CONFIG_DIR = tmpRoot
try {
const store = new AttachmentStore()
const target = store.resolvePath('feishu', 'session', 'fixture.png')
expect(target).toBe(path.join(tmpRoot, 'im-downloads', 'feishu', 'session', 'fixture.png'))
await store.write(target, Buffer.from('fixture'))
expect(await fs.readFile(target, 'utf8')).toBe('fixture')
} finally {
if (previous === undefined) delete process.env.CLAUDE_CONFIG_DIR
else process.env.CLAUDE_CONFIG_DIR = previous
}
})
it('writes a buffer and returns the absolute path', async () => {
const store = new AttachmentStore({ root: tmpRoot, retentionMs: 60_000 })
const target = store.resolvePath('feishu', 'sess-1', 'hello.png')
@@ -24,6 +24,14 @@ describe('ImageBlockWatcher', () => {
}
})
it('recognizes Windows drive and UNC paths after a data directory move', () => {
for (const sourcePath of [String.raw`D:\data\attachment.png`, 'D:/data/attachment.png', String.raw`\\server\share\attachment.png`]) {
const watcher = new ImageBlockWatcher()
expect(watcher.feed(`![moved](${sourcePath})`)[0]?.source).toEqual({ kind: 'path', path: sourcePath })
}
expect(new ImageBlockWatcher().feed('![relative](D:attachment.png)')).toEqual([])
})
it('extracts a markdown image with file:// URL as path', () => {
const w = new ImageBlockWatcher()
const out = w.feed('![x](file:///var/img/x.png)')
+12 -3
View File
@@ -11,6 +11,7 @@
*/
import * as fs from 'node:fs/promises'
import { adapterMigrationLifecycle } from '../migration-lifecycle.js'
import * as fsSync from 'node:fs'
import type { Dirent } from 'node:fs'
import * as path from 'node:path'
@@ -29,7 +30,7 @@ const DEFAULT_RETENTION_MS = 24 * 60 * 60 * 1000
const DEFAULT_ORPHAN_GRACE_MS = 10 * 60 * 1000
function defaultRoot(): string {
return path.join(os.homedir(), '.claude', 'im-downloads')
return path.join(process.env.CLAUDE_CONFIG_DIR || path.join(os.homedir(), '.claude'), 'im-downloads')
}
/** Strip path separators / .. / control chars from a filename. */
@@ -70,7 +71,11 @@ export class AttachmentStore {
}
/** Write atomically: stream to {target}.part, then rename. */
async write(target: string, data: Buffer): Promise<string> {
write(target: string, data: Buffer): Promise<string> {
return adapterMigrationLifecycle.track(this.writeOnce(target, data))
}
private async writeOnce(target: string, data: Buffer): Promise<string> {
await fs.mkdir(path.dirname(target), { recursive: true })
const tmp = `${target}.${process.pid}.${Date.now()}.part`
await fs.writeFile(tmp, data)
@@ -79,7 +84,11 @@ export class AttachmentStore {
}
/** Remove files older than retentionMs. Returns summary. */
async gc(): Promise<{ removed: number; bytes: number }> {
gc(): Promise<{ removed: number; bytes: number }> {
return adapterMigrationLifecycle.track(this.gcOnce())
}
private async gcOnce(): Promise<{ removed: number; bytes: number }> {
let removed = 0
let bytes = 0
const now = Date.now()
@@ -35,7 +35,7 @@ function classify(target: string): PendingUpload['source'] | null {
if (target.startsWith('http://') || target.startsWith('https://')) {
return { kind: 'url', url: target }
}
if (target.startsWith('/')) {
if (target.startsWith('/') || /^[a-zA-Z]:[\\/]/.test(target) || /^\\\\[^\\]+\\[^\\]+/.test(target)) {
return { kind: 'path', path: target }
}
return null // relative paths — skip, we can't resolve them safely
+5 -1
View File
@@ -6,11 +6,15 @@
* 参考 openclaw-lark chat-queue.ts 的 Promise 链设计。
*/
import { adapterMigrationLifecycle } from './migration-lifecycle.js'
const queues = new Map<string, Promise<void>>()
export async function enqueue(chatId: string, fn: () => Promise<void>): Promise<void> {
if (adapterMigrationLifecycle.isQuiescing) return
const prev = queues.get(chatId) ?? Promise.resolve()
const next = prev.then(fn, () => fn()).catch((err) => {
const invoke = () => adapterMigrationLifecycle.isQuiescing ? undefined : fn()
const next = adapterMigrationLifecycle.track(prev.then(invoke, invoke)).catch((err) => {
console.error(`[ChatQueue] Error in task for chat ${chatId}:`, err)
})
queues.set(chatId, next)
+2 -1
View File
@@ -18,6 +18,7 @@
import * as path from 'node:path'
import { enqueue } from './chat-queue.js'
import { adapterMigrationLifecycle } from './migration-lifecycle.js'
import {
getConfiguredWorkDir,
type AdapterConfig,
@@ -388,7 +389,7 @@ export class ImChatRuntime {
this.getBuffer(chatId).append(msg.text)
if (this.port.sendImage) {
for (const pending of this.getImageWatcher(chatId).feed(msg.text)) {
void this.dispatchOutboundImage(chatId, pending)
void adapterMigrationLifecycle.track(this.dispatchOutboundImage(chatId, pending))
}
}
}
+41
View File
@@ -0,0 +1,41 @@
import { timingSafeEqual } from 'node:crypto'
import type { Readable, Writable } from 'node:stream'
import { adapterMigrationLifecycle } from './migration-lifecycle.js'
/** An inherited anonymous pipe is available only to the owning desktop host. */
export function installAdapterMigrationControl(options: {
token: string
input?: Readable
output?: Writable
quiesce?: () => Promise<void>
exit?: (code: number) => void
}): void {
const input = options.input ?? process.stdin
const output = options.output ?? process.stdout
const exit = options.exit ?? (code => process.exit(code))
let buffer = ''
let stopping = false
input.on('data', chunk => {
buffer += chunk.toString()
if (buffer.length > 8192) {
buffer = ''
return
}
const lines = buffer.split('\n')
buffer = lines.pop() ?? ''
for (const line of lines) {
let request: { type?: string; token?: string; requestId?: string }
try { request = JSON.parse(line) } catch { continue }
if (stopping || request.type !== 'migration_quiesce' || typeof request.token !== 'string' || typeof request.requestId !== 'string') continue
const actual = Buffer.from(request.token)
const expected = Buffer.from(options.token)
if (actual.length !== expected.length || !timingSafeEqual(actual, expected)) continue
stopping = true
void (options.quiesce ?? (() => adapterMigrationLifecycle.quiesce()))().then(() => {
output.write(JSON.stringify({ type: 'migration_quiesced', requestId: request.requestId }) + '\n', () => exit(0))
}, () => {
output.write(JSON.stringify({ type: 'migration_quiesce_failed', requestId: request.requestId }) + '\n', () => exit(1))
})
}
})
}
+51
View File
@@ -0,0 +1,51 @@
/** Shared maintenance boundary for sidecar-owned IM transports and disk writes. */
export class AdapterMigrationLifecycle {
private quiescing = false
private pending = new Set<Promise<unknown>>()
private shutdown = new Set<() => Promise<void> | void>()
private failures: unknown[] = []
private stopping: Promise<void> | null = null
get isQuiescing(): boolean { return this.quiescing }
registerShutdown(cleanup: () => Promise<void> | void): void { this.shutdown.add(cleanup) }
track<T>(operation: Promise<T>): Promise<T> {
this.pending.add(operation)
void operation.then(() => this.pending.delete(operation), error => {
this.pending.delete(operation)
if (this.quiescing) this.failures.push(error)
})
return operation
}
quiesce(): Promise<void> {
if (this.stopping) return this.stopping
this.quiescing = true
this.stopping = this.drain()
return this.stopping
}
private async drain(): Promise<void> {
const cleanupResults = await Promise.allSettled([...this.shutdown].map(cleanup => Promise.resolve().then(cleanup)))
for (const result of cleanupResults) {
if (result.status === 'rejected') this.failures.push(result.reason)
}
while (this.pending.size > 0) await Promise.allSettled([...this.pending])
if (this.failures.length > 0) throw this.failures[0]
}
}
export const adapterMigrationLifecycle = new AdapterMigrationLifecycle()
export function registerAdapterShutdown(cleanup: () => Promise<void> | void): void {
adapterMigrationLifecycle.registerShutdown(cleanup)
const stop = () => {
void adapterMigrationLifecycle.quiesce().then(() => process.exit(0), error => {
console.error('[Adapter] Shutdown failed', error)
process.exit(1)
})
}
process.once('SIGINT', stop)
process.once('SIGTERM', stop)
}
+3 -2
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle } from './migration-lifecycle.js'
/**
* WebSocket Bridge
*
@@ -210,7 +211,7 @@ export class WsBridge {
// races where a later message reads stale map entries set up by an
// earlier-but-still-in-flight handler.
const prev = this.handlerChains.get(chatId) ?? Promise.resolve()
const next = prev
const next = adapterMigrationLifecycle.track(prev
.catch(() => {}) // upstream errors must not poison the chain
.then(() => {
// Resetting a chat cannot cancel promises already queued for its old
@@ -222,7 +223,7 @@ export class WsBridge {
})
.catch((err) => {
console.error(`[WsBridge] Handler error on ${chatId}:`, err)
})
}))
this.handlerChains.set(chatId, next)
})
+11 -16
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle, registerAdapterShutdown } from '../common/migration-lifecycle.js'
/**
* DingTalk Adapter for Claude Code Desktop.
*
@@ -674,6 +675,12 @@ async function start(): Promise<void> {
keepAlive: true,
} as any)
registerAdapterShutdown(async () => {
bridge.destroy()
dedup.destroy()
await client.disconnect()
})
client.registerCallbackListener(TOPIC_ROBOT, async (res: any) => {
const messageId = res.headers?.messageId
if (messageId) {
@@ -685,7 +692,7 @@ async function start(): Promise<void> {
if (!data) return
if (data.msgId && !dedup.tryRecord(`body:${data.msgId}`)) return
await handleRobotMessage(data)
await adapterMigrationLifecycle.track(handleRobotMessage(data))
})
client.registerCallbackListener(TOPIC_CARD, async (res: any) => {
@@ -695,28 +702,16 @@ async function start(): Promise<void> {
if (!dedup.tryRecord(`card:${messageId}`)) return
}
await handleCardCallback(res.data ?? res)
await adapterMigrationLifecycle.track(handleCardCallback(res.data ?? res))
})
if (adapterMigrationLifecycle.isQuiescing) return
await client.connect()
console.log(`[DingTalk] Stream connected. Server: ${config.serverUrl}`)
const shutdown = async () => {
console.log('[DingTalk] Shutting down...')
bridge.destroy()
dedup.destroy()
try {
await client.disconnect()
} catch {
// ignore
}
process.exit(0)
}
process.once('SIGINT', () => void shutdown())
process.once('SIGTERM', () => void shutdown())
}
if (import.meta.main || process.argv.includes('--dingtalk')) start().catch((err) => {
if (import.meta.main || process.argv.includes('--dingtalk')) adapterMigrationLifecycle.track(start()).catch((err) => {
console.error('[DingTalk] Fatal:', err instanceof Error ? err.message : err)
process.exit(1)
})
+7 -5
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle, registerAdapterShutdown } from '../common/migration-lifecycle.js'
/**
* 飞书 (Feishu/Lark) Adapter for Claude Code Desktop
*
@@ -1254,6 +1255,7 @@ async function start(): Promise<void> {
console.log(`[Feishu] App ID: ${config.feishu.appId}`)
await resolveBotOpenId()
if (adapterMigrationLifecycle.isQuiescing) return
const dispatcher = new Lark.EventDispatcher({
encryptKey: config.feishu.encryptKey,
@@ -1263,14 +1265,14 @@ async function start(): Promise<void> {
dispatcher.register({
'im.message.receive_v1': async (data: any) => {
try {
await handleMessage(data)
await adapterMigrationLifecycle.track(handleMessage(data))
} catch (err) {
console.error('[Feishu] Message handler error:', err)
}
},
'card.action.trigger': async (data: any) => {
try {
return await handleCardAction(data)
return await adapterMigrationLifecycle.track(handleCardAction(data))
} catch (err) {
console.error('[Feishu] Card action error:', err)
}
@@ -1288,16 +1290,16 @@ async function start(): Promise<void> {
console.log('[Feishu] Bot is running! (WebSocket connected)')
}
if (import.meta.main || process.argv.includes('--feishu')) start().catch((err) => {
if (import.meta.main || process.argv.includes('--feishu')) adapterMigrationLifecycle.track(start()).catch((err) => {
console.error('[Feishu] Failed to start:', err)
process.exit(1)
})
if (import.meta.main || process.argv.includes('--feishu')) process.on('SIGINT', () => {
if (import.meta.main || process.argv.includes('--feishu')) registerAdapterShutdown(async () => {
console.log('[Feishu] Shutting down...')
bridge.destroy()
dedup.destroy()
process.exit(0)
wsClient?.close({ force: true })
})
export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, handleMessage, handleCardAction, larkClient, prepareNewSession }
+2 -1
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle } from '../common/migration-lifecycle.js'
/**
* Feishu media service — wraps im.messageResource / im.image / im.file
* so adapters/feishu/index.ts stays focused on flow control.
@@ -101,7 +102,7 @@ export class FeishuMediaService {
})
if (typeof resp?.writeFile === 'function') {
await resp.writeFile(target)
await adapterMigrationLifecycle.track(Promise.resolve(resp.writeFile(target)))
} else if (resp?.data instanceof Buffer) {
await this.store.write(target, resp.data)
} else if (resp instanceof Buffer) {
+4 -4
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle, registerAdapterShutdown } from '../common/migration-lifecycle.js'
/**
* QQ Adapter for Claude Code Desktop
*
@@ -259,19 +260,18 @@ bot.on('error', (err: Error) => console.error('[QQ] Connection error:', err.mess
console.log('[QQ] Starting adapter...')
console.log(`[QQ] Server: ${config.serverUrl}`)
console.log(`[QQ] App: ${config.qq.appId}`)
void bot.start().catch((err) => {
void adapterMigrationLifecycle.track(bot.start()).catch((err) => {
console.error('[QQ] Failed to start:', err instanceof Error ? err.message : err)
process.exit(1)
})
process.on('SIGINT', () => {
registerAdapterShutdown(async () => {
console.log('[QQ] Shutting down...')
try {
bot.stop()
await bot.stop()
} catch {
// Best-effort: the process is exiting either way.
}
bridge.destroy()
dedup.destroy()
process.exit(0)
})
+9 -8
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle, registerAdapterShutdown } from '../common/migration-lifecycle.js'
/**
* Slack Adapter for Claude Code Desktop
*
@@ -250,7 +251,8 @@ const socket = new SlackSocketMode({
replyThreads.set(payload.chatId, payload.threadTs)
void (async () => {
if (adapterMigrationLifecycle.isQuiescing) return
void adapterMigrationLifecycle.track((async () => {
try {
await runtime.handleInbound({
chatId: payload.chatId,
@@ -264,14 +266,14 @@ const socket = new SlackSocketMode({
} catch (err) {
console.error('[Slack] Failed to prepare inbound message:', err)
}
})()
})())
},
})
console.log('[Slack] Starting adapter...')
console.log(`[Slack] Server: ${config.serverUrl}`)
void (async () => {
void adapterMigrationLifecycle.track((async () => {
try {
const identity = await api.authTest()
botUserId = identity.userId || undefined
@@ -280,13 +282,12 @@ void (async () => {
console.error('[Slack] auth.test failed:', err instanceof Error ? err.message : err)
process.exit(1)
}
await socket.start()
})()
if (!adapterMigrationLifecycle.isQuiescing) await socket.start()
})())
process.on('SIGINT', () => {
registerAdapterShutdown(async () => {
console.log('[Slack] Shutting down...')
socket.stop()
await socket.stop()
bridge.destroy()
dedup.destroy()
process.exit(0)
})
@@ -6,6 +6,7 @@ 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 { adapterMigrationLifecycle } from '../../common/migration-lifecycle.js'
// Import the actual entrypoint with isolated configuration. Telegram API calls
// terminate in grammY's documented transformer; HTTP and WS use loopback only.
@@ -357,7 +358,7 @@ describe('Telegram entrypoint session routing', () => {
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')
const previousListeners = new Map(['SIGINT', 'SIGTERM'].map(signal => [signal, process.listeners(signal)]))
try {
entry.startTelegramAdapter()
await eventually(() => expect(apiCalls.some((call) => call.method === 'setMyCommands')).toBe(true))
@@ -366,11 +367,54 @@ describe('Telegram entrypoint session routing', () => {
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)
for (const [signal, listeners] of previousListeners) {
for (const listener of process.listeners(signal)) {
if (!listeners.includes(listener)) process.removeListener(signal, listener)
}
}
start.mockRestore()
gc.mockRestore()
}
})
it('blocks real update ingress during migration and its registered shutdown waits for polling to stop', async () => {
let cleanup!: () => Promise<void> | void
let release!: () => void
const pendingStop = new Promise<void>(resolve => { release = resolve })
const register = spyOn(adapterMigrationLifecycle, 'registerShutdown').mockImplementation(callback => { cleanup = callback })
const start = spyOn(entry.bot, 'start').mockResolvedValue()
const running = spyOn(entry.bot, 'isRunning').mockReturnValue(true)
const stop = spyOn(entry.bot, 'stop').mockImplementation(async () => { await pendingStop })
const gc = spyOn(AttachmentStore.prototype, 'gc').mockResolvedValue({ removed: 0, bytes: 0 })
const destroy = spyOn(WsBridge.prototype, 'destroy')
const previousListeners = new Map(['SIGINT', 'SIGTERM'].map(signal => [signal, process.listeners(signal)]))
try {
entry.startTelegramAdapter()
expect(register).toHaveBeenCalledTimes(1)
Object.defineProperty(adapterMigrationLifecycle, 'isQuiescing', { configurable: true, get: () => true })
const previousRequests = requests.length
const previousMessages = messages.length
await text(710, 'Should remain unadmitted')
expect(requests).toHaveLength(previousRequests)
expect(messages).toHaveLength(previousMessages)
let completed = false
const stopping = Promise.resolve(cleanup()).then(() => { completed = true })
await Promise.resolve()
expect(stop).toHaveBeenCalledTimes(1)
expect(completed).toBe(false)
expect(destroy).not.toHaveBeenCalled()
release()
await stopping
expect(destroy).toHaveBeenCalledTimes(1)
} finally {
release()
Reflect.deleteProperty(adapterMigrationLifecycle, 'isQuiescing')
for (const [signal, listeners] of previousListeners) {
for (const listener of process.listeners(signal)) {
if (!listeners.includes(listener)) process.removeListener(signal, listener)
}
}
for (const spy of [destroy, gc, stop, running, start, register]) spy.mockRestore()
}
})
})
+10 -7
View File
@@ -1,3 +1,4 @@
import { adapterMigrationLifecycle, registerAdapterShutdown } from '../common/migration-lifecycle.js'
/**
* Telegram Adapter for Claude Code Desktop
*
@@ -51,6 +52,7 @@ if (!config.telegram.botToken) {
}
export const bot = new Bot(config.telegram.botToken)
bot.use((_ctx, next) => adapterMigrationLifecycle.isQuiescing ? Promise.resolve() : adapterMigrationLifecycle.track(next()))
const bridge = new WsBridge(config.serverUrl, 'tg')
const streamDelivery = new TelegramStreamDelivery(bot.api)
const dedup = new MessageDedup()
@@ -491,7 +493,7 @@ const isAuthorizedTelegramUser = (userId: number) => isAllowedUser('telegram', u
registerAuthorizedTelegramCommand(bot, 'stop', isAuthorizedTelegramUser, (ctx) => {
const chatId = String(ctx.chat!.id)
void (async () => {
void adapterMigrationLifecycle.track((async () => {
const result = await ensureExistingSession(chatId)
if (result.status !== 'restored') {
await ctx.reply(result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null))
@@ -499,7 +501,7 @@ registerAuthorizedTelegramCommand(bot, 'stop', isAuthorizedTelegramUser, (ctx) =
}
bridge.sendStopGeneration(chatId)
await ctx.reply('⏹ 已发送停止信号。')
})()
})())
})
registerAuthorizedTelegramCommand(bot, 'status', isAuthorizedTelegramUser, async (ctx) => {
@@ -509,7 +511,7 @@ registerAuthorizedTelegramCommand(bot, 'status', isAuthorizedTelegramUser, async
registerAuthorizedTelegramCommand(bot, 'clear', isAuthorizedTelegramUser, (ctx) => {
const chatId = String(ctx.chat!.id)
void (async () => {
void adapterMigrationLifecycle.track((async () => {
const result = await ensureExistingSession(chatId)
if (result.status !== 'restored') {
await ctx.reply(result.status === 'unavailable' ? SESSION_RECONNECT_NOTICE : formatImStatus(null))
@@ -523,7 +525,7 @@ registerAuthorizedTelegramCommand(bot, 'clear', isAuthorizedTelegramUser, (ctx)
}
getRuntimeState(chatId).state = 'thinking'
await ctx.reply('🧹 已清空当前会话上下文。')
})()
})())
})
for (const command of ['allow', 'always', 'allow-always', 'deny'] as const) {
@@ -714,10 +716,11 @@ export function startTelegramAdapter(): void {
})
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', () => {
registerAdapterShutdown(async () => {
console.log('[Telegram] Shutting down...')
stopTelegramAdapter()
process.exit(0)
if (bot.isRunning()) await bot.stop()
bridge.destroy()
dedup.destroy()
})
}
+3 -2
View File
@@ -1,3 +1,4 @@
import { registerAdapterShutdown } from '../common/migration-lifecycle.js'
import * as path from 'node:path'
import { WsBridge, type ServerMessage, type AttachmentRef } from '../common/ws-bridge.js'
import { MessageDedup } from '../common/message-dedup.js'
@@ -616,6 +617,7 @@ async function pollLoop(): Promise<void> {
timeoutMs: GET_UPDATES_TIMEOUT_MS,
})
if (resp.get_updates_buf) getUpdatesBuf = resp.get_updates_buf
if (stopped) return
const hasRetError = typeof resp.ret === 'number' && resp.ret !== 0
const hasErrCode = typeof resp.errcode === 'number' && resp.errcode !== 0
if (hasRetError || hasErrCode) {
@@ -646,13 +648,12 @@ console.log('[WeChat] Starting adapter...')
console.log(`[WeChat] Account: ${accountId}`)
if (import.meta.main || process.argv.includes('--wechat')) void pollLoop()
if (import.meta.main || process.argv.includes('--wechat')) process.on('SIGINT', () => {
if (import.meta.main || process.argv.includes('--wechat')) registerAdapterShutdown(() => {
console.log('[WeChat] Shutting down...')
stopped = true
typingController.destroy()
bridge.destroy()
dedup.destroy()
process.exit(0)
})
export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, routeUserMessage, startNewSession, typingController }
+2 -2
View File
@@ -1,3 +1,4 @@
import { registerAdapterShutdown } from '../common/migration-lifecycle.js'
/**
* 企业微信 (Enterprise WeChat / WeCom) Adapter for Claude Code Desktop
*
@@ -255,7 +256,7 @@ console.log(`[WeCom] Server: ${config.serverUrl}`)
console.log(`[WeCom] Bot: ${config.wecom.botId}`)
client.connect()
process.on('SIGINT', () => {
registerAdapterShutdown(() => {
console.log('[WeCom] Shutting down...')
try {
client.disconnect()
@@ -264,5 +265,4 @@ process.on('SIGINT', () => {
}
bridge.destroy()
dedup.destroy()
process.exit(0)
})
+102 -2
View File
@@ -1,10 +1,14 @@
import { describe, expect, it } from 'bun:test'
import { describe, expect, it, spyOn } from 'bun:test'
import { EventEmitter } from 'node:events'
import * as baileys from '@whiskeysockets/baileys'
import { AdapterMigrationLifecycle, adapterMigrationLifecycle } from '../../common/migration-lifecycle.js'
import * as fs from 'node:fs'
import * as os from 'node:os'
import * as path from 'node:path'
import {
clearWhatsAppAuth,
closeWhatsAppSocket,
createWhatsAppSocket,
getWhatsAppDisconnectStatus,
hasWhatsAppAuth,
isWhatsAppLoggedOut,
@@ -58,6 +62,102 @@ describe('whatsapp session helpers', () => {
})
it('returns immediately when no credential save is queued', async () => {
await expect(waitForWhatsAppCredsSave(makeTempAuthDir())).resolves.toBeUndefined()
const authDir = makeTempAuthDir()
try { await expect(waitForWhatsAppCredsSave(authDir)).resolves.toBeUndefined() } finally { fs.rmSync(authDir, { recursive: true, force: true }) }
})
it('drains queued credentials and signal-key writes from the actual socket binding before migration', async () => {
const authDir = makeTempAuthDir()
const lifecycle = new AdapterMigrationLifecycle()
const events = new EventEmitter()
let releaseCreds!: () => void
let releaseKeys!: () => void
const pendingCreds = new Promise<void>(resolve => { releaseCreds = resolve })
const pendingKeys = new Promise<void>(resolve => { releaseKeys = resolve })
let saves = 0
let authKeys: any
fs.writeFileSync(path.join(authDir, 'creds.json'), '{"old":true}')
const auth = spyOn(baileys, 'useMultiFileAuthState').mockResolvedValue({
state: { creds: {} as any, keys: { get: async () => ({}), set: async () => { await pendingKeys } } },
saveCreds: async () => {
if (++saves === 1) await pendingCreds
fs.writeFileSync(path.join(authDir, 'creds.json'), JSON.stringify({ saves }))
},
})
const version = spyOn(baileys, 'fetchLatestBaileysVersion').mockResolvedValue({ version: [2, 3, 4], isLatest: true })
const socket = spyOn(baileys, 'makeWASocket').mockImplementation((options: any) => {
authKeys = options.auth.keys
return { ev: events, ws: new EventEmitter() } as any
})
const track = spyOn(adapterMigrationLifecycle, 'track').mockImplementation(operation => lifecycle.track(operation))
try {
await createWhatsAppSocket({ authDir })
events.emit('creds.update', {})
events.emit('creds.update', {})
const keys = authKeys.set({ 'pre-key': { 'fixture-key': { data: 'fixture' } } })
lifecycle.registerShutdown(() => waitForWhatsAppCredsSave(authDir))
let drained = false
const stopping = lifecycle.quiesce().then(() => { drained = true })
await Promise.resolve()
expect(drained).toBe(false)
releaseCreds()
await waitForWhatsAppCredsSave(authDir)
expect(saves).toBe(2)
expect(drained).toBe(false)
releaseKeys()
await Promise.all([keys, stopping])
expect(JSON.parse(fs.readFileSync(path.join(authDir, 'creds.json'), 'utf8'))).toEqual({ saves: 2 })
expect(fs.existsSync(path.join(authDir, 'creds.json.bak'))).toBe(true)
} finally {
releaseCreds()
releaseKeys()
for (const spy of [track, socket, version, auth]) spy.mockRestore()
await waitForWhatsAppCredsSave(authDir)
fs.rmSync(authDir, { recursive: true, force: true })
}
})
it('keeps the migration barrier pending for a credential save that has not started yet', async () => {
const authDir = makeTempAuthDir()
const lifecycle = new AdapterMigrationLifecycle()
const events = new EventEmitter()
let first!: () => void
let second!: () => void
const pending = [new Promise<void>(resolve => { first = resolve }), new Promise<void>(resolve => { second = resolve })]
let saves = 0
const auth = spyOn(baileys, 'useMultiFileAuthState').mockResolvedValue({
state: { creds: {} as any, keys: { get: async () => ({}), set: async () => {} } },
saveCreds: async () => {
await pending[saves++]
fs.writeFileSync(path.join(authDir, 'creds.json'), JSON.stringify({ saves }))
},
})
const version = spyOn(baileys, 'fetchLatestBaileysVersion').mockResolvedValue({ version: [2, 3, 4], isLatest: true })
const socket = spyOn(baileys, 'makeWASocket').mockReturnValue({ ev: events, ws: new EventEmitter() } as any)
const track = spyOn(adapterMigrationLifecycle, 'track').mockImplementation(operation => lifecycle.track(operation))
let stopping: Promise<void> | undefined
try {
await createWhatsAppSocket({ authDir })
events.emit('creds.update', {})
events.emit('creds.update', {})
await Promise.resolve()
expect(saves).toBe(1)
let drained = false
stopping = lifecycle.quiesce().then(() => { drained = true })
first()
for (let tick = 0; tick < 15; tick++) await Promise.resolve()
expect(saves).toBe(2)
expect(drained).toBe(false)
second()
await stopping
expect(drained).toBe(true)
} finally {
first()
second()
await stopping
await waitForWhatsAppCredsSave(authDir)
for (const spy of [track, socket, version, auth]) spy.mockRestore()
fs.rmSync(authDir, { recursive: true, force: true })
}
})
})
+13 -5
View File
@@ -1,3 +1,5 @@
import { adapterMigrationLifecycle, registerAdapterShutdown } from '../common/migration-lifecycle.js'
import { waitForWhatsAppCredsSave } from './session.js'
/**
* WhatsApp Adapter for Claude Code Desktop
*
@@ -605,16 +607,21 @@ export function useWhatsAppSocket(socket: WhatsAppSocket): void {
}
async function startSocket(): Promise<void> {
if (shuttingDown || adapterMigrationLifecycle.isQuiescing) return
if (reconnectTimer) {
clearTimeout(reconnectTimer)
reconnectTimer = null
}
useWhatsAppSocket(await createWhatsAppSocket({ authDir }))
if (shuttingDown || adapterMigrationLifecycle.isQuiescing) {
closeWhatsAppSocket(sock, 'data migration')
return
}
sock.ev.on('messages.upsert', ({ type, messages }) => {
if (type !== 'notify') return
if (type !== 'notify' || adapterMigrationLifecycle.isQuiescing) return
for (const message of messages) {
void handleIncomingMessage(message)
void adapterMigrationLifecycle.track(handleIncomingMessage(message))
}
})
@@ -635,13 +642,14 @@ async function startSocket(): Promise<void> {
}
function scheduleReconnect(): void {
if (shuttingDown || adapterMigrationLifecycle.isQuiescing) return
if (reconnectTimer) return
const delay = Math.min(RECONNECT_MAX_MS, RECONNECT_BASE_MS * 2 ** reconnectAttempts)
reconnectAttempts += 1
console.warn(`[WhatsApp] Connection closed. Reconnecting in ${delay}ms...`)
reconnectTimer = setTimeout(() => {
reconnectTimer = null
startSocket().catch((err) => {
adapterMigrationLifecycle.track(startSocket()).catch((err) => {
console.error('[WhatsApp] Reconnect failed:', err instanceof Error ? err.message : err)
scheduleReconnect()
})
@@ -655,14 +663,14 @@ console.log(`[WhatsApp] Allowed users: ${config.whatsapp.allowedUsers.length ===
if (import.meta.main || process.argv.includes('--whatsapp')) await startSocket()
if (import.meta.main || process.argv.includes('--whatsapp')) process.on('SIGINT', () => {
if (import.meta.main || process.argv.includes('--whatsapp')) registerAdapterShutdown(async () => {
console.log('[WhatsApp] Shutting down...')
shuttingDown = true
if (reconnectTimer) clearTimeout(reconnectTimer)
closeWhatsAppSocket(sock, 'SIGINT')
bridge.destroy()
dedup.destroy()
process.exit(0)
await waitForWhatsAppCredsSave(authDir)
})
export { bridge, dedup, sessionStore, sessionSelectionController, handleServerMessage, getRuntimeState, clearTransientChatState, createSessionForChat, showProjectPicker, routeUserMessage, startNewSession }
+5 -2
View File
@@ -1,4 +1,5 @@
import * as fs from 'node:fs'
import { adapterMigrationLifecycle } from '../common/migration-lifecycle.js'
import * as path from 'node:path'
import {
DisconnectReason,
@@ -48,6 +49,8 @@ export async function createWhatsAppSocket(options: {
const logger = makeBaileysLogger(options.verbose ? 'info' : 'silent')
const { state, saveCreds } = await useMultiFileAuthState(authDir)
const setKeys = state.keys.set.bind(state.keys)
state.keys.set = (...args) => adapterMigrationLifecycle.track(Promise.resolve(setKeys(...args)))
const { version } = await fetchLatestBaileysVersion()
const sock = makeWASocket({
auth: {
@@ -108,8 +111,8 @@ function maybeRestoreCredsFromBackup(authDir: string): void {
function enqueueSaveCreds(authDir: string, saveCreds: () => Promise<void> | void): void {
const resolved = path.resolve(authDir)
const prev = credsSaveQueues.get(resolved) ?? Promise.resolve()
const next = prev
.then(() => safeSaveCreds(resolved, saveCreds))
const save = adapterMigrationLifecycle.track(prev.then(() => safeSaveCreds(resolved, saveCreds)))
const next = save
.catch((err) => {
console.warn('[WhatsApp] Failed to save credentials:', err instanceof Error ? err.message : err)
})