Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
70 changes: 60 additions & 10 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,58 +11,108 @@ interface Env {
ADMIN_TOKEN: string
}

const RECONNECT_DELAY_MS = 30_000

export class Bot extends DurableObject<Env> {
private socket?: HostWASocket
private status: BotStatus = { state: 'stopped' }
private starting?: Promise<void>
private retryAt = 0
private readonly initialization: Promise<void>

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<Response> {
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<void> {
if (this.socket) return
private start(allowPairing: boolean): Promise<void> {
if (this.socket) return Promise.resolve()
if (this.starting) return this.starting

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve pairing intent when sharing a start

If a pre-registration socket closes and, after the cooldown, a /status retry overlaps with POST /start while authentication is pending, the status request installs a non-pairing starting promise and this branch makes the explicit start share it. Because the stored credentials are still unregistered, openSocket(false) exits with stopped, so the POST succeeds without creating a socket or QR; preserve or upgrade allowPairing when a pairing request joins an in-progress automatic attempt.

Useful? React with 👍 / 👎.

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<void> {
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
}
}

Expand Down
131 changes: 93 additions & 38 deletions test/index.test.ts
Original file line number Diff line number Diff line change
@@ -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(),
Expand All @@ -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<string, unknown>()) => ({
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<typeof createStorage>) => ({
storage,
blockConcurrencyWhile: async (callback: () => Promise<void>) => callback()
}) as unknown as DurableObjectState

const createEnv = () => ({ ADMIN_TOKEN: 'test', BOT: {} as DurableObjectNamespace<Bot> })

const latestConnectionUpdate = () => host.on.mock.calls
.filter(([event]) => event === 'connection.update')
.at(-1)?.[1] as (event: HostBaileysEventMap['connection.update']) => Promise<void> | 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<string, unknown>()
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<Bot> }
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> | 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<string, unknown>([['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<Bot>
})
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<void>
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<Bot>
})
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)
})
Loading