repositories / will
will
owned by admin
src/agent/telegram.int.test.ts
Rawimport 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<string, unknown> }[]
botId: number
botUsername: string
}
async function startFakeTelegram(): Promise<FakeTelegram> {
const updates: TgUpdateRaw[] = []
const sent: { method: string; body: Record<string, unknown> }[] = []
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<string, unknown> = {}
try {
body = data === '' ? {} : (JSON.parse(data) as Record<string, unknown>)
} 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<void>((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()
}
})