feat: harden workflow gates and expose token usage

This commit is contained in:
2026-07-21 12:14:26 +02:00
parent 8aa1676811
commit 2ff63d371f
20 changed files with 604 additions and 53 deletions
+64 -1
View File
@@ -1,7 +1,7 @@
import { renderHook } from "@testing-library/react";
import { act, useLayoutEffect } from "react";
import { FakeEventSource } from "../test/fakeEventSource";
import { useSessionStream } from "./useSessionStream";
import { createStreamEventCoalescer, useSessionStream } from "./useSessionStream";
import { useSessionStore } from "../store/sessionStore";
beforeEach(() => {
@@ -10,6 +10,49 @@ beforeEach(() => {
useSessionStore.getState().resetSession();
});
test("coalesces a burst of text deltas into bounded store updates", () => {
vi.useFakeTimers();
try {
const applied: Array<{ type: string; text?: string }> = [];
const coalescer = createStreamEventCoalescer((event) => applied.push(event), 100);
for (let i = 0; i < 100; i += 1) {
coalescer.push({ type: "text_delta", text: "x" });
}
expect(applied).toEqual([{ type: "text_delta", text: "x" }]);
vi.advanceTimersByTime(100);
expect(applied).toEqual([
{ type: "text_delta", text: "x" },
{ type: "text_delta", text: "x".repeat(99) },
]);
coalescer.dispose();
} finally {
vi.useRealTimers();
}
});
test("flushes pending stream text before a structural event", () => {
vi.useFakeTimers();
try {
const applied: Array<{ type: string; text?: string }> = [];
const coalescer = createStreamEventCoalescer((event) => applied.push(event), 100);
coalescer.push({ type: "text_delta", text: "a" });
coalescer.push({ type: "text_delta", text: "b" });
coalescer.push({ type: "system_event", event: "turn_end" });
expect(applied).toEqual([
{ type: "text_delta", text: "a" },
{ type: "text_delta", text: "b" },
{ type: "system_event", event: "turn_end" },
]);
coalescer.dispose();
} finally {
vi.useRealTimers();
}
});
test("opens an EventSource and feeds NAMED events to the store", () => {
renderHook(() => useSessionStream("s1"));
const es = FakeEventSource.instances[0];
@@ -56,6 +99,26 @@ test("feeds named activity_event events to the tool activity timeline", () => {
});
});
test("feeds named usage events to the session token counters", () => {
renderHook(() => useSessionStream("s1"));
const es = FakeEventSource.instances[0];
act(() =>
es.emitNamed("usage", {
type: "usage",
usage: { input: 1_000, cacheRead: 2_000, output: 300, totalTokens: 3_300, contextWindow: 200_000 },
})
);
expect(useSessionStore.getState().tokenUsage).toEqual({
input: 1_000,
cacheRead: 2_000,
output: 300,
totalTokens: 3_300,
contextWindow: 200_000,
});
});
test("also feeds UNNAMED (default message) events to the store", () => {
renderHook(() => useSessionStream("s1"));
const es = FakeEventSource.instances[0];
+66 -1
View File
@@ -4,6 +4,68 @@ import { joinBackendPath } from "../api/runtime-config";
import { useSessionStore } from "../store/sessionStore";
import type { StreamEvent } from "../api/types";
const STREAM_UPDATE_INTERVAL_MS = 100;
function isStreamingDelta(
event: StreamEvent,
): event is Extract<StreamEvent, { type: "text_delta" | "activity_delta" }> {
return event.type === "text_delta" || event.type === "activity_delta";
}
export function createStreamEventCoalescer(
applyEvent: (event: StreamEvent) => void,
intervalMs = STREAM_UPDATE_INTERVAL_MS,
) {
let pending: Extract<StreamEvent, { type: "text_delta" | "activity_delta" }> | null = null;
let timer: ReturnType<typeof setTimeout> | null = null;
let disposed = false;
const flushPending = () => {
if (!pending) return;
const event = pending;
pending = null;
applyEvent(event);
};
const schedule = () => {
timer = setTimeout(() => {
timer = null;
if (!pending || disposed) return;
flushPending();
schedule();
}, intervalMs);
};
return {
push(event: StreamEvent) {
if (disposed) return;
if (!isStreamingDelta(event)) {
flushPending();
if (timer) clearTimeout(timer);
timer = null;
applyEvent(event);
return;
}
if (!timer) {
applyEvent(event);
schedule();
return;
}
if (pending && pending.type !== event.type) flushPending();
pending = pending
? { ...event, text: pending.text + event.text }
: event;
},
dispose() {
if (disposed) return;
if (timer) clearTimeout(timer);
timer = null;
flushPending();
disposed = true;
},
};
}
export function useSessionStream(
sessionId: string | null,
generation = 0,
@@ -36,6 +98,7 @@ export function useSessionStream(
: "";
const es = new EventSource(joinBackendPath(BASE, `/sessions/${sessionId}/events${query}`));
const identity = { source: es, sessionId, cursorResetEpoch };
const coalescer = createStreamEventCoalescer(applyEvent);
activeSource.current = identity;
es.onopen = () => {
if (activeSource.current === identity) setConnected(true);
@@ -48,7 +111,7 @@ export function useSessionStream(
if (activeSource.current !== identity) return;
if (ev.lastEventId) cursor.current.lastEventId = ev.lastEventId;
try {
applyEvent(JSON.parse(ev.data) as StreamEvent);
coalescer.push(JSON.parse(ev.data) as StreamEvent);
} catch {
/* ignore malformed */
}
@@ -62,6 +125,7 @@ export function useSessionStream(
"text_delta",
"activity_delta",
"activity_event",
"usage",
"info",
"system_event",
] as const;
@@ -71,6 +135,7 @@ export function useSessionStream(
for (const name of namedEvents) es.removeEventListener(name, handle);
es.onmessage = null;
es.close();
coalescer.dispose();
if (activeSource.current === identity) {
activeSource.current = null;
setConnected(false);