SPB Git

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%
9.8 KB · 254 lines typescript
Raw Blame History
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