diff --git a/README.md b/README.md index 58d3eca..0940794 100644 --- a/README.md +++ b/README.md @@ -41,9 +41,9 @@ curl -fsS http://localhost:8787/status -H "Authorization: Bearer YOUR_TOKEN" \ The renderer runs locally. The QR is short-lived and must be kept private. Repeat the command if it expires. The response contains the QR only while pairing. Send `ping` from another WhatsApp account to receive `pong`. Group messages and messages sent by the bot are ignored. -If `/status` reports `closed`, call `/start` again to reconnect. After an object restart, it reports `stopped`; call `/start` to reconnect using its stored credentials. +After an object restart, the first authenticated request automatically attempts to reconnect with registered credentials already in Durable Object storage. A fresh, unpaired object remains `stopped` until you call `/start` to begin QR pairing. After a non-logout disconnect or a failed connection attempt, `/status` and `/start` wait for a 30-second cooldown before retrying. The cooldown is held in memory and resets on an object restart. Concurrent requests share an in-progress connection attempt. Automatic recovery requires registered credentials; use `/start` if pairing is needed. -If `/status` reports `logged_out`, the bot has cleared its stored credentials. Call `/start` and pair again. +Recovery requires a request to trigger each reconnect attempt. This does not guarantee an unattended, uninterrupted WhatsApp connection. If `/status` reports `logged_out`, WhatsApp has revoked the session and the bot clears its stored credentials. Call `/start` and pair again. A failure to clear credentials changes the status to `closed` with an error. The Durable Object stores credentials and Signal state. Treat its storage as a secret. Do not share status responses, QR codes, or backups. diff --git a/src/index.ts b/src/index.ts index 3cabefd..f054b60 100644 --- a/src/index.ts +++ b/src/index.ts @@ -11,58 +11,108 @@ interface Env { ADMIN_TOKEN: string } +const RECONNECT_DELAY_MS = 30_000 + export class Bot extends DurableObject { private socket?: HostWASocket private status: BotStatus = { state: 'stopped' } + private starting?: Promise + private retryAt = 0 + private readonly initialization: Promise constructor(ctx: DurableObjectState, env: Env) { super(ctx, env) + this.initialization = this.ctx.blockConcurrencyWhile(async () => { + try { + await this.start(false) + } catch (error) { + console.error('Could not restore WhatsApp connection', JSON.stringify({ errorClass: error instanceof Error ? error.name : typeof error })) + this.status = { state: 'closed', error: 'Could not restore stored session' } + this.retryAt = Date.now() + RECONNECT_DELAY_MS + } + }) } async fetch(request: Request): Promise { + await this.initialization const url = new URL(request.url) if (request.method === 'POST' && url.pathname === '/start') { - await this.start() + await this.start(true) + return Response.json(this.status) + } + if (request.method === 'GET' && url.pathname === '/status') { + if (this.status.state === 'closed' && Date.now() >= this.retryAt) { + try { + await this.start(false) + } catch (error) { + console.error('Could not reconnect WhatsApp session', JSON.stringify({ errorClass: error instanceof Error ? error.name : typeof error })) + this.status = { state: 'closed', error: 'Could not reconnect stored session' } + this.retryAt = Date.now() + RECONNECT_DELAY_MS + } + } return Response.json(this.status) } - if (request.method === 'GET' && url.pathname === '/status') return Response.json(this.status) return new Response('Not found', { status: 404 }) } - private async start(): Promise { - if (this.socket) return + private start(allowPairing: boolean): Promise { + if (this.socket) return Promise.resolve() + if (this.starting) return this.starting + if (this.status.state === 'closed' && Date.now() < this.retryAt) return Promise.resolve() + + const starting = this.openSocket(allowPairing).catch(error => { + this.status = { state: 'closed', error: 'Could not start WhatsApp connection' } + this.retryAt = Date.now() + RECONNECT_DELAY_MS + throw error + }) + this.starting = starting + return starting.finally(() => { + if (this.starting === starting) this.starting = undefined + }) + } + + private async openSocket(allowPairing: boolean): Promise { initSync({ module: wasm }) const store = createStore(this.ctx.storage) const auth = await createAuthenticationState(store) + if (!allowPairing && !auth.creds.registered) { + if (this.status.state === 'closed') this.status = { state: 'stopped' } + return + } + // Retain the explicit protocol version used when WhatsApp rejected the former preview's default. const socket = makeWASocket({ auth, version: [2, 3000, 1043857760], logger: undefined }) this.socket = socket socket.ev.on('connection.update', async ({ connection, qr, lastDisconnect }) => { console.log('connection state', connection ?? (qr ? 'qr' : 'updated')) if (qr) this.status = { state: 'waiting_for_qr', qr } - else if (connection === 'open') this.status = { state: 'connected' } - else if (connection === 'connecting') this.status = { state: 'connecting' } + else if (connection === 'open') { + this.status = { state: 'connected' } + this.retryAt = 0 + } else if (connection === 'connecting') this.status = { state: 'connecting' } else if (connection === 'close') { const statusCode = (lastDisconnect?.error as (Error & { output?: { statusCode?: number } }) | undefined)?.output?.statusCode + this.socket = undefined if (statusCode === 401) { - this.status = { state: 'closed' } + this.status = { state: 'logged_out' } try { await this.ctx.storage.deleteAll() - this.status = { state: 'logged_out' } } catch (error) { - console.error('Could not clear logged-out credentials', error) + console.error('Could not clear logged-out credentials', JSON.stringify({ errorClass: error instanceof Error ? error.name : typeof error })) this.status = { state: 'closed', error: 'Could not clear logged-out credentials' } + this.retryAt = Date.now() + RECONNECT_DELAY_MS } } else { this.status = { state: 'closed', error: lastDisconnect?.error?.message } + this.retryAt = Date.now() + RECONNECT_DELAY_MS } - this.socket = undefined } }) socket.ev.on('messages.upsert', event => { void handleMessages(socket, event).catch(error => console.error('Message handler failed', JSON.stringify({ errorClass: error instanceof Error ? error.name : typeof error }))) }) this.status = { state: 'connecting' } + this.retryAt = 0 } } diff --git a/test/index.test.ts b/test/index.test.ts index 1591a1f..16b3db7 100644 --- a/test/index.test.ts +++ b/test/index.test.ts @@ -1,5 +1,5 @@ -import { beforeEach, expect, it, vi } from 'vitest' -import type { HostBaileysEventMap } from '@oxidezap/baileyrs/host' +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import type { HostBaileysEventMap, HostStoreCallbacks } from '@oxidezap/baileyrs/host' const host = vi.hoisted(() => ({ on: vi.fn(), @@ -21,60 +21,115 @@ vi.mock('@oxidezap/baileyrs/host', () => ({ import { Bot } from '../src/index.ts' +const request = (path: '/start' | '/status', method: 'GET' | 'POST' = 'GET') => + new Request(`https://bot${path}`, { method }) + +const createStorage = (values = new Map()) => ({ + values, + get: async (key: string) => values.get(key), + put: async (key: string, value: unknown) => { values.set(key, value) }, + deleteAll: vi.fn(async () => { values.clear() }) +}) + +const createContext = (storage: ReturnType) => ({ + storage, + blockConcurrencyWhile: async (callback: () => Promise) => callback() +}) as unknown as DurableObjectState + +const createEnv = () => ({ ADMIN_TOKEN: 'test', BOT: {} as DurableObjectNamespace }) + +const latestConnectionUpdate = () => host.on.mock.calls + .filter(([event]) => event === 'connection.update') + .at(-1)?.[1] as (event: HostBaileysEventMap['connection.update']) => Promise | void + +afterEach(() => vi.restoreAllMocks()) + beforeEach(() => { vi.resetAllMocks() host.makeSocket.mockReturnValue({ ev: { on: host.on } }) - host.authenticate.mockResolvedValue({ creds: { registered: true } }) + host.authenticate.mockResolvedValue({ creds: { registered: false } }) }) -it('keeps connection status with the socket across object restarts', async () => { - const values = new Map() - const storage = { - get: async (key: string) => values.get(key), - put: async (key: string, value: unknown) => { values.set(key, value) } - } - const ctx = { storage } as unknown as DurableObjectState - const env = { ADMIN_TOKEN: 'test', BOT: {} as DurableObjectNamespace } +it('restores persisted credentials on a new instance without requiring POST /start', async () => { + host.authenticate.mockImplementation(async (store: HostStoreCallbacks) => ({ + creds: { registered: Boolean(await store.get('auth', 'registered')) } + })) + const storage = createStorage() + const ctx = createContext(storage) + const env = createEnv() const first = new Bot(ctx, env) - await first.fetch(new Request('https://bot/start', { method: 'POST' })) - const update = host.on.mock.calls.find(([event]) => event === 'connection.update')?.[1] as - (event: HostBaileysEventMap['connection.update']) => Promise | void - await update({ connection: 'open' }) - expect(await (await first.fetch(new Request('https://bot/status'))).json()).toEqual({ state: 'connected' }) + + expect(await (await first.fetch(request('/status'))).json()).toEqual({ state: 'stopped' }) + expect(host.makeSocket).not.toHaveBeenCalled() + + await first.fetch(request('/start', 'POST')) + await latestConnectionUpdate()({ qr: 'pairing-qr' }) + expect(await (await first.fetch(request('/status'))).json()).toEqual({ state: 'waiting_for_qr', qr: 'pairing-qr' }) + await storage.put('auth:registered', new Uint8Array([1])) // pairing persisted credentials in the shared durable store + await latestConnectionUpdate()({ connection: 'open' }) + expect(await (await first.fetch(request('/status'))).json()).toEqual({ state: 'connected' }) + const restarted = new Bot(ctx, env) - expect(await (await restarted.fetch(new Request('https://bot/status'))).json()).toEqual({ state: 'stopped' }) + expect(await (await restarted.fetch(request('/status'))).json()).toEqual({ state: 'connecting' }) + expect(host.authenticate).toHaveBeenCalledTimes(3) + expect(host.makeSocket).toHaveBeenCalledTimes(2) + expect(await (await restarted.fetch(request('/status'))).json()).toEqual({ state: 'connecting' }) + expect(host.makeSocket).toHaveBeenCalledTimes(2) +}) + +it('retries a closed authenticated socket on status after a cooldown, without duplicate starts', async () => { + const now = vi.spyOn(Date, 'now').mockReturnValue(1_000_000) + host.authenticate.mockResolvedValue({ creds: { registered: true } }) + const bot = new Bot(createContext(createStorage()), createEnv()) + await bot.fetch(request('/status')) + expect(host.makeSocket).toHaveBeenCalledTimes(1) + + await latestConnectionUpdate()({ connection: 'close', lastDisconnect: { + date: new Date(), error: new Error('temporary disconnect') + } }) + await bot.fetch(request('/status')) + expect(host.makeSocket).toHaveBeenCalledTimes(1) + + now.mockReturnValue(1_031_000) + const [first, second] = await Promise.all([ + bot.fetch(request('/status')), + bot.fetch(request('/status')) + ]) + expect(await first.json()).toEqual({ state: 'connecting' }) + expect(await second.json()).toEqual({ state: 'connecting' }) + expect(host.makeSocket).toHaveBeenCalledTimes(2) +}) + +it('bounds repeated /start attempts when socket creation fails', async () => { + host.makeSocket.mockImplementationOnce(() => { throw new Error('temporary socket failure') }) + const bot = new Bot(createContext(createStorage()), createEnv()) + await expect(bot.fetch(request('/start', 'POST'))).rejects.toThrow('temporary socket failure') + expect(await (await bot.fetch(request('/start', 'POST'))).json()).toEqual({ + state: 'closed', error: 'Could not start WhatsApp connection' + }) + expect(host.makeSocket).toHaveBeenCalledOnce() }) it('clears revoked credentials before allowing a new pairing attempt', async () => { const values = new Map([['device:device', new Uint8Array([1])]]) - const storage = { - get: async (key: string) => values.get(key), - deleteAll: vi.fn(async () => { values.clear() }) - } - const bot = new Bot({ storage } as unknown as DurableObjectState, { - ADMIN_TOKEN: 'test', BOT: {} as DurableObjectNamespace - }) - await bot.fetch(new Request('https://bot/start', { method: 'POST' })) - const update = host.on.mock.calls.find(([event]) => event === 'connection.update')?.[1] as - (event: HostBaileysEventMap['connection.update']) => Promise - await update({ connection: 'close', lastDisconnect: { + host.authenticate.mockResolvedValue({ creds: { registered: true } }) + const storage = createStorage(values) + const bot = new Bot(createContext(storage), createEnv()) + await bot.fetch(request('/status')) + await latestConnectionUpdate()({ connection: 'close', lastDisconnect: { date: new Date(), error: Object.assign(new Error('logged out'), { output: { statusCode: 401 } }) } }) expect(storage.deleteAll).toHaveBeenCalledOnce() expect(values.size).toBe(0) - expect(await (await bot.fetch(new Request('https://bot/status'))).json()).toEqual({ state: 'logged_out' }) - await bot.fetch(new Request('https://bot/start', { method: 'POST' })) + expect(await (await bot.fetch(request('/status'))).json()).toEqual({ state: 'logged_out' }) + await bot.fetch(request('/start', 'POST')) expect(host.authenticate).toHaveBeenCalledTimes(2) + expect(host.makeSocket).toHaveBeenCalledTimes(2) }) it('passes native credentials to the socket without a stale snapshot', async () => { - const storage = { - get: async () => JSON.stringify({ registered: false }), - put: vi.fn() - } - const bot = new Bot({ storage } as unknown as DurableObjectState, { - ADMIN_TOKEN: 'test', BOT: {} as DurableObjectNamespace - }) - await bot.fetch(new Request('https://bot/start', { method: 'POST' })) + host.authenticate.mockResolvedValue({ creds: { registered: true } }) + const bot = new Bot(createContext(createStorage()), createEnv()) + await bot.fetch(request('/status')) expect(host.makeSocket.mock.calls[0]?.[0].auth.creds.registered).toBe(true) })