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() private constructor(file: string) { this.#file = file } static async open(file: string): Promise { 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, ): Promise { 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 { 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) }