splitwise / src /lib /server /ledger /ledger.ts
assafvayner's picture
assafvayner HF Staff
test(ledger): prove compaction crash safety with fault injection and against the real bucket
8d2f1b6
Raw History Blame Contribute Delete
11.8 kB
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, K extends PropertyKey> = T extends unknown ? Omit<T, K> : never;
/** A LedgerEvent before the Ledger assigns id, ts and actor. */
export type NewEvent = DistributiveOmit<LedgerEvent, 'id' | 'ts' | 'actor'>;
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<unknown> = 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<string>,
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<string, LedgerEvent>,
private readonly clock: Clock,
private readonly newId: () => string
) {}
static async load(store: BucketStore, opts: LedgerOptions = {}): Promise<Ledger> {
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<string, LedgerEvent>();
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<Expense> {
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<Expense> {
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<void> {
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<Settlement> {
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<Settlement> {
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<void> {
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<CompactionResult> {
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<string, LedgerEvent>();
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<T>(fn: () => Promise<T>): Promise<T> {
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<LedgerEvent> {
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;
}
}