import { afterAll, beforeAll, beforeEach, describe, expect, test } from "bun:test"; import { type BunTestDb, createTestDb } from "../../bun-db/__tests__/bun-test-db"; import { asRawClient } from "../../db/query"; import { ensureTemporalPolyfill } from "../../time/polyfill"; import { generateId as uuid } from "../../utils"; import { append, createEventsTable, IdempotentAppendConflictError, loadAggregate, loadAggregateAsOf, loadAllEventsByType, loadEventsAfterVersion, type StoredEvent, streamAllEventsByType, VersionConflictError, } from "../index"; let testDb: BunTestDb; const tenantA = uuid(); const tenantB = uuid(); const userA = uuid(); beforeAll(async () => { await ensureTemporalPolyfill(); testDb = await createTestDb(); await createEventsTable(testDb.db); }); afterAll(async () => { await testDb.cleanup(); }); beforeEach(async () => { await asRawClient(testDb.db).unsafe(`TRUNCATE kumiko_events RESTART IDENTITY`); }); describe("event-store: append + load", () => { test("append first event writes version=1 and round-trips via loadAggregate", async () => { const aggregateId = uuid(); const stored = await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Buy milk" }, metadata: { userId: userA }, }); expect(stored.version).toBe(1); expect(stored.id).toBeDefined(); const events = await loadAggregate(testDb.db, aggregateId, tenantA); expect(events).toHaveLength(1); expect(events[0]?.type).toBe("task.created"); expect(events[0]?.payload).toEqual({ title: "Buy milk" }); expect(events[0]?.metadata.userId).toBe(userA); }); test("kumiko-framework#1490: subsequent append doesn't rely on globalThis.Temporal", async () => { const aggregateId = uuid(); const first = await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Buy milk" }, metadata: { userId: userA }, }); const savedGlobal = (globalThis as { Temporal?: unknown }).Temporal; delete (globalThis as { Temporal?: unknown }).Temporal; try { const second = await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: first.version, type: "task.updated", payload: { title: "Buy oat milk" }, metadata: { userId: userA }, }); expect(second.version).toBe(2); expect(second.createdAt).toBeDefined(); } finally { if (savedGlobal === undefined) delete (globalThis as { Temporal?: unknown }).Temporal; else (globalThis as { Temporal?: unknown }).Temporal = savedGlobal; } }); test("subsequent appends increment version and are ordered", async () => { const aggregateId = uuid(); const base = { aggregateId, aggregateType: "task", tenantId: tenantA, metadata: { userId: userA }, }; await append(testDb.db, { ...base, expectedVersion: 0, type: "task.created", payload: { title: "T" }, }); await append(testDb.db, { ...base, expectedVersion: 1, type: "task.updated", payload: { title: "T2" }, }); await append(testDb.db, { ...base, expectedVersion: 2, type: "task.completed", payload: {}, }); const events = await loadAggregate(testDb.db, aggregateId, tenantA); expect(events.map((e) => e.version)).toEqual([1, 2, 3]); expect(events.map((e) => e.type)).toEqual(["task.created", "task.updated", "task.completed"]); }); }); describe("event-store: idempotency-key conflict", () => { test("second append with a reused idempotencyKey (new aggregate+version) throws IdempotentAppendConflictError", async () => { const first = uuid(); const second = uuid(); const key = uuid(); await append(testDb.db, { aggregateId: first, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Orig" }, metadata: { userId: userA, idempotencyKey: key }, }); // A retried command that re-runs after the Redis idempotency guard // missed its window: different aggregate, fresh expectedVersion=0 — the // aggregate-version unique index has nothing to say about this pair, so // only the idempotency-key index catches the duplicate. await expect( append(testDb.db, { aggregateId: second, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Retry" }, metadata: { userId: userA, idempotencyKey: key }, }), ).rejects.toThrow(IdempotentAppendConflictError); const events = await loadAggregate(testDb.db, second, tenantA); expect(events).toHaveLength(0); }); test("subsequent-event append (expectedVersion > 0) also enforces the idempotency key", async () => { const aggregateId = uuid(); const key = uuid(); await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Orig" }, metadata: { userId: userA }, }); await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 1, type: "task.updated", payload: { title: "V2" }, metadata: { userId: userA, idempotencyKey: key }, }); // Retry of the v2 update: goes through insertSubsequentEventRow's raw // INSERT ... SELECT ... WHERE EXISTS path, not insertFirstEvent — the // idempotency index must catch it there too. await expect( append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 2, type: "task.updated", payload: { title: "V3-retry" }, metadata: { userId: userA, idempotencyKey: key }, }), ).rejects.toThrow(IdempotentAppendConflictError); const events = await loadAggregate(testDb.db, aggregateId, tenantA); expect(events).toHaveLength(2); }); test("same idempotencyKey on a different tenant does not conflict", async () => { const key = uuid(); const a = await append(testDb.db, { aggregateId: uuid(), aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: {}, metadata: { userId: userA, idempotencyKey: key }, }); const b = await append(testDb.db, { aggregateId: uuid(), aggregateType: "task", tenantId: tenantB, expectedVersion: 0, type: "task.created", payload: {}, metadata: { userId: userA, idempotencyKey: key }, }); expect(a.version).toBe(1); expect(b.version).toBe(1); }); test("omitting idempotencyKey allows unlimited appends, unchanged from before", async () => { // Same tenant for both, unlike the different-aggregateId version this // replaces — Postgres treats every NULL as distinct even in a plain // unique index, so this does NOT exercise the partial index's `WHERE // ... IS NOT NULL` clause specifically (that's provable only by a // duplicate NON-null key, covered above). What this does pin: a fully- // omitted key and an explicit `idempotencyKey: undefined` both serialize // to a JSON-absent key (JSON.stringify drops undefined) and must behave // identically, rather than one silently colliding on `"idempotencyKey":null`. const omitted = await append(testDb.db, { aggregateId: uuid(), aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: {}, metadata: { userId: userA }, }); const explicitUndefined = await append(testDb.db, { aggregateId: uuid(), aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: {}, metadata: { userId: userA, idempotencyKey: undefined }, }); expect(omitted.version).toBe(1); expect(explicitUndefined.version).toBe(1); }); }); describe("event-store: optimistic concurrency", () => { test("wrong expectedVersion throws VersionConflictError (no write)", async () => { const aggregateId = uuid(); await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Orig" }, metadata: { userId: userA }, }); // Stale writer: thinks predecessor is at v0 — but v1 already exists. await expect( append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.updated", payload: { title: "Stale" }, metadata: { userId: userA }, }), ).rejects.toThrow(VersionConflictError); const events = await loadAggregate(testDb.db, aggregateId, tenantA); expect(events).toHaveLength(1); }); test("concurrent writers at same expectedVersion: exactly one wins", async () => { const aggregateId = uuid(); await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Orig" }, metadata: { userId: userA }, }); const update = (label: string) => append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 1, type: "task.updated", payload: { title: label }, metadata: { userId: userA }, }); const results = await Promise.allSettled([update("A"), update("B")]); const fulfilled = results.filter((r) => r.status === "fulfilled"); const rejected = results.filter((r) => r.status === "rejected"); expect(fulfilled).toHaveLength(1); expect(rejected).toHaveLength(1); expect((rejected[0] as PromiseRejectedResult).reason).toBeInstanceOf(VersionConflictError); const events = await loadAggregate(testDb.db, aggregateId, tenantA); expect(events).toHaveLength(2); }); }); describe("event-store: tenant isolation", () => { test("cross-tenant append at same aggregateId is rejected", async () => { const aggregateId = uuid(); await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: {}, metadata: { userId: userA }, }); // Tenant B tries to write v2 against A's v1 — predecessor check must fail. await expect( append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantB, expectedVersion: 1, type: "task.updated", payload: {}, metadata: { userId: userA }, }), ).rejects.toThrow(VersionConflictError); const eventsA = await loadAggregate(testDb.db, aggregateId, tenantA); const eventsB = await loadAggregate(testDb.db, aggregateId, tenantB); expect(eventsA).toHaveLength(1); expect(eventsB).toHaveLength(0); }); }); describe("event-store: requestId is a trace marker (no DB-level uniqueness)", () => { test("same (tenant, requestId) twice → both events persist, no collision", async () => { // Idempotency is an HTTP-level concern, handled via Redis in // pipeline/idempotency.ts before the command executes. The events-table // imposes no uniqueness on metadata.requestId — a single request may // write N events (CRUD + ctx.appendEvent + saga follow-ups), all // carrying the same requestId as a trace marker. const aggregateId1 = uuid(); const aggregateId2 = uuid(); const requestId = uuid(); const first = await append(testDb.db, { aggregateId: aggregateId1, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "First" }, metadata: { userId: userA, requestId }, }); const second = await append(testDb.db, { aggregateId: aggregateId2, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "Second" }, metadata: { userId: userA, requestId }, }); expect(first.metadata.requestId).toBe(requestId); expect(second.metadata.requestId).toBe(requestId); expect(first.aggregateId).not.toBe(second.aggregateId); }); test("metadata.headers (Marten free key/value) round-trips via append + load", async () => { const aggregateId = uuid(); const headers = { abTestBucket: "control", sdkVersion: 42, betaFeatures: true, }; await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "with headers" }, metadata: { userId: userA, headers }, }); // Subsequent event uses the WHERE-EXISTS raw-SQL path — make sure // headers survive that route too, not just the typed insertFirstEvent. await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 1, type: "task.updated", payload: { title: "v2" }, metadata: { userId: userA, headers: { ...headers, sdkVersion: 43 } }, }); const events = await loadAggregate(testDb.db, aggregateId, tenantA); expect(events).toHaveLength(2); expect(events[0]?.metadata.headers).toEqual(headers); expect(events[1]?.metadata.headers).toEqual({ ...headers, sdkVersion: 43 }); }); }); describe("event-store: asOf + after-version reads", () => { test("loadAggregateAsOf excludes events after the timestamp", async () => { const aggregateId = uuid(); const e1 = await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: {}, metadata: { userId: userA }, }); // Ensure the second event is strictly after. await new Promise((r) => setTimeout(r, 5)); await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: 1, type: "task.updated", payload: {}, metadata: { userId: userA }, }); const atT1 = await loadAggregateAsOf(testDb.db, aggregateId, tenantA, e1.createdAt); expect(atT1).toHaveLength(1); expect(atT1[0]?.version).toBe(1); }); test("loadEventsAfterVersion returns only events strictly > given version", async () => { const aggregateId = uuid(); for (let v = 0; v < 3; v++) { await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: v, type: v === 0 ? "task.created" : "task.updated", payload: { n: v }, metadata: { userId: userA }, }); } const after1 = await loadEventsAfterVersion(testDb.db, aggregateId, tenantA, 1); expect(after1.map((e) => e.version)).toEqual([2, 3]); }); }); describe("event-store: loadAllEventsByType", () => { // Backbone of projection-rebuild replay: all events of one aggregateType, // cross-tenant, in chronological order. Flagged as untested in a prior // audit — these tests close the gap. test("returns only events of the requested aggregateType", async () => { const taskId = uuid(); const invoiceId = uuid(); await append(testDb.db, { aggregateId: taskId, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { title: "T" }, metadata: { userId: userA }, }); await append(testDb.db, { aggregateId: invoiceId, aggregateType: "invoice", tenantId: tenantA, expectedVersion: 0, type: "invoice.created", payload: { amount: 42 }, metadata: { userId: userA }, }); const taskEvents = await loadAllEventsByType(testDb.db, "task"); const invoiceEvents = await loadAllEventsByType(testDb.db, "invoice"); expect(taskEvents).toHaveLength(1); expect(taskEvents[0]?.aggregateId).toBe(taskId); expect(taskEvents[0]?.type).toBe("task.created"); expect(invoiceEvents).toHaveLength(1); expect(invoiceEvents[0]?.aggregateId).toBe(invoiceId); }); test("spans all tenants — projection rebuild must see every row", async () => { // Rebuilds run system-scoped (cross-tenant) because a projection table // can hold data from many tenants. Missing this would leak tenant B's // absence into tenant A's projection snapshot after a rebuild. const aggA = uuid(); const aggB = uuid(); await append(testDb.db, { aggregateId: aggA, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { owner: "A" }, metadata: { userId: userA }, }); await append(testDb.db, { aggregateId: aggB, aggregateType: "task", tenantId: tenantB, expectedVersion: 0, type: "task.created", payload: { owner: "B" }, metadata: { userId: userA }, }); const all = await loadAllEventsByType(testDb.db, "task"); expect(all).toHaveLength(2); const tenants = new Set(all.map((e) => e.tenantId)); expect(tenants).toEqual(new Set([tenantA, tenantB])); }); test("ordered by (createdAt, id) for deterministic replay", async () => { // Projection rebuild applies events in the order they were written. // The ordering is part of the contract — without it, different // replays would produce different projection states. const a1 = uuid(); const a2 = uuid(); const a3 = uuid(); const e1 = await append(testDb.db, { aggregateId: a1, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { n: 1 }, metadata: { userId: userA }, }); await new Promise((r) => setTimeout(r, 5)); const e2 = await append(testDb.db, { aggregateId: a2, aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { n: 2 }, metadata: { userId: userA }, }); await new Promise((r) => setTimeout(r, 5)); const e3 = await append(testDb.db, { aggregateId: a3, aggregateType: "task", tenantId: tenantB, expectedVersion: 0, type: "task.created", payload: { n: 3 }, metadata: { userId: userA }, }); const all = await loadAllEventsByType(testDb.db, "task"); expect(all.map((e) => e.id)).toEqual([e1.id, e2.id, e3.id]); // createdAt strictly non-decreasing for (let i = 1; i < all.length; i++) { const prev = all[i - 1]; const cur = all[i]; if (!prev || !cur) throw new Error("unreachable"); expect(Temporal.Instant.compare(prev.createdAt, cur.createdAt)).toBeLessThanOrEqual(0); } }); test("returns empty array when no events of that type exist", async () => { const events = await loadAllEventsByType(testDb.db, "nonexistent-type"); expect(events).toEqual([]); }); test("includes every event of an aggregate — multiple versions in order", async () => { // A single aggregate with multiple versions must appear in order in the // replay stream, otherwise projection-apply sees events out-of-sequence. const aggregateId = uuid(); for (let v = 0; v < 4; v++) { await append(testDb.db, { aggregateId, aggregateType: "task", tenantId: tenantA, expectedVersion: v, type: v === 0 ? "task.created" : "task.updated", payload: { v }, metadata: { userId: userA }, }); } const all = await loadAllEventsByType(testDb.db, "task"); expect(all).toHaveLength(4); expect(all.map((e) => e.version)).toEqual([1, 2, 3, 4]); expect(all.map((e) => (e.payload as { v: number }).v)).toEqual([0, 1, 2, 3]); }); test("throws once rows exceed the given rowLimit", async () => { for (let v = 0; v < 3; v++) { await append(testDb.db, { aggregateId: uuid(), aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { v }, metadata: { userId: userA }, }); } await expect(loadAllEventsByType(testDb.db, "task", 2)).rejects.toThrow(/exceeds 2 rows/); }); test("does not throw when rows equal the given rowLimit", async () => { for (let v = 0; v < 2; v++) { await append(testDb.db, { aggregateId: uuid(), aggregateType: "task", tenantId: tenantA, expectedVersion: 0, type: "task.created", payload: { v }, metadata: { userId: userA }, }); } const all = await loadAllEventsByType(testDb.db, "task", 2); expect(all).toHaveLength(2); }); }); describe("event-store: streamAllEventsByType (memory-bounded iteration)", () => { test("yields every event in id order across multiple batches", async () => { // Seed 25 events; with batchSize=10 that's 3 batches (10+10+5). // Verifies cursor advance (no skipping between batches) and final // empty-batch termination. for (let i = 0; i < 25; i++) { await append(testDb.db, { aggregateId: uuid(), aggregateType: "stream-task", tenantId: tenantA, expectedVersion: 0, type: "stream-task.created", payload: { i }, metadata: { userId: userA }, }); } const collected: Array<{ id: string; i: number }> = []; for await (const event of streamAllEventsByType(testDb.db, "stream-task", 10)) { collected.push({ id: event.id, i: (event.payload as { i: number }).i }); } expect(collected).toHaveLength(25); // Reihenfolge nach events.id (= chronological commit order). expect(collected.map((c) => c.i)).toEqual(Array.from({ length: 25 }, (_, n) => n)); // ids strict aufsteigend (bigserial monotonic). for (let i = 1; i < collected.length; i++) { expect(BigInt(collected[i]!.id)).toBeGreaterThan(BigInt(collected[i - 1]!.id)); } }); test("empty store yields nothing", async () => { const collected: StoredEvent[] = []; for await (const event of streamAllEventsByType(testDb.db, "nonexistent", 10)) { collected.push(event); } expect(collected).toEqual([]); }); test("filters by aggregateType — other types stay unstreamed", async () => { await append(testDb.db, { aggregateId: uuid(), aggregateType: "stream-included", tenantId: tenantA, expectedVersion: 0, type: "stream-included.x", payload: {}, metadata: { userId: userA }, }); await append(testDb.db, { aggregateId: uuid(), aggregateType: "stream-excluded", tenantId: tenantA, expectedVersion: 0, type: "stream-excluded.x", payload: {}, metadata: { userId: userA }, }); const yielded: string[] = []; for await (const event of streamAllEventsByType(testDb.db, "stream-included")) { yielded.push(event.aggregateType); } expect(yielded).toEqual(["stream-included"]); }); test("per-yield abort: stops at exactly the event after abort, regardless of batch size", async () => { // Seed 25 events. Use batchSize=10 so an abort at length=5 lands // mid-batch — verifies the per-yield check, not just batch-boundary // semantics. The generator throws on its next yield after abort, so // collected.length stays at 5. for (let i = 0; i < 25; i++) { await append(testDb.db, { aggregateId: uuid(), aggregateType: "stream-abort", tenantId: tenantA, expectedVersion: 0, type: "stream-abort.x", payload: { i }, metadata: { userId: userA }, }); } const controller = new AbortController(); const collected: StoredEvent[] = []; let thrown: unknown; try { for await (const event of streamAllEventsByType( testDb.db, "stream-abort", 10, controller.signal, )) { collected.push(event); if (collected.length === 5) controller.abort(); } } catch (e) { thrown = e; } expect(collected.length).toBe(5); expect(thrown).toBeInstanceOf(Error); expect((thrown as Error).name).toBe("AbortError"); }); test("pre-aborted signal throws before any rows are fetched", async () => { const controller = new AbortController(); controller.abort(); let thrown: unknown; try { for await (const _event of streamAllEventsByType( testDb.db, "stream-included", 10, controller.signal, )) { // unreachable } } catch (e) { thrown = e; } expect((thrown as Error).name).toBe("AbortError"); }); }); describe("event-store: jsonb encoding of payload/metadata", () => { // Regression für die Bun.SQL-Doppelkodierung (Prod-Incident 2026-06-11): // insertSubsequentEventRow band stringifyJson(payload) an ::jsonb — Bun // kodiert einen JS-String für jsonb erneut, das Ergebnis ist ein jsonb- // STRING-Skalar statt einem Objekt. Der typed Read-Pfad parsed Strings // beim Laden zurück (bun-db/query.ts), deshalb blieb das in allen // loadAggregate-Tests unsichtbar — nur SQL-Konsumenten (payload->>'x', // GDPR-Pipeline, MSP-Replays über raw rows) sahen kaputte Daten. Darum // prüft dieser Test das Spalten-Encoding direkt in SQL. test("first AND subsequent append store payload/metadata as jsonb objects", async () => { const aggregateId = uuid(); const base = { aggregateId, aggregateType: "task", tenantId: tenantA, metadata: { userId: userA }, }; await append(testDb.db, { ...base, expectedVersion: 0, type: "task.created", payload: { title: "T" }, }); await append(testDb.db, { ...base, expectedVersion: 1, type: "task.updated", payload: { title: "T2", nested: { deep: true } }, }); const rows = (await asRawClient(testDb.db).unsafe( `SELECT version, jsonb_typeof(payload) AS payload_type, jsonb_typeof(metadata) AS metadata_type, payload->>'title' AS title FROM kumiko_events WHERE aggregate_id = $1 ORDER BY version`, [aggregateId], )) as ReadonlyArray<{ version: number; payload_type: string; metadata_type: string; title: string | null; }>; expect(rows).toHaveLength(2); for (const row of rows) { expect(row.payload_type).toBe("object"); expect(row.metadata_type).toBe("object"); } // SQL-seitiger Feldzugriff funktioniert nur auf echten Objekten — // genau der Pfad, der mit String-Skalaren null lieferte. expect(rows[1]?.title).toBe("T2"); }); });