Luigit
repositories / will

will

owned by admin

src/agent/outbox.ts

Raw
import { mkdir, readFile, rename, writeFile } from 'node:fs/promises'
import { dirname, join } from 'node:path'
import type { ChatMessageIn, MediaRef } from '../shared/events.ts'
import type { ServiceClient } from './service-client.ts'
import type { TelegramClient } from './telegram-client.ts'

/**
 * At-least-once outbound Telegram delivery (spec AC18).
 * Records are persisted before the first send attempt and retried with backoff;
 * rare duplicates after ambiguous network failures are accepted.
 */

export interface OutboxRecord {
  id: string
  chatId: number
  text?: string | undefined
  mediaDigest?: string | undefined
  mediaType?: string | undefined
  mediaName?: string | undefined
  replyToMessageId?: number | undefined
  status: 'pending' | 'sent' | 'failed'
  attempts: number
  nextAttemptAt: number
  sentMessageId?: number | undefined
}

export class Outbox {
  readonly #file: string
  readonly #records = new Map<string, OutboxRecord>()

  private constructor(file: string) {
    this.#file = file
  }

  static async open(file: string): Promise<Outbox> {
    const out = new Outbox(file)
    try {
      for (const line of (await readFile(file, 'utf8')).split('\n')) {
        if (line === '') continue
        const rec = JSON.parse(line) as OutboxRecord
        if (rec.status !== 'sent') out.#records.set(rec.id, rec)
      }
    } catch {
      // fresh outbox
    }
    return out
  }

  async enqueue(
    rec: Omit<OutboxRecord, 'status' | 'attempts' | 'nextAttemptAt'>,
  ): Promise<OutboxRecord> {
    const full: OutboxRecord = { ...rec, status: 'pending', attempts: 0, nextAttemptAt: Date.now() }
    this.#records.set(full.id, full)
    await this.#persist()
    return full
  }

  pending(): OutboxRecord[] {
    return [...this.#records.values()].filter((r) => r.status === 'pending')
  }

  /**
   * Flush due records. Delivery is at-least-once; the service chat mirror is
   * ingested after a successful send (idempotent by message id).
   */
  async flush(opts: {
    tg: TelegramClient
    service: ServiceClient
    now?: () => number
    maxAttempts?: number
  }): Promise<{ sent: number; failed: number }> {
    const now = opts.now ?? Date.now
    const maxAttempts = opts.maxAttempts ?? 8
    let sent = 0
    let failed = 0
    for (const rec of this.pending()) {
      if (rec.nextAttemptAt > now()) continue
      try {
        const delivered = rec.mediaDigest
          ? await opts.tg.sendMedia(rec.chatId, {
              bytes: await opts.service.getMedia(rec.mediaDigest),
              type: rec.mediaType ?? 'application/octet-stream',
              ...(rec.mediaName !== undefined ? { fileName: rec.mediaName } : {}),
            })
          : await opts.tg.sendMessage(rec.chatId, rec.text ?? '', rec.replyToMessageId)
        rec.status = 'sent'
        rec.sentMessageId = delivered.messageId
        sent++
        const mirror: ChatMessageIn = {
          id: rec.id,
          chatId: rec.chatId,
          messageId: delivered.messageId,
          direction: 'out',
          ts: new Date().toISOString(),
          ...(rec.replyToMessageId !== undefined ? { replyToMessageId: rec.replyToMessageId } : {}),
          ...(rec.text !== undefined ? { text: rec.text } : {}),
          ...(rec.mediaDigest
            ? {
                media: [
                  {
                    digest: rec.mediaDigest,
                    type: rec.mediaType ?? 'application/octet-stream',
                    size: 0,
                  } satisfies MediaRef,
                ],
              }
            : {}),
        }
        await opts.service.ingestChat([mirror])
      } catch {
        rec.attempts++
        if (rec.attempts >= maxAttempts) {
          rec.status = 'failed'
          failed++
        } else {
          rec.nextAttemptAt = now() + backoffMs(rec.attempts)
        }
      }
    }
    await this.#persist()
    return { sent, failed }
  }

  async #persist(): Promise<void> {
    await mkdir(dirname(this.#file), { recursive: true })
    const tmp = join(dirname(this.#file), `.outbox-${process.pid}.tmp`)
    const live = [...this.#records.values()].filter((r) => r.status !== 'sent')
    const failedKeep = [...this.#records.values()].filter((r) => r.status === 'failed')
    const all = [...live, ...failedKeep.filter((r) => !live.includes(r))]
    const lines = all.map((r) => JSON.stringify(r))
    await writeFile(tmp, lines.length > 0 ? `${lines.join('\n')}\n` : '', 'utf8')
    await rename(tmp, this.#file)
  }
}

function backoffMs(attempt: number): number {
  return Math.min(60_000 * 2 ** (attempt - 1), 3_600_000)
}