Skip to content

Commit 02c9207

Browse files
Merge pull request #305 from corbitsdev/stack/w2a-cl-5171-auto-spans
Automatic core spans from reactor events (CL-5171)
2 parents faa36e7 + 073d6de commit 02c9207

3 files changed

Lines changed: 590 additions & 0 deletions

File tree

src/perf/reactor-spans.test.ts

Lines changed: 324 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,324 @@
1+
import { afterEach, describe, expect, test } from "bun:test";
2+
import type { ReactorEmittedEvent } from "@intx/inference";
3+
import { clear, snapshot, type PerfSpan } from "./index.js";
4+
import { createPerfReactorObserver } from "./reactor-spans.js";
5+
import { createTurnContextCollector } from "../session/hooks.js";
6+
7+
afterEach(() => {
8+
clear();
9+
});
10+
11+
function event(type: string, data: unknown = {}): ReactorEmittedEvent {
12+
return { type, seq: 1, data } as ReactorEmittedEvent;
13+
}
14+
15+
function byName(spans: PerfSpan[], name: string): PerfSpan[] {
16+
return spans.filter((s) => s.name === name);
17+
}
18+
19+
function completed(spans: PerfSpan[]): PerfSpan[] {
20+
return spans.filter((s) => s.endNs !== undefined);
21+
}
22+
23+
const emptyUsage = { input: 10, output: 5, cacheRead: 0, cacheWrite: 0, thinking: 0 };
24+
const source = { provider: "test-provider", model: "test-model" };
25+
26+
function inferenceDone(content: unknown[] = [{ type: "text", text: "hi" }]): ReactorEmittedEvent {
27+
return event("inference.done", {
28+
turn: { role: "assistant", content, model: "test-model", timestamp: 0 },
29+
usage: emptyUsage,
30+
source,
31+
});
32+
}
33+
34+
describe("createPerfReactorObserver", () => {
35+
test("turn nests inference with ttft and stream when deltas exist", () => {
36+
const obs = createPerfReactorObserver();
37+
38+
obs.observe(event("inference.start", { model: "test-model" }));
39+
obs.observe(event("inference.text.delta", { token: "Hello", partial: { text: "Hello" } }));
40+
obs.observe(event("inference.text.delta", { token: " world", partial: { text: "Hello world" } }));
41+
obs.observe(inferenceDone());
42+
43+
const spans = completed(snapshot());
44+
const turns = byName(spans, "turn");
45+
const inferences = byName(spans, "inference");
46+
const ttfts = byName(spans, "inference.ttft");
47+
const streams = byName(spans, "inference.stream");
48+
49+
expect(turns).toHaveLength(1);
50+
expect(inferences).toHaveLength(1);
51+
expect(ttfts).toHaveLength(1);
52+
expect(streams).toHaveLength(1);
53+
54+
const turn = turns[0]!;
55+
const inference = inferences[0]!;
56+
const ttft = ttfts[0]!;
57+
const stream = streams[0]!;
58+
59+
expect(turn.parentId).toBeUndefined();
60+
expect(inference.parentId).toBe(turn.id);
61+
expect(ttft.parentId).toBe(inference.id);
62+
expect(stream.parentId).toBe(inference.id);
63+
64+
// Ordering: ttft ends at/before stream starts; stream ends at/before inference ends.
65+
expect(ttft.endNs! <= stream.startNs).toBe(true);
66+
expect(stream.endNs! <= inference.endNs!).toBe(true);
67+
expect(inference.endNs! <= turn.endNs!).toBe(true);
68+
});
69+
70+
test("tool spans nest under turn after inference.done with tool_calls", () => {
71+
const obs = createPerfReactorObserver();
72+
73+
obs.observe(event("inference.start", { model: "test-model" }));
74+
obs.observe(event("inference.text.delta", { token: "x", partial: { text: "x" } }));
75+
obs.observe(
76+
inferenceDone([
77+
{ type: "tool_call", id: "call-1", name: "run_shell", arguments: {} },
78+
]),
79+
);
80+
obs.observe(event("tool.start", { call: { id: "call-1", name: "run_shell", arguments: {} } }));
81+
obs.observe(event("tool.done", { result: { callId: "call-1", content: "ok" } }));
82+
83+
const spans = completed(snapshot());
84+
const turn = byName(spans, "turn")[0]!;
85+
const inference = byName(spans, "inference")[0]!;
86+
const tools = byName(spans, "tool");
87+
88+
expect(tools).toHaveLength(1);
89+
expect(tools[0]!.parentId).toBe(turn.id);
90+
expect(tools[0]!.tags?.tool_id).toBe("call-1");
91+
expect(inference.parentId).toBe(turn.id);
92+
expect(turn.endNs).toBeDefined();
93+
});
94+
95+
test("multiple turns produce separate top-level turn spans", () => {
96+
const obs = createPerfReactorObserver();
97+
98+
for (let i = 0; i < 2; i += 1) {
99+
obs.observe(event("inference.start", { model: "test-model" }));
100+
obs.observe(event("inference.text.delta", { token: "a", partial: { text: "a" } }));
101+
obs.observe(inferenceDone());
102+
}
103+
104+
const spans = completed(snapshot());
105+
const turns = byName(spans, "turn");
106+
const inferences = byName(spans, "inference");
107+
108+
expect(turns).toHaveLength(2);
109+
expect(inferences).toHaveLength(2);
110+
expect(turns.every((t) => t.parentId === undefined)).toBe(true);
111+
expect(inferences[0]!.parentId).toBe(turns[0]!.id);
112+
expect(inferences[1]!.parentId).toBe(turns[1]!.id);
113+
});
114+
115+
test("inference without stream deltas has turn + inference only (no stream)", () => {
116+
const obs = createPerfReactorObserver();
117+
118+
obs.observe(event("inference.start", { model: "test-model" }));
119+
obs.observe(inferenceDone());
120+
121+
const spans = completed(snapshot());
122+
expect(byName(spans, "turn")).toHaveLength(1);
123+
expect(byName(spans, "inference")).toHaveLength(1);
124+
// TTFT still closes at done when no first-token event arrived.
125+
expect(byName(spans, "inference.ttft")).toHaveLength(1);
126+
expect(byName(spans, "inference.stream")).toHaveLength(0);
127+
});
128+
129+
test("thinking.delta counts as first token for TTFT", () => {
130+
const obs = createPerfReactorObserver();
131+
132+
obs.observe(event("inference.start", { model: "test-model" }));
133+
obs.observe(event("inference.thinking.delta", { token: "hmm", partial: { text: "" } }));
134+
obs.observe(inferenceDone());
135+
136+
const spans = completed(snapshot());
137+
expect(byName(spans, "inference.ttft")).toHaveLength(1);
138+
expect(byName(spans, "inference.stream")).toHaveLength(1);
139+
});
140+
141+
test("blocked tool.done without tool.start still records a tool span", () => {
142+
const obs = createPerfReactorObserver();
143+
144+
obs.observe(event("inference.start", { model: "test-model" }));
145+
obs.observe(
146+
inferenceDone([
147+
{ type: "tool_call", id: "blocked-1", name: "run_shell", arguments: {} },
148+
]),
149+
);
150+
obs.observe(
151+
event("tool.done", {
152+
result: { callId: "blocked-1", content: "blocked", isError: true },
153+
}),
154+
);
155+
156+
const tools = byName(completed(snapshot()), "tool");
157+
expect(tools).toHaveLength(1);
158+
expect(tools[0]!.tags?.tool_id).toBe("blocked-1");
159+
expect(byName(completed(snapshot()), "turn")[0]!.endNs).toBeDefined();
160+
});
161+
162+
test("reset closes open spans and clears state", () => {
163+
const obs = createPerfReactorObserver();
164+
obs.observe(event("inference.start", { model: "test-model" }));
165+
expect(snapshot().some((s) => s.endNs === undefined)).toBe(true);
166+
167+
obs.reset();
168+
169+
const spans = snapshot();
170+
expect(spans.every((s) => s.endNs !== undefined)).toBe(true);
171+
expect(byName(spans, "turn")).toHaveLength(1);
172+
173+
// Next turn is independent.
174+
obs.observe(event("inference.start", { model: "test-model" }));
175+
obs.observe(inferenceDone());
176+
expect(byName(completed(snapshot()), "turn")).toHaveLength(2);
177+
});
178+
179+
test("abandon mid-inference then new start closes prior turn with no orphans", () => {
180+
const obs = createPerfReactorObserver();
181+
182+
obs.observe(event("inference.start", { model: "test-model" }));
183+
obs.observe(event("inference.text.delta", { token: "partial", partial: { text: "partial" } }));
184+
// Interrupt: no inference.done / error — next start must abandon.
185+
obs.observe(event("inference.start", { model: "test-model" }));
186+
obs.observe(inferenceDone());
187+
188+
const spans = snapshot();
189+
expect(spans.every((s) => s.endNs !== undefined)).toBe(true);
190+
191+
const turns = byName(completed(spans), "turn");
192+
const inferences = byName(completed(spans), "inference");
193+
expect(turns).toHaveLength(2);
194+
expect(inferences).toHaveLength(2);
195+
expect(inferences[0]!.parentId).toBe(turns[0]!.id);
196+
expect(inferences[1]!.parentId).toBe(turns[1]!.id);
197+
// First turn abandoned before second opened — not nested.
198+
expect(turns[0]!.endNs! <= turns[1]!.startNs).toBe(true);
199+
});
200+
201+
test("abandon mid-tool then new start closes open tools and prior turn", () => {
202+
const obs = createPerfReactorObserver();
203+
204+
obs.observe(event("inference.start", { model: "test-model" }));
205+
obs.observe(
206+
inferenceDone([
207+
{ type: "tool_call", id: "call-1", name: "run_shell", arguments: {} },
208+
]),
209+
);
210+
obs.observe(event("tool.start", { call: { id: "call-1", name: "run_shell", arguments: {} } }));
211+
// Interrupt mid-tool: no tool.done — next inference.start must not nest.
212+
obs.observe(event("inference.start", { model: "next-model" }));
213+
obs.observe(inferenceDone());
214+
215+
const spans = snapshot();
216+
expect(spans.every((s) => s.endNs !== undefined)).toBe(true);
217+
218+
const turns = byName(completed(spans), "turn");
219+
const tools = byName(completed(spans), "tool");
220+
const inferences = byName(completed(spans), "inference");
221+
222+
expect(turns).toHaveLength(2);
223+
expect(tools).toHaveLength(1);
224+
expect(tools[0]!.parentId).toBe(turns[0]!.id);
225+
expect(inferences).toHaveLength(2);
226+
expect(inferences[0]!.parentId).toBe(turns[0]!.id);
227+
expect(inferences[1]!.parentId).toBe(turns[1]!.id);
228+
// Second inference must not nest under the abandoned turn.
229+
expect(inferences[1]!.parentId).not.toBe(turns[0]!.id);
230+
expect(turns[0]!.endNs! <= turns[1]!.startNs).toBe(true);
231+
});
232+
233+
test("inference.error mid-turn then new start leaves no open spans", () => {
234+
const obs = createPerfReactorObserver();
235+
236+
obs.observe(event("inference.start", { model: "test-model" }));
237+
obs.observe(event("inference.text.delta", { token: "x", partial: { text: "x" } }));
238+
obs.observe(event("inference.error", { error: { message: "timeout" } }));
239+
240+
// Error with no pending tools closes the turn.
241+
expect(snapshot().every((s) => s.endNs !== undefined)).toBe(true);
242+
243+
obs.observe(event("inference.start", { model: "retry-model" }));
244+
obs.observe(inferenceDone());
245+
246+
const spans = snapshot();
247+
expect(spans.every((s) => s.endNs !== undefined)).toBe(true);
248+
const turns = byName(completed(spans), "turn");
249+
expect(turns).toHaveLength(2);
250+
expect(byName(completed(spans), "inference")[1]!.parentId).toBe(turns[1]!.id);
251+
});
252+
});
253+
254+
describe("turn collector durationMs unchanged with perf observer", () => {
255+
test("coarse durationMs still reported for a completed turn", () => {
256+
// createTurnContextCollector stamps cycleStartedAt on construction, then
257+
// inference.start re-stamps it, then completePending reads finish.
258+
const times = [1_000, 1_000, 1_250];
259+
let i = 0;
260+
const now = (): number => times[Math.min(i++, times.length - 1)]!;
261+
262+
const completedTurns: { durationMs: number }[] = [];
263+
const collector = createTurnContextCollector((ctx) => {
264+
completedTurns.push({ durationMs: ctx.durationMs });
265+
}, now);
266+
const obs = createPerfReactorObserver();
267+
268+
// Mirror the sink: both observers see the same events.
269+
const feed = (e: ReactorEmittedEvent): void => {
270+
collector.observe(e);
271+
obs.observe(e);
272+
};
273+
274+
feed(event("inference.start", { model: "test-model" }));
275+
feed(event("inference.text.delta", { token: "hi", partial: { text: "hi" } }));
276+
feed(inferenceDone());
277+
278+
expect(completedTurns).toHaveLength(1);
279+
expect(completedTurns[0]!.durationMs).toBe(250);
280+
expect(collector.getTurnCount()).toBe(1);
281+
282+
// Perf spans still present and nested.
283+
const spans = completed(snapshot());
284+
expect(byName(spans, "turn")).toHaveLength(1);
285+
expect(byName(spans, "inference")).toHaveLength(1);
286+
});
287+
288+
test("durationMs with tools waits until tool.done", () => {
289+
// Construction stamps cycleStartedAt, inference.start re-stamps, tool.done
290+
// completePending reads finish.
291+
const times = [5_000, 5_000, 5_400];
292+
let i = 0;
293+
const now = (): number => times[Math.min(i++, times.length - 1)]!;
294+
295+
const completedTurns: { durationMs: number }[] = [];
296+
const collector = createTurnContextCollector((ctx) => {
297+
completedTurns.push({ durationMs: ctx.durationMs });
298+
}, now);
299+
const obs = createPerfReactorObserver();
300+
301+
const feed = (e: ReactorEmittedEvent): void => {
302+
collector.observe(e);
303+
obs.observe(e);
304+
};
305+
306+
feed(event("inference.start", { model: "m" }));
307+
feed(
308+
inferenceDone([
309+
{ type: "tool_call", id: "c1", name: "read_file", arguments: {} },
310+
]),
311+
);
312+
expect(completedTurns).toHaveLength(0);
313+
314+
feed(event("tool.start", { call: { id: "c1", name: "read_file", arguments: {} } }));
315+
feed(event("tool.done", { result: { callId: "c1", content: "ok" } }));
316+
317+
expect(completedTurns).toHaveLength(1);
318+
expect(completedTurns[0]!.durationMs).toBe(400);
319+
320+
const spans = completed(snapshot());
321+
expect(byName(spans, "tool")).toHaveLength(1);
322+
expect(byName(spans, "tool")[0]!.parentId).toBe(byName(spans, "turn")[0]!.id);
323+
});
324+
});

0 commit comments

Comments
 (0)