Luigit
repositories / will

will

owned by admin

src/agent/telegram.int.test.ts

Raw
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<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()
  }
})