spb/khaelor Public
KHAELOR — a terminal-native autonomous engineering agent powered by Anthropic.
TypeScript 82.9%
HTML 14.9%
CSS 1.1%
JavaScript 0.7%
1/**2 * KHAELOR3 * File: tests/session/bus.test.ts4 * Description: Event bus tests — typed/wildcard subscription, write-ahead ordering, coalescing contract.5 *6 * Author: Simon-Pierre Boucher7 * Contact: contact@spboucher.ai8 */910import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";11import { Coalescer, SessionEventBus, blockKey, COALESCER_TRUNCATION_MARKER } from "../../src/session/bus.js";12import type { CoalescedFrame, DurableAppender } from "../../src/session/bus.js";13import { EVENT_SCHEMA_VERSION } from "../../src/session/events.js";14import type { DurableEvent, DurableEventInput, KhaelorEvent } from "../../src/session/events.js";15import { SESSION_ID } from "./fixtures.js";1617function makeBus(onHandlerError?: (error: unknown, event: KhaelorEvent) => void) {18 return new SessionEventBus(19 onHandlerError ? { sessionId: SESSION_ID, onHandlerError } : { sessionId: SESSION_ID },20 );21}2223describe("SessionEventBus", () => {24 it("delivers typed subscriptions only for the matching type", () => {25 const bus = makeBus();26 const renamed: string[] = [];27 const other: string[] = [];28 bus.on("session.renamed", (e) => renamed.push(e.payload.title));29 bus.on("user.message-created", (e) => other.push(e.payload.text));30 bus.publishDurable({ type: "session.renamed", payload: { title: "T1" } });31 expect(renamed).toEqual(["T1"]);32 expect(other).toEqual([]);33 });3435 it("wildcard subscribers see durable and ephemeral events", () => {36 const bus = makeBus();37 const seen: string[] = [];38 bus.onAny((e) => seen.push(e.type));39 bus.publishDurable({ type: "session.renamed", payload: { title: "x" } });40 bus.publishEphemeral({41 type: "model.text-delta",42 payload: { requestId: "r1", blockIndex: 0, text: "hi" },43 });44 expect(seen).toEqual(["session.renamed", "model.text-delta"]);45 });4647 it("assigns envelope fields (v, id, sessionId, seq, ts) with gapless seq", () => {48 const bus = makeBus();49 const a = bus.publishDurable({ type: "session.renamed", payload: { title: "a" } });50 const b = bus.publishDurable({ type: "session.renamed", payload: { title: "b" } });51 expect(a.v).toBe(EVENT_SCHEMA_VERSION);52 expect(a.sessionId).toBe(SESSION_ID);53 expect(a.seq).toBe(1);54 expect(b.seq).toBe(2);55 expect(a.id).not.toBe(b.id);56 });5758 it("write-ahead: the appender records the event before subscribers see it", () => {59 const appended: DurableEvent[] = [];60 let seq = 0;61 const appender: DurableAppender = {62 append(input: DurableEventInput): DurableEvent {63 const event = {64 v: EVENT_SCHEMA_VERSION,65 id: `id-${seq + 1}`,66 sessionId: SESSION_ID,67 seq: ++seq,68 ts: 1,69 type: input.type,70 payload: input.payload,71 } as DurableEvent;72 appended.push(event);73 return event;74 },75 };76 const bus = new SessionEventBus({ sessionId: SESSION_ID, appender });77 let appendedAtDelivery = -1;78 bus.on("session.renamed", () => {79 appendedAtDelivery = appended.length;80 });81 const event = bus.publishDurable({ type: "session.renamed", payload: { title: "x" } });82 expect(appendedAtDelivery).toBe(1); // already appended when the handler ran83 expect(appended[0]).toBe(event);84 });8586 it("rejects ephemeral types on publishDurable and vice versa", () => {87 const bus = makeBus();88 expect(() =>89 bus.publishDurable({ type: "model.text-delta", payload: {} } as never),90 ).toThrow(/not a durable/);91 expect(() =>92 bus.publishEphemeral({ type: "session.renamed", payload: {} } as never),93 ).toThrow(/not an ephemeral/);94 });9596 it("a throwing handler never blocks other subscribers", () => {97 const errors: unknown[] = [];98 const bus = makeBus((error) => errors.push(error));99 const seen: string[] = [];100 bus.on("session.renamed", () => {101 throw new Error("bad handler");102 });103 bus.on("session.renamed", (e) => seen.push(e.payload.title));104 bus.publishDurable({ type: "session.renamed", payload: { title: "still delivered" } });105 expect(seen).toEqual(["still delivered"]);106 expect(errors).toHaveLength(1);107 });108109 it("unsubscribe detaches the handler (leak check)", () => {110 const bus = makeBus();111 const seen: string[] = [];112 const un1 = bus.on("session.renamed", (e) => seen.push(e.payload.title));113 const un2 = bus.onAny(() => seen.push("any"));114 expect(bus.subscriberCount()).toBe(2);115 un1();116 un2();117 expect(bus.subscriberCount()).toBe(0);118 bus.publishDurable({ type: "session.renamed", payload: { title: "x" } });119 expect(seen).toEqual([]);120 });121});122123describe("Coalescer (~16 ms contract)", () => {124 beforeEach(() => {125 vi.useFakeTimers();126 });127 afterEach(() => {128 vi.useRealTimers();129 });130131 function setup(opts?: { windowMs?: number; maxAppendPerKey?: number }) {132 const bus = makeBus();133 const coalescer = new Coalescer(bus, opts ?? {});134 const frames: CoalescedFrame[] = [];135 coalescer.subscribe((f) => frames.push(f));136 return { bus, coalescer, frames };137 }138139 const textDelta = (text: string, blockIndex = 0) =>140 ({141 type: "model.text-delta",142 payload: { requestId: "r1", blockIndex, text },143 }) as const;144145 it("a delta storm produces at most one frame per window, with concatenated appends", () => {146 const { bus, frames } = setup();147 for (let i = 0; i < 100; i++) bus.publishEphemeral(textDelta(`${i};`));148 expect(frames).toHaveLength(0); // nothing until the window elapses149 vi.advanceTimersByTime(16);150 expect(frames).toHaveLength(1);151 const appended = frames[0]!.textAppends.get(blockKey("r1", 0));152 expect(appended).toBe(Array.from({ length: 100 }, (_, i) => `${i};`).join(""));153 vi.advanceTimersByTime(160);154 expect(frames).toHaveLength(1); // no idle wake-ups, no empty frames155 });156157 it("keeps separate keys separate (blocks, tools, processes)", () => {158 const { bus, frames } = setup();159 bus.publishEphemeral(textDelta("a", 0));160 bus.publishEphemeral(textDelta("b", 1));161 bus.publishEphemeral({ type: "tool.output", payload: { toolUseId: "t1", chunk: "out1" } });162 bus.publishEphemeral({163 type: "process.output",164 payload: { processId: "p1", stream: "stdout", chunk: "proc1" },165 });166 vi.advanceTimersByTime(16);167 const frame = frames[0]!;168 expect(frame.textAppends.get(blockKey("r1", 0))).toBe("a");169 expect(frame.textAppends.get(blockKey("r1", 1))).toBe("b");170 expect(frame.toolOutputAppends.get("t1")).toBe("out1");171 expect(frame.processOutputAppends.get("p1")).toBe("proc1");172 });173174 it("settlement flushes immediately: buffered deltas arrive in the same frame, before the durable", () => {175 const { bus, frames } = setup();176 bus.publishEphemeral(textDelta("Hel"));177 bus.publishEphemeral(textDelta("lo"));178 expect(frames).toHaveLength(0);179 const settled = bus.publishDurable({180 type: "model.text-block-completed",181 payload: { requestId: "r1", blockIndex: 0, text: "Hello" },182 });183 // Immediate flush — no timer needed.184 expect(frames).toHaveLength(1);185 const frame = frames[0]!;186 expect(frame.textAppends.get(blockKey("r1", 0))).toBe("Hello".slice(0, 3) + "lo");187 expect(frame.durables).toEqual([settled]);188 });189190 it("tool input previews accumulate across frames until tool.requested settles them", () => {191 const { bus, frames } = setup();192 bus.publishEphemeral({193 type: "model.tool-input-delta",194 payload: { requestId: "r1", blockIndex: 0, toolUseId: "t1", partialJson: '{"pa' },195 });196 vi.advanceTimersByTime(16);197 expect(frames[0]!.toolInputPreviews.get("t1")).toBe('{"pa');198 bus.publishEphemeral({199 type: "model.tool-input-delta",200 payload: { requestId: "r1", blockIndex: 0, toolUseId: "t1", partialJson: 'th": "x"}' },201 });202 vi.advanceTimersByTime(16);203 // Latest ACCUMULATED preview, not just this window's chunk.204 expect(frames[1]!.toolInputPreviews.get("t1")).toBe('{"path": "x"}');205 bus.publishDurable({206 type: "tool.requested",207 payload: { requestId: "r1", blockIndex: 0, toolUseId: "t1", toolName: "read", input: { path: "x" } },208 });209 vi.advanceTimersByTime(16);210 expect(frames).toHaveLength(3);211 expect(frames[2]!.toolInputPreviews.has("t1")).toBe(false); // settlement replaces212 });213214 it("Interrupted forces an immediate flush", () => {215 const { bus, frames } = setup();216 bus.publishEphemeral(textDelta("partial"));217 bus.publishDurable({218 type: "user.interrupted",219 payload: { scope: "turn", pendingToolUseIds: [] },220 });221 expect(frames).toHaveLength(1);222 expect(frames[0]!.textAppends.get(blockKey("r1", 0))).toBe("partial");223 expect(frames[0]!.durables[0]!.type).toBe("user.interrupted");224 });225226 it("non-settling durable events are delivered on the window cadence", () => {227 const { bus, frames } = setup();228 const event = bus.publishDurable({ type: "session.renamed", payload: { title: "x" } });229 expect(frames).toHaveLength(0);230 vi.advanceTimersByTime(16);231 expect(frames).toHaveLength(1);232 expect(frames[0]!.durables).toEqual([event]);233 });234235 it("caps per-key appends within a window and marks the truncation", () => {236 const { bus, frames } = setup({ maxAppendPerKey: 8 });237 bus.publishEphemeral(textDelta("12345678"));238 bus.publishEphemeral(textDelta("OVERFLOW"));239 vi.advanceTimersByTime(16);240 const appended = frames[0]!.textAppends.get(blockKey("r1", 0))!;241 expect(appended.startsWith("12345678")).toBe(true);242 expect(appended.endsWith(COALESCER_TRUNCATION_MARKER)).toBe(true);243 expect(appended.length).toBe(8 + COALESCER_TRUNCATION_MARKER.length);244 });245246 it("dispose detaches from the bus and stops delivering frames", () => {247 const { bus, coalescer, frames } = setup();248 bus.publishEphemeral(textDelta("x"));249 coalescer.dispose();250 vi.advanceTimersByTime(160);251 expect(frames).toHaveLength(0);252 });253});254