import { ulid as defaultUlid } from '../../domain/ulid'; import type { Expense, ExpenseBody, Settlement, SettlementBody } from '../../domain/types'; import { decodeText, type BucketStore } from '../storage/store'; import { ConflictError, NotFoundError, ReadOnlyError } from './errors'; import { decodeEvent, encodeEvent, encodeJsonl, EVENT_KEY_RE, EVENTS_PREFIX, eventKey, LEDGER_PATH, LedgerEventSchema, parseJsonl, type LedgerEvent } from './events'; import { applyEvent, foldEvents, type FoldedState } from './fold'; export type Clock = () => Date; export interface LedgerOptions { clock?: Clock; newId?: () => string; } export interface LedgerStats { pendingEvents: number; ledgerEvents: number; ledgerExists: boolean; } type DistributiveOmit = T extends unknown ? Omit : never; /** A LedgerEvent before the Ledger assigns id, ts and actor. */ export type NewEvent = DistributiveOmit; export interface CompactionResult { /** Events appended to the ledger. */ appended: number; /** Event files deleted. */ deleted: number; backedUp: boolean; } /** * In-memory view of all expenses and settlements, backed by an append-only event log. * All mutations go through a single async queue so writes never interleave. */ export class Ledger { readOnly = false; private queue: Promise = Promise.resolve(); private lastBackupDay: string | null = null; private constructor( private readonly store: BucketStore, private readonly state: FoldedState, /** Ids of events already present in ledger.jsonl. */ private readonly ledgerIds: Set, private ledgerExists: boolean, /** ledger.jsonl ends in a partial line, so the next compaction rewrites it instead of appending. */ private ledgerTruncated: boolean, /** Event files in events/ that have not been compacted yet, keyed by object path. */ private readonly pending: Map, private readonly clock: Clock, private readonly newId: () => string ) {} static async load(store: BucketStore, opts: LedgerOptions = {}): Promise { const clock = opts.clock ?? (() => new Date()); const newId = opts.newId ?? defaultUlid; const listing = await store.list(EVENTS_PREFIX); let ledgerBytes = await store.get(LEDGER_PATH); const pending = new Map(); let sawMissing = false; for (const entry of listing) { if (!EVENT_KEY_RE.test(entry.path)) continue; const bytes = await store.get(entry.path); if (!bytes) { sawMissing = true; continue; } try { pending.set(entry.path, decodeEvent(bytes)); } catch (e) { throw new Error(`Corrupt event file ${entry.path}`, { cause: e }); } } // A compactor deleted an event after publishing a ledger that contains it: re-read the ledger. if (sawMissing) ledgerBytes = await store.get(LEDGER_PATH); const parsed = ledgerBytes ? parseJsonl(decodeText(ledgerBytes)) : { events: [], truncated: false }; const ledgerIds = new Set(parsed.events.map((e) => e.id)); const sortedPending = [...pending.values()].sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)); const state = foldEvents([...parsed.events, ...sortedPending]); return new Ledger(store, state, ledgerIds, ledgerBytes !== null, parsed.truncated, pending, clock, newId); } // ---- reads ------------------------------------------------------------- getExpense(id: string): Expense | undefined { return this.state.expenses.get(id); } listExpenses(): Expense[] { return [...this.state.expenses.values()]; } getSettlement(id: string): Settlement | undefined { return this.state.settlements.get(id); } listSettlements(): Settlement[] { return [...this.state.settlements.values()]; } history(entityId: string): LedgerEvent[] { return [...(this.state.history.get(entityId) ?? [])]; } stats(): LedgerStats { return { pendingEvents: this.pending.size, ledgerEvents: this.ledgerIds.size, ledgerExists: this.ledgerExists }; } // ---- writes ------------------------------------------------------------ /** `id` may be supplied so callers can store a photo under the expense id before the event is written. */ createExpense(actor: string, body: ExpenseBody, id = `e_${this.newId()}`): Promise { return this.run(async () => { if (this.state.latest.has(id)) throw new ConflictError(id, this.state.latest.get(id)!.version); const now = this.clock().toISOString(); const expense: Expense = { ...body, id, version: 1, createdBy: actor, createdAt: now, updatedBy: actor, updatedAt: now }; const written = await this.writeEvent({ kind: 'expense', op: 'upsert', entityId: id, version: 1, data: expense }, actor, now); return written.kind === 'expense' && written.op === 'upsert' ? written.data : expense; }); } updateExpense(actor: string, id: string, body: ExpenseBody, expectedVersion: number): Promise { return this.run(async () => { const current = this.requireExpense(id, expectedVersion); const now = this.clock().toISOString(); const expense: Expense = { ...body, id, version: current.version + 1, createdBy: current.createdBy, createdAt: current.createdAt, updatedBy: actor, updatedAt: now }; const written = await this.writeEvent({ kind: 'expense', op: 'upsert', entityId: id, version: expense.version, data: expense }, actor, now); return written.kind === 'expense' && written.op === 'upsert' ? written.data : expense; }); } deleteExpense(actor: string, id: string, expectedVersion: number): Promise { return this.run(async () => { const current = this.requireExpense(id, expectedVersion); await this.writeEvent({ kind: 'expense', op: 'delete', entityId: id, version: current.version + 1 }, actor, this.clock().toISOString()); }); } createSettlement(actor: string, body: SettlementBody): Promise { return this.run(async () => { const now = this.clock().toISOString(); const settlement: Settlement = { ...body, id: `s_${this.newId()}`, version: 1, createdBy: actor, createdAt: now, updatedBy: actor, updatedAt: now }; const written = await this.writeEvent({ kind: 'settlement', op: 'upsert', entityId: settlement.id, version: 1, data: settlement }, actor, now); return written.kind === 'settlement' && written.op === 'upsert' ? written.data : settlement; }); } updateSettlement(actor: string, id: string, body: SettlementBody, expectedVersion: number): Promise { return this.run(async () => { const current = this.requireSettlement(id, expectedVersion); const now = this.clock().toISOString(); const settlement: Settlement = { ...body, id, version: current.version + 1, createdBy: current.createdBy, createdAt: current.createdAt, updatedBy: actor, updatedAt: now }; const written = await this.writeEvent({ kind: 'settlement', op: 'upsert', entityId: id, version: settlement.version, data: settlement }, actor, now); return written.kind === 'settlement' && written.op === 'upsert' ? written.data : settlement; }); } deleteSettlement(actor: string, id: string, expectedVersion: number): Promise { return this.run(async () => { const current = this.requireSettlement(id, expectedVersion); await this.writeEvent({ kind: 'settlement', op: 'delete', entityId: id, version: current.version + 1 }, actor, this.clock().toISOString()); }); } // ---- compaction -------------------------------------------------------- /** * Appends every pending event to ledger.jsonl, then deletes the event files. A ledger whose * last line was truncated by an interrupted append is rewritten in full instead. * Safe to crash at any point: nothing acknowledged is ever lost, and duplicates are removed by the fold. */ compact(): Promise { return this.run(async () => { if (this.readOnly) throw new ReadOnlyError(); const listing = await this.store.list(EVENTS_PREFIX); const keys = listing.map((e) => e.path).filter((p) => EVENT_KEY_RE.test(p)).sort(); // A delete may have succeeded while reporting failure; trust the bucket over `pending`. const present = new Set(keys); for (const key of this.pending.keys()) if (!present.has(key)) this.pending.delete(key); const toAppend: LedgerEvent[] = []; const toDelete: string[] = []; for (const key of keys) { let event = this.pending.get(key); if (!event) { const bytes = await this.store.get(key); if (!bytes) continue; event = decodeEvent(bytes); applyEvent(this.state, event); this.pending.set(key, event); } toDelete.push(key); if (!this.ledgerIds.has(event.id)) toAppend.push(event); } if (toDelete.length === 0 && !this.ledgerTruncated) return { appended: 0, deleted: 0, backedUp: false }; let backedUp = false; if (toAppend.length > 0 || this.ledgerTruncated) { const day = this.clock().toISOString().slice(0, 10); if (this.ledgerExists && this.lastBackupDay !== day) { await this.store.copy(LEDGER_PATH, `backups/ledger-${day}.jsonl`); this.lastBackupDay = day; backedUp = true; } if (this.ledgerTruncated) { await this.store.put(LEDGER_PATH, encodeJsonl([...this.compactedLedgerEvents(), ...toAppend]), 'application/x-ndjson'); this.ledgerTruncated = false; } else if (this.ledgerExists) { await this.store.append(LEDGER_PATH, encodeJsonl(toAppend)); } else { await this.store.put(LEDGER_PATH, encodeJsonl(toAppend), 'application/x-ndjson'); } this.ledgerExists = true; for (const e of toAppend) this.ledgerIds.add(e.id); } if (toDelete.length > 0) { await this.store.delete(toDelete); for (const key of toDelete) this.pending.delete(key); } return { appended: toAppend.length, deleted: toDelete.length, backedUp }; }); } /** Every event already durable in ledger.jsonl, recovered from the folded history for a full rewrite. */ private compactedLedgerEvents(): LedgerEvent[] { const byId = new Map(); for (const events of this.state.history.values()) { for (const event of events) if (this.ledgerIds.has(event.id)) byId.set(event.id, event); } return [...byId.values()].sort((a, b) => (a.id < b.id ? -1 : a.id > b.id ? 1 : 0)); } // ---- internals --------------------------------------------------------- /** * Serializes every mutation onto one chain. Non-reentrant: a queued function must never call * another queued method, or it deadlocks. Rejections propagate to the caller only; an unawaited * call is swallowed here to keep the chain alive, so callers must await to see failures. */ private run(fn: () => Promise): Promise { const next = this.queue.then(fn, fn); this.queue = next.catch(() => undefined); return next; } private requireExpense(id: string, expectedVersion: number): Expense { const current = this.state.expenses.get(id); if (!current) throw new NotFoundError(id); if (current.version !== expectedVersion) throw new ConflictError(id, current.version); return current; } private requireSettlement(id: string, expectedVersion: number): Settlement { const current = this.state.settlements.get(id); if (!current) throw new NotFoundError(id); if (current.version !== expectedVersion) throw new ConflictError(id, current.version); return current; } private async writeEvent(partial: NewEvent, actor: string, ts: string): Promise { if (this.readOnly) throw new ReadOnlyError(); const event = LedgerEventSchema.parse({ ...partial, id: this.newId(), ts, actor }); const key = eventKey(event.id); await this.store.put(key, encodeEvent(event), 'application/json'); applyEvent(this.state, event); this.pending.set(key, event); return event; } }