import assert from 'node:assert/strict' import { mkdtemp } from 'node:fs/promises' import { createServer, type Server } from 'node:http' import { tmpdir } from 'node:os' import { join } from 'node:path' import { test } from 'node:test' import { startService } from '../service/main.ts' import type { TgUpdateRaw } from '../shared/telegram.ts' import { Outbox } from './outbox.ts' import { WakeQueue } from './queue-store.ts' import { createServiceClient, type ServiceClient } from './service-client.ts' import { AuditSpool } from './spool.ts' import { AgentState } from './state.ts' import { TelegramClient } from './telegram-client.ts' import { createTelegramRuntime } from './telegram-runtime.ts' interface FakeTelegram { server: Server url: string updates: TgUpdateRaw[] sent: { method: string; body: Record }[] botId: number botUsername: string } async function startFakeTelegram(): Promise { const updates: TgUpdateRaw[] = [] const sent: { method: string; body: Record }[] = [] const state = { botId: 42, botUsername: 'will_bot' } const server = createServer((req, res) => { const url = new URL(req.url ?? '/', 'http://localhost') const parts = url.pathname.split('/') const method = parts[parts.length - 1] ?? '' let data = '' req.on('data', (c: Buffer) => { data += c.toString('utf8') }) req.on('end', () => { let body: Record = {} try { body = data === '' ? {} : (JSON.parse(data) as Record) } catch { body = {} } const json = (payload: unknown) => { res.writeHead(200, { 'content-type': 'application/json' }) res.end(JSON.stringify(payload)) } if (url.pathname.includes('/file/bot')) { res.writeHead(200, { 'content-type': 'image/png' }) res.end(Buffer.from('fake-png-bytes')) return } switch (method) { case 'getMe': return json({ ok: true, result: { id: state.botId, username: state.botUsername } }) case 'getUpdates': { const offset = (body.offset as number | undefined) ?? 0 return json({ ok: true, result: updates.filter((u) => u.update_id >= offset) }) } case 'getFile': return json({ ok: true, result: { file_id: 'f1', file_path: 'files/media.png' } }) default: sent.push({ method, body }) return json({ ok: true, result: { message_id: 900 + sent.length } }) } }) }) await new Promise((res) => server.listen(0, '127.0.0.1', () => res())) const url = `http://127.0.0.1:${(server.address() as { port: number }).port}` return { server, url, updates, sent, botId: state.botId, botUsername: state.botUsername, } } test('fake telegram sanity', async () => { const fake = await startFakeTelegram() const tg = new TelegramClient('TOKEN', { apiBase: fake.url, fileBase: fake.url }) const me = await tg.getMe() assert.equal(me.username, 'will_bot') fake.server.close() }) async function setupStack() { const dir = await mkdtemp(join(tmpdir(), 'tg-rt-')) const webDist = join(dir, 'web-dist') await import('node:fs/promises').then((fs) => fs.mkdir(webDist, { recursive: true })) const service = await startService({ dataDir: dir, publicAddr: '127.0.0.1:0', internalAddr: '127.0.0.1:0', webDist, docsDist: webDist, mediaMaxBytes: 1000 << 10, storageWarnPct: 20, storageStopPct: 10, workerUrl: join(import.meta.dirname, '../service/worker.ts'), }) const intAddr = service.internalServer.address() if (intAddr === null || typeof intAddr !== 'object') throw new Error('no addr') const client: ServiceClient = createServiceClient(`http://127.0.0.1:${intAddr.port}`) const queue = await WakeQueue.open(join(dir, 'agent', 'wakeq.jsonl')) const spool = await AuditSpool.open(join(dir, 'agent', 'spool')) const state = await AgentState.load(join(dir, 'agent', 'state.json')) return { dir, service, client, queue, spool, state } } test('mention wakes, chatter archives without wake, other chats ignored', async () => { const fake = await startFakeTelegram() const { service, client, queue, spool, state } = await setupStack() try { const tg = new TelegramClient('TOKEN', { apiBase: fake.url, fileBase: fake.url }) const rt = createTelegramRuntime({ tg, service: client, queue, spool, state, chatId: -100, oliverId: '1', }) await rt.identifyBot() fake.updates.push( { update_id: 1, message: { message_id: 10, date: 1700000000, chat: { id: -100 }, from: { id: 7, username: 'friend' }, text: 'hey @will_bot whats up', entities: [{ type: 'mention', offset: 4, length: 9 }], }, }, { update_id: 2, message: { message_id: 11, date: 1700000060, chat: { id: -100 }, from: { id: 8, username: 'other' }, text: 'just chatting', }, }, { update_id: 3, message: { message_id: 12, date: 1700000120, chat: { id: -999 }, from: { id: 9, username: 'intruder' }, text: 'let me in @will_bot', entities: [{ type: 'mention', offset: 11, length: 9 }], }, }, ) const summary = await rt.pollOnce() assert.equal(summary.archived, 2) assert.equal(summary.wakes, 1) assert.equal(summary.ignored, 1) assert.equal(state.data.telegramOffset, 4) const pending = queue.peek() assert.equal(pending.length, 1) assert.equal(pending[0]?.kind, 'chat') const payload = pending[0]?.payload as { text: string; triggerReason: string } assert.equal(payload.triggerReason, 'mention') const page = await client.recall({ kind: 'chat', limit: 10 }) assert.equal(page.rows.length, 2) } finally { await service.close() fake.server.close() } }) test('media message downloads into service store and references digest', async () => { const fake = await startFakeTelegram() const { service, client, queue, spool, state } = await setupStack() try { const tg = new TelegramClient('TOKEN', { apiBase: fake.url, fileBase: fake.url }) const rt = createTelegramRuntime({ tg, service: client, queue, spool, state, chatId: -100, oliverId: '1', }) await rt.identifyBot() fake.updates.push({ update_id: 5, message: { message_id: 20, date: 1700000000, chat: { id: -100 }, from: { id: 7, username: 'friend' }, caption: '@will_bot what is this', entities: [{ type: 'mention', offset: 0, length: 9 }], photo: [ { file_id: 'small', width: 1, height: 1 }, { file_id: 'big', width: 2, height: 2 }, ], }, }) const summary = await rt.pollOnce() assert.equal(summary.wakes, 1) const payload = queue.peek()[0]?.payload as { media: { digest: string }[] } assert.ok(payload.media[0]?.digest) const status = await client.status() assert.equal(status.counts.media, 1) assert.equal(status.counts.chat, 1) } finally { await service.close() fake.server.close() } }) test('/hide via reply tombstones deterministically without a model wake', async () => { const fake = await startFakeTelegram() const { service, client, queue, spool, state } = await setupStack() try { const tg = new TelegramClient('TOKEN', { apiBase: fake.url, fileBase: fake.url }) const rt = createTelegramRuntime({ tg, service: client, queue, spool, state, chatId: -100, oliverId: '1', }) await rt.identifyBot() fake.updates.push( { update_id: 1, message: { message_id: 30, date: 1700000000, chat: { id: -100 }, from: { id: 7, username: 'friend' }, text: 'embarrassing message', }, }, { update_id: 2, message: { message_id: 31, date: 1700000100, chat: { id: -100 }, from: { id: 7, username: 'friend' }, text: '/hide', reply_to_message: { message_id: 30, from: { id: 7 } }, }, }, ) const summary = await rt.pollOnce() assert.equal(summary.hidden, 1) assert.equal(summary.wakes, 0) const page = await client.recall({ kind: 'chat', limit: 10 }) assert.equal(page.rows.length, 1, 'hidden message excluded from recall') } finally { await service.close() fake.server.close() } }) test('outbox delivers at-least-once and mirrors outbound to service', async () => { const fake = await startFakeTelegram() const { dir, service, client } = await setupStack() try { const outbox = await Outbox.open(join(dir, 'agent', 'outbox.jsonl')) const tg = new TelegramClient('TOKEN', { apiBase: fake.url, fileBase: fake.url }) await outbox.enqueue({ id: 'out-1', chatId: -100, text: 'hello from will' }) const result = await outbox.flush({ tg, service: client }) assert.equal(result.sent, 1) assert.equal(fake.sent.filter((s) => s.method === 'sendMessage').length, 1) const page = await client.recall({ kind: 'chat', limit: 10 }) assert.equal(page.rows.length, 1) // ambiguous failure: first send throws, retry succeeds (at-least-once) await outbox.enqueue({ id: 'out-2', chatId: -100, text: 'retry me' }) let fail = true const realSend = tg.sendMessage.bind(tg) const flaky = new TelegramClient('TOKEN', { apiBase: fake.url, fileBase: fake.url }) flaky.sendMessage = async (chatId: number, text: string, reply?: number) => { if (fail) { fail = false throw new Error('network hiccup') } return realSend(chatId, text, reply) } let clockNow = Date.now() const first = await outbox.flush({ tg: flaky, service: client, now: () => clockNow }) assert.equal(first.sent, 0) clockNow += 120_000 // past the first retry backoff const second = await outbox.flush({ tg: flaky, service: client, now: () => clockNow }) assert.equal(second.sent, 1) } finally { await service.close() fake.server.close() } })