File size: 11,792 Bytes
3bfa0fc
 
 
 
664d562
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
664d562
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
664d562
 
 
8d2f1b6
664d562
3bfa0fc
 
 
 
664d562
 
3bfa0fc
664d562
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
8d2f1b6
3bfa0fc
 
664d562
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
664d562
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
 
664d562
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
664d562
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
664d562
 
3bfa0fc
 
 
 
 
 
 
 
664d562
 
 
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
 
664d562
3bfa0fc
 
664d562
 
 
 
 
 
 
3bfa0fc
664d562
 
 
 
 
 
 
 
3bfa0fc
 
 
 
664d562
 
 
 
3bfa0fc
 
 
 
664d562
 
 
 
 
 
 
 
 
3bfa0fc
 
664d562
 
 
 
 
3bfa0fc
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
664d562
3bfa0fc
664d562
3bfa0fc
 
 
 
664d562
3bfa0fc
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
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;
	}
}