Skip to content

Commit 978ca10

Browse files
committed
Wire automatic PerfTrace spans from reactor events
Maps reactor events to nested perf spans: turn contains inference, which contains inference.ttft and inference.stream when stream deltas arrive; tool spans nest under turn. Wired always-on in createRunSink alongside the existing turn collector so coarse durationMs is unchanged.
1 parent 73158b2 commit 978ca10

3 files changed

Lines changed: 518 additions & 0 deletions

File tree

src/perf/reactor-spans.test.ts

Lines changed: 250 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,250 @@
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+
180+
describe("turn collector durationMs unchanged with perf observer", () => {
181+
test("coarse durationMs still reported for a completed turn", () => {
182+
// createTurnContextCollector stamps cycleStartedAt on construction, then
183+
// inference.start re-stamps it, then completePending reads finish.
184+
const times = [1_000, 1_000, 1_250];
185+
let i = 0;
186+
const now = (): number => times[Math.min(i++, times.length - 1)]!;
187+
188+
const completedTurns: { durationMs: number }[] = [];
189+
const collector = createTurnContextCollector((ctx) => {
190+
completedTurns.push({ durationMs: ctx.durationMs });
191+
}, now);
192+
const obs = createPerfReactorObserver();
193+
194+
// Mirror the sink: both observers see the same events.
195+
const feed = (e: ReactorEmittedEvent): void => {
196+
collector.observe(e);
197+
obs.observe(e);
198+
};
199+
200+
feed(event("inference.start", { model: "test-model" }));
201+
feed(event("inference.text.delta", { token: "hi", partial: { text: "hi" } }));
202+
feed(inferenceDone());
203+
204+
expect(completedTurns).toHaveLength(1);
205+
expect(completedTurns[0]!.durationMs).toBe(250);
206+
expect(collector.getTurnCount()).toBe(1);
207+
208+
// Perf spans still present and nested.
209+
const spans = completed(snapshot());
210+
expect(byName(spans, "turn")).toHaveLength(1);
211+
expect(byName(spans, "inference")).toHaveLength(1);
212+
});
213+
214+
test("durationMs with tools waits until tool.done", () => {
215+
// Construction stamps cycleStartedAt, inference.start re-stamps, tool.done
216+
// completePending reads finish.
217+
const times = [5_000, 5_000, 5_400];
218+
let i = 0;
219+
const now = (): number => times[Math.min(i++, times.length - 1)]!;
220+
221+
const completedTurns: { durationMs: number }[] = [];
222+
const collector = createTurnContextCollector((ctx) => {
223+
completedTurns.push({ durationMs: ctx.durationMs });
224+
}, now);
225+
const obs = createPerfReactorObserver();
226+
227+
const feed = (e: ReactorEmittedEvent): void => {
228+
collector.observe(e);
229+
obs.observe(e);
230+
};
231+
232+
feed(event("inference.start", { model: "m" }));
233+
feed(
234+
inferenceDone([
235+
{ type: "tool_call", id: "c1", name: "read_file", arguments: {} },
236+
]),
237+
);
238+
expect(completedTurns).toHaveLength(0);
239+
240+
feed(event("tool.start", { call: { id: "c1", name: "read_file", arguments: {} } }));
241+
feed(event("tool.done", { result: { callId: "c1", content: "ok" } }));
242+
243+
expect(completedTurns).toHaveLength(1);
244+
expect(completedTurns[0]!.durationMs).toBe(400);
245+
246+
const spans = completed(snapshot());
247+
expect(byName(spans, "tool")).toHaveLength(1);
248+
expect(byName(spans, "tool")[0]!.parentId).toBe(byName(spans, "turn")[0]!.id);
249+
});
250+
});

0 commit comments

Comments
 (0)