repositories / will
will
owned by admin
src/agent/outbox.ts
Rawimport { 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)
}