/** * KHAELOR * File: tests/session/bus.test.ts * Description: Event bus tests — typed/wildcard subscription, write-ahead ordering, coalescing contract. * * Author: Simon-Pierre Boucher * Contact: contact@spboucher.ai */ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { Coalescer, SessionEventBus, blockKey, COALESCER_TRUNCATION_MARKER } from "../../src/session/bus.js"; import type { CoalescedFrame, DurableAppender } from "../../src/session/bus.js"; import { EVENT_SCHEMA_VERSION } from "../../src/session/events.js"; import type { DurableEvent, DurableEventInput, KhaelorEvent } from "../../src/session/events.js"; import { SESSION_ID } from "./fixtures.js"; function makeBus(onHandlerError?: (error: unknown, event: KhaelorEvent) => void) { return new SessionEventBus( onHandlerError ? { sessionId: SESSION_ID, onHandlerError } : { sessionId: SESSION_ID }, ); } describe("SessionEventBus", () => { it("delivers typed subscriptions only for the matching type", () => { const bus = makeBus(); const renamed: string[] = []; const other: string[] = []; bus.on("session.renamed", (e) => renamed.push(e.payload.title)); bus.on("user.message-created", (e) => other.push(e.payload.text)); bus.publishDurable({ type: "session.renamed", payload: { title: "T1" } }); expect(renamed).toEqual(["T1"]); expect(other).toEqual([]); }); it("wildcard subscribers see durable and ephemeral events", () => { const bus = makeBus(); const seen: string[] = []; bus.onAny((e) => seen.push(e.type)); bus.publishDurable({ type: "session.renamed", payload: { title: "x" } }); bus.publishEphemeral({ type: "model.text-delta", payload: { requestId: "r1", blockIndex: 0, text: "hi" }, }); expect(seen).toEqual(["session.renamed", "model.text-delta"]); }); it("assigns envelope fields (v, id, sessionId, seq, ts) with gapless seq", () => { const bus = makeBus(); const a = bus.publishDurable({ type: "session.renamed", payload: { title: "a" } }); const b = bus.publishDurable({ type: "session.renamed", payload: { title: "b" } }); expect(a.v).toBe(EVENT_SCHEMA_VERSION); expect(a.sessionId).toBe(SESSION_ID); expect(a.seq).toBe(1); expect(b.seq).toBe(2); expect(a.id).not.toBe(b.id); }); it("write-ahead: the appender records the event before subscribers see it", () => { const appended: DurableEvent[] = []; let seq = 0; const appender: DurableAppender = { append(input: DurableEventInput): DurableEvent { const event = { v: EVENT_SCHEMA_VERSION, id: `id-${seq + 1}`, sessionId: SESSION_ID, seq: ++seq, ts: 1, type: input.type, payload: input.payload, } as DurableEvent; appended.push(event); return event; }, }; const bus = new SessionEventBus({ sessionId: SESSION_ID, appender }); let appendedAtDelivery = -1; bus.on("session.renamed", () => { appendedAtDelivery = appended.length; }); const event = bus.publishDurable({ type: "session.renamed", payload: { title: "x" } }); expect(appendedAtDelivery).toBe(1); // already appended when the handler ran expect(appended[0]).toBe(event); }); it("rejects ephemeral types on publishDurable and vice versa", () => { const bus = makeBus(); expect(() => bus.publishDurable({ type: "model.text-delta", payload: {} } as never), ).toThrow(/not a durable/); expect(() => bus.publishEphemeral({ type: "session.renamed", payload: {} } as never), ).toThrow(/not an ephemeral/); }); it("a throwing handler never blocks other subscribers", () => { const errors: unknown[] = []; const bus = makeBus((error) => errors.push(error)); const seen: string[] = []; bus.on("session.renamed", () => { throw new Error("bad handler"); }); bus.on("session.renamed", (e) => seen.push(e.payload.title)); bus.publishDurable({ type: "session.renamed", payload: { title: "still delivered" } }); expect(seen).toEqual(["still delivered"]); expect(errors).toHaveLength(1); }); it("unsubscribe detaches the handler (leak check)", () => { const bus = makeBus(); const seen: string[] = []; const un1 = bus.on("session.renamed", (e) => seen.push(e.payload.title)); const un2 = bus.onAny(() => seen.push("any")); expect(bus.subscriberCount()).toBe(2); un1(); un2(); expect(bus.subscriberCount()).toBe(0); bus.publishDurable({ type: "session.renamed", payload: { title: "x" } }); expect(seen).toEqual([]); }); }); describe("Coalescer (~16 ms contract)", () => { beforeEach(() => { vi.useFakeTimers(); }); afterEach(() => { vi.useRealTimers(); }); function setup(opts?: { windowMs?: number; maxAppendPerKey?: number }) { const bus = makeBus(); const coalescer = new Coalescer(bus, opts ?? {}); const frames: CoalescedFrame[] = []; coalescer.subscribe((f) => frames.push(f)); return { bus, coalescer, frames }; } const textDelta = (text: string, blockIndex = 0) => ({ type: "model.text-delta", payload: { requestId: "r1", blockIndex, text }, }) as const; it("a delta storm produces at most one frame per window, with concatenated appends", () => { const { bus, frames } = setup(); for (let i = 0; i < 100; i++) bus.publishEphemeral(textDelta(`${i};`)); expect(frames).toHaveLength(0); // nothing until the window elapses vi.advanceTimersByTime(16); expect(frames).toHaveLength(1); const appended = frames[0]!.textAppends.get(blockKey("r1", 0)); expect(appended).toBe(Array.from({ length: 100 }, (_, i) => `${i};`).join("")); vi.advanceTimersByTime(160); expect(frames).toHaveLength(1); // no idle wake-ups, no empty frames }); it("keeps separate keys separate (blocks, tools, processes)", () => { const { bus, frames } = setup(); bus.publishEphemeral(textDelta("a", 0)); bus.publishEphemeral(textDelta("b", 1)); bus.publishEphemeral({ type: "tool.output", payload: { toolUseId: "t1", chunk: "out1" } }); bus.publishEphemeral({ type: "process.output", payload: { processId: "p1", stream: "stdout", chunk: "proc1" }, }); vi.advanceTimersByTime(16); const frame = frames[0]!; expect(frame.textAppends.get(blockKey("r1", 0))).toBe("a"); expect(frame.textAppends.get(blockKey("r1", 1))).toBe("b"); expect(frame.toolOutputAppends.get("t1")).toBe("out1"); expect(frame.processOutputAppends.get("p1")).toBe("proc1"); }); it("settlement flushes immediately: buffered deltas arrive in the same frame, before the durable", () => { const { bus, frames } = setup(); bus.publishEphemeral(textDelta("Hel")); bus.publishEphemeral(textDelta("lo")); expect(frames).toHaveLength(0); const settled = bus.publishDurable({ type: "model.text-block-completed", payload: { requestId: "r1", blockIndex: 0, text: "Hello" }, }); // Immediate flush — no timer needed. expect(frames).toHaveLength(1); const frame = frames[0]!; expect(frame.textAppends.get(blockKey("r1", 0))).toBe("Hello".slice(0, 3) + "lo"); expect(frame.durables).toEqual([settled]); }); it("tool input previews accumulate across frames until tool.requested settles them", () => { const { bus, frames } = setup(); bus.publishEphemeral({ type: "model.tool-input-delta", payload: { requestId: "r1", blockIndex: 0, toolUseId: "t1", partialJson: '{"pa' }, }); vi.advanceTimersByTime(16); expect(frames[0]!.toolInputPreviews.get("t1")).toBe('{"pa'); bus.publishEphemeral({ type: "model.tool-input-delta", payload: { requestId: "r1", blockIndex: 0, toolUseId: "t1", partialJson: 'th": "x"}' }, }); vi.advanceTimersByTime(16); // Latest ACCUMULATED preview, not just this window's chunk. expect(frames[1]!.toolInputPreviews.get("t1")).toBe('{"path": "x"}'); bus.publishDurable({ type: "tool.requested", payload: { requestId: "r1", blockIndex: 0, toolUseId: "t1", toolName: "read", input: { path: "x" } }, }); vi.advanceTimersByTime(16); expect(frames).toHaveLength(3); expect(frames[2]!.toolInputPreviews.has("t1")).toBe(false); // settlement replaces }); it("Interrupted forces an immediate flush", () => { const { bus, frames } = setup(); bus.publishEphemeral(textDelta("partial")); bus.publishDurable({ type: "user.interrupted", payload: { scope: "turn", pendingToolUseIds: [] }, }); expect(frames).toHaveLength(1); expect(frames[0]!.textAppends.get(blockKey("r1", 0))).toBe("partial"); expect(frames[0]!.durables[0]!.type).toBe("user.interrupted"); }); it("non-settling durable events are delivered on the window cadence", () => { const { bus, frames } = setup(); const event = bus.publishDurable({ type: "session.renamed", payload: { title: "x" } }); expect(frames).toHaveLength(0); vi.advanceTimersByTime(16); expect(frames).toHaveLength(1); expect(frames[0]!.durables).toEqual([event]); }); it("caps per-key appends within a window and marks the truncation", () => { const { bus, frames } = setup({ maxAppendPerKey: 8 }); bus.publishEphemeral(textDelta("12345678")); bus.publishEphemeral(textDelta("OVERFLOW")); vi.advanceTimersByTime(16); const appended = frames[0]!.textAppends.get(blockKey("r1", 0))!; expect(appended.startsWith("12345678")).toBe(true); expect(appended.endsWith(COALESCER_TRUNCATION_MARKER)).toBe(true); expect(appended.length).toBe(8 + COALESCER_TRUNCATION_MARKER.length); }); it("dispose detaches from the bus and stops delivering frames", () => { const { bus, coalescer, frames } = setup(); bus.publishEphemeral(textDelta("x")); coalescer.dispose(); vi.advanceTimersByTime(160); expect(frames).toHaveLength(0); }); });