Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions apps/cli/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -24,17 +24,17 @@
"@cyrus/utils": "workspace:*",
"@ff-labs/fff-node": "^0.9.6",
"@orpc/server": "catalog:rpc",
"@soorya-u/better-auth-ws-ticket": "catalog:auth",
"@t3-oss/env-core": "catalog:env",
"@tursodatabase/database": "^0.6.1",
"@soorya-u/better-auth-ws-ticket": "catalog:auth",
"async-mutex": "^0.5.0",
"better-auth": "catalog:auth",
"better-result": "catalog:core",
"commander": "^15.0.0",
"diff": "^7.0.0",
"es-git": "^0.7.0",
"extract-zip": "^2.0.1",
"evlog": "catalog:observability",
"extract-zip": "^2.0.1",
"node-datachannel": "^0.32.3",
"parse-diff": "^0.11.1",
"proper-lockfile": "^4.1.2",
Expand Down
97 changes: 97 additions & 0 deletions apps/cli/src/queue/bus.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
import { describe, expect, test } from "bun:test";
import type { ChatChunk } from "@cyrus/schemas/rtc/chat";
import { createThreadEventBus } from "./bus";

function persistedChunk(seq: number, turnId = "turn-1"): ChatChunk {
return {
threadId: "thread-1",
turnId,
seq,
event: { type: "message_completed", text: "done", messageId: "m1" },
};
}

function tokenChunk(turnId = "turn-1"): ChatChunk {
return {
threadId: "thread-1",
turnId,
seq: 0,
event: { type: "token", text: "hi", messageId: "m1" },
};
}

function createReader(iterator: AsyncGenerator<ChatChunk>) {
let pending = iterator.next();
return {
async next(): Promise<ChatChunk> {
const result = await pending;
pending = iterator.next();
if (result.done) throw new Error("iterator ended unexpectedly");
return result.value;
},
};
}

describe("thread event bus", () => {
test("stamps ephemeral chunks with the last known persisted anchor and an increasing sub", async () => {
const bus = createThreadEventBus();
bus.watch("peer-1", "thread-1");
const reader = createReader(bus.subscribe("peer-1"));

bus.publish(persistedChunk(5));
bus.publish(tokenChunk());
bus.publish(tokenChunk());

const first = await reader.next();
expect(first.seq).toBe(5);
expect(first.sub).toBeUndefined();

expect(await reader.next()).toMatchObject({ seq: 5, sub: 1 });
expect(await reader.next()).toMatchObject({ seq: 5, sub: 2 });
});

test("does not replay ephemeral chunks on reconnect, but does replay persisted ones", async () => {
const bus = createThreadEventBus();
bus.watch("peer-1", "thread-1");
const first = createReader(bus.subscribe("peer-1"));

bus.publish(persistedChunk(5));
bus.publish(tokenChunk());
await first.next();
await first.next();

// Simulate a reconnect: same peerId subscribes again.
const second = createReader(bus.subscribe("peer-1"));
const replayed = await second.next();
expect(replayed.seq).toBe(5);
expect(replayed.sub).toBeUndefined();

bus.publish(tokenChunk());
const live = await second.next();
// The sub counter keeps advancing across a reconnect (it's per-anchor,
// not per-connection) — the first ephemeral chunk before the reconnect
// already claimed sub 1.
expect(live).toMatchObject({ sub: 2 });
});

test("advances the anchor and resets sub once new content persists", async () => {
const bus = createThreadEventBus();
bus.watch("peer-1", "thread-1");
const reader = createReader(bus.subscribe("peer-1"));

bus.publish(persistedChunk(5));
bus.publish(tokenChunk());
bus.publish(persistedChunk(6));
bus.publish(tokenChunk());

const a = await reader.next();
expect(a.seq).toBe(5);
expect(a.sub).toBeUndefined();
expect(await reader.next()).toMatchObject({ seq: 5, sub: 1 });

const b = await reader.next();
expect(b.seq).toBe(6);
expect(b.sub).toBeUndefined();
expect(await reader.next()).toMatchObject({ seq: 6, sub: 1 });
});
});
52 changes: 42 additions & 10 deletions apps/cli/src/queue/bus.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,15 @@ function isTerminalEvent(event: ChatChunk["event"]): boolean {
return event.type === "turn_completed" || event.type === "turn_interrupted";
}

type ThreadCursor = {
anchor: number;
nextSub: number;
};

function isUnstampedPlaceholder(chunk: ChatChunk): boolean {
return chunk.seq === 0;
}

export function createThreadEventBus(
options: CreateThreadEventBusOptions = {}
): ThreadEventBus {
Expand All @@ -26,8 +35,33 @@ export function createThreadEventBus(
const watchedThreads = new Map<string, Set<string>>();
const activeTurnLogs = new Map<string, ChatChunk[]>();
const turnThreads = new Map<string, string>();
const threadCursors = new Map<string, ThreadCursor>();
let closed = false;

function stampChunk(chunk: ChatChunk): ChatChunk {
let cursor = threadCursors.get(chunk.threadId);
if (!cursor) {
cursor = { anchor: 0, nextSub: 1 };
threadCursors.set(chunk.threadId, cursor);
}

if (!isUnstampedPlaceholder(chunk)) {
if (chunk.seq > cursor.anchor) {
cursor.anchor = chunk.seq;
cursor.nextSub = 1;
}
return chunk;
}

const stamped: ChatChunk = {
...chunk,
seq: cursor.anchor,
sub: cursor.nextSub,
};
cursor.nextSub += 1;
return stamped;
}

function getWatchedThreads(peerId: string): Set<string> {
let set = watchedThreads.get(peerId);
if (!set) {
Expand Down Expand Up @@ -62,11 +96,6 @@ export function createThreadEventBus(
}

function trimTurnLog(log: ChatChunk[]): void {
while (log.length > maxChunksPerTurn) {
const deltaIndex = log.findIndex((chunk) => chunk.seq === 0);
if (deltaIndex === -1) break;
log.splice(deltaIndex, 1);
}
while (log.length > maxChunksPerTurn) log.shift();
}

Expand Down Expand Up @@ -105,13 +134,15 @@ export function createThreadEventBus(
publish(chunk) {
if (closed) return;

const terminal = isTerminalEvent(chunk.event);
if (!terminal) {
appendToTurnLog(chunk);
const stamped = stampChunk(chunk);
const terminal = isTerminalEvent(stamped.event);

if (!terminal && stamped.sub === undefined) {
appendToTurnLog(stamped);
}
fanOut(chunk);
fanOut(stamped);
if (terminal) {
evictTurnLog(chunk.turnId);
evictTurnLog(stamped.turnId);
}
},

Expand Down Expand Up @@ -186,6 +217,7 @@ export function createThreadEventBus(
watchedThreads.clear();
activeTurnLogs.clear();
turnThreads.clear();
threadCursors.clear();
},
};
}
Loading