diff --git a/src/perf/reactor-spans.test.ts b/src/perf/reactor-spans.test.ts new file mode 100644 index 000000000..1b1a980f0 --- /dev/null +++ b/src/perf/reactor-spans.test.ts @@ -0,0 +1,324 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import type { ReactorEmittedEvent } from "@intx/inference"; +import { clear, snapshot, type PerfSpan } from "./index.js"; +import { createPerfReactorObserver } from "./reactor-spans.js"; +import { createTurnContextCollector } from "../session/hooks.js"; + +afterEach(() => { + clear(); +}); + +function event(type: string, data: unknown = {}): ReactorEmittedEvent { + return { type, seq: 1, data } as ReactorEmittedEvent; +} + +function byName(spans: PerfSpan[], name: string): PerfSpan[] { + return spans.filter((s) => s.name === name); +} + +function completed(spans: PerfSpan[]): PerfSpan[] { + return spans.filter((s) => s.endNs !== undefined); +} + +const emptyUsage = { input: 10, output: 5, cacheRead: 0, cacheWrite: 0, thinking: 0 }; +const source = { provider: "test-provider", model: "test-model" }; + +function inferenceDone(content: unknown[] = [{ type: "text", text: "hi" }]): ReactorEmittedEvent { + return event("inference.done", { + turn: { role: "assistant", content, model: "test-model", timestamp: 0 }, + usage: emptyUsage, + source, + }); +} + +describe("createPerfReactorObserver", () => { + test("turn nests inference with ttft and stream when deltas exist", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(event("inference.text.delta", { token: "Hello", partial: { text: "Hello" } })); + obs.observe(event("inference.text.delta", { token: " world", partial: { text: "Hello world" } })); + obs.observe(inferenceDone()); + + const spans = completed(snapshot()); + const turns = byName(spans, "turn"); + const inferences = byName(spans, "inference"); + const ttfts = byName(spans, "inference.ttft"); + const streams = byName(spans, "inference.stream"); + + expect(turns).toHaveLength(1); + expect(inferences).toHaveLength(1); + expect(ttfts).toHaveLength(1); + expect(streams).toHaveLength(1); + + const turn = turns[0]!; + const inference = inferences[0]!; + const ttft = ttfts[0]!; + const stream = streams[0]!; + + expect(turn.parentId).toBeUndefined(); + expect(inference.parentId).toBe(turn.id); + expect(ttft.parentId).toBe(inference.id); + expect(stream.parentId).toBe(inference.id); + + // Ordering: ttft ends at/before stream starts; stream ends at/before inference ends. + expect(ttft.endNs! <= stream.startNs).toBe(true); + expect(stream.endNs! <= inference.endNs!).toBe(true); + expect(inference.endNs! <= turn.endNs!).toBe(true); + }); + + test("tool spans nest under turn after inference.done with tool_calls", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(event("inference.text.delta", { token: "x", partial: { text: "x" } })); + obs.observe( + inferenceDone([ + { type: "tool_call", id: "call-1", name: "run_shell", arguments: {} }, + ]), + ); + obs.observe(event("tool.start", { call: { id: "call-1", name: "run_shell", arguments: {} } })); + obs.observe(event("tool.done", { result: { callId: "call-1", content: "ok" } })); + + const spans = completed(snapshot()); + const turn = byName(spans, "turn")[0]!; + const inference = byName(spans, "inference")[0]!; + const tools = byName(spans, "tool"); + + expect(tools).toHaveLength(1); + expect(tools[0]!.parentId).toBe(turn.id); + expect(tools[0]!.tags?.tool_id).toBe("call-1"); + expect(inference.parentId).toBe(turn.id); + expect(turn.endNs).toBeDefined(); + }); + + test("multiple turns produce separate top-level turn spans", () => { + const obs = createPerfReactorObserver(); + + for (let i = 0; i < 2; i += 1) { + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(event("inference.text.delta", { token: "a", partial: { text: "a" } })); + obs.observe(inferenceDone()); + } + + const spans = completed(snapshot()); + const turns = byName(spans, "turn"); + const inferences = byName(spans, "inference"); + + expect(turns).toHaveLength(2); + expect(inferences).toHaveLength(2); + expect(turns.every((t) => t.parentId === undefined)).toBe(true); + expect(inferences[0]!.parentId).toBe(turns[0]!.id); + expect(inferences[1]!.parentId).toBe(turns[1]!.id); + }); + + test("inference without stream deltas has turn + inference only (no stream)", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(inferenceDone()); + + const spans = completed(snapshot()); + expect(byName(spans, "turn")).toHaveLength(1); + expect(byName(spans, "inference")).toHaveLength(1); + // TTFT still closes at done when no first-token event arrived. + expect(byName(spans, "inference.ttft")).toHaveLength(1); + expect(byName(spans, "inference.stream")).toHaveLength(0); + }); + + test("thinking.delta counts as first token for TTFT", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(event("inference.thinking.delta", { token: "hmm", partial: { text: "" } })); + obs.observe(inferenceDone()); + + const spans = completed(snapshot()); + expect(byName(spans, "inference.ttft")).toHaveLength(1); + expect(byName(spans, "inference.stream")).toHaveLength(1); + }); + + test("blocked tool.done without tool.start still records a tool span", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe( + inferenceDone([ + { type: "tool_call", id: "blocked-1", name: "run_shell", arguments: {} }, + ]), + ); + obs.observe( + event("tool.done", { + result: { callId: "blocked-1", content: "blocked", isError: true }, + }), + ); + + const tools = byName(completed(snapshot()), "tool"); + expect(tools).toHaveLength(1); + expect(tools[0]!.tags?.tool_id).toBe("blocked-1"); + expect(byName(completed(snapshot()), "turn")[0]!.endNs).toBeDefined(); + }); + + test("reset closes open spans and clears state", () => { + const obs = createPerfReactorObserver(); + obs.observe(event("inference.start", { model: "test-model" })); + expect(snapshot().some((s) => s.endNs === undefined)).toBe(true); + + obs.reset(); + + const spans = snapshot(); + expect(spans.every((s) => s.endNs !== undefined)).toBe(true); + expect(byName(spans, "turn")).toHaveLength(1); + + // Next turn is independent. + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(inferenceDone()); + expect(byName(completed(snapshot()), "turn")).toHaveLength(2); + }); + + test("abandon mid-inference then new start closes prior turn with no orphans", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(event("inference.text.delta", { token: "partial", partial: { text: "partial" } })); + // Interrupt: no inference.done / error — next start must abandon. + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(inferenceDone()); + + const spans = snapshot(); + expect(spans.every((s) => s.endNs !== undefined)).toBe(true); + + const turns = byName(completed(spans), "turn"); + const inferences = byName(completed(spans), "inference"); + expect(turns).toHaveLength(2); + expect(inferences).toHaveLength(2); + expect(inferences[0]!.parentId).toBe(turns[0]!.id); + expect(inferences[1]!.parentId).toBe(turns[1]!.id); + // First turn abandoned before second opened — not nested. + expect(turns[0]!.endNs! <= turns[1]!.startNs).toBe(true); + }); + + test("abandon mid-tool then new start closes open tools and prior turn", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe( + inferenceDone([ + { type: "tool_call", id: "call-1", name: "run_shell", arguments: {} }, + ]), + ); + obs.observe(event("tool.start", { call: { id: "call-1", name: "run_shell", arguments: {} } })); + // Interrupt mid-tool: no tool.done — next inference.start must not nest. + obs.observe(event("inference.start", { model: "next-model" })); + obs.observe(inferenceDone()); + + const spans = snapshot(); + expect(spans.every((s) => s.endNs !== undefined)).toBe(true); + + const turns = byName(completed(spans), "turn"); + const tools = byName(completed(spans), "tool"); + const inferences = byName(completed(spans), "inference"); + + expect(turns).toHaveLength(2); + expect(tools).toHaveLength(1); + expect(tools[0]!.parentId).toBe(turns[0]!.id); + expect(inferences).toHaveLength(2); + expect(inferences[0]!.parentId).toBe(turns[0]!.id); + expect(inferences[1]!.parentId).toBe(turns[1]!.id); + // Second inference must not nest under the abandoned turn. + expect(inferences[1]!.parentId).not.toBe(turns[0]!.id); + expect(turns[0]!.endNs! <= turns[1]!.startNs).toBe(true); + }); + + test("inference.error mid-turn then new start leaves no open spans", () => { + const obs = createPerfReactorObserver(); + + obs.observe(event("inference.start", { model: "test-model" })); + obs.observe(event("inference.text.delta", { token: "x", partial: { text: "x" } })); + obs.observe(event("inference.error", { error: { message: "timeout" } })); + + // Error with no pending tools closes the turn. + expect(snapshot().every((s) => s.endNs !== undefined)).toBe(true); + + obs.observe(event("inference.start", { model: "retry-model" })); + obs.observe(inferenceDone()); + + const spans = snapshot(); + expect(spans.every((s) => s.endNs !== undefined)).toBe(true); + const turns = byName(completed(spans), "turn"); + expect(turns).toHaveLength(2); + expect(byName(completed(spans), "inference")[1]!.parentId).toBe(turns[1]!.id); + }); +}); + +describe("turn collector durationMs unchanged with perf observer", () => { + test("coarse durationMs still reported for a completed turn", () => { + // createTurnContextCollector stamps cycleStartedAt on construction, then + // inference.start re-stamps it, then completePending reads finish. + const times = [1_000, 1_000, 1_250]; + let i = 0; + const now = (): number => times[Math.min(i++, times.length - 1)]!; + + const completedTurns: { durationMs: number }[] = []; + const collector = createTurnContextCollector((ctx) => { + completedTurns.push({ durationMs: ctx.durationMs }); + }, now); + const obs = createPerfReactorObserver(); + + // Mirror the sink: both observers see the same events. + const feed = (e: ReactorEmittedEvent): void => { + collector.observe(e); + obs.observe(e); + }; + + feed(event("inference.start", { model: "test-model" })); + feed(event("inference.text.delta", { token: "hi", partial: { text: "hi" } })); + feed(inferenceDone()); + + expect(completedTurns).toHaveLength(1); + expect(completedTurns[0]!.durationMs).toBe(250); + expect(collector.getTurnCount()).toBe(1); + + // Perf spans still present and nested. + const spans = completed(snapshot()); + expect(byName(spans, "turn")).toHaveLength(1); + expect(byName(spans, "inference")).toHaveLength(1); + }); + + test("durationMs with tools waits until tool.done", () => { + // Construction stamps cycleStartedAt, inference.start re-stamps, tool.done + // completePending reads finish. + const times = [5_000, 5_000, 5_400]; + let i = 0; + const now = (): number => times[Math.min(i++, times.length - 1)]!; + + const completedTurns: { durationMs: number }[] = []; + const collector = createTurnContextCollector((ctx) => { + completedTurns.push({ durationMs: ctx.durationMs }); + }, now); + const obs = createPerfReactorObserver(); + + const feed = (e: ReactorEmittedEvent): void => { + collector.observe(e); + obs.observe(e); + }; + + feed(event("inference.start", { model: "m" })); + feed( + inferenceDone([ + { type: "tool_call", id: "c1", name: "read_file", arguments: {} }, + ]), + ); + expect(completedTurns).toHaveLength(0); + + feed(event("tool.start", { call: { id: "c1", name: "read_file", arguments: {} } })); + feed(event("tool.done", { result: { callId: "c1", content: "ok" } })); + + expect(completedTurns).toHaveLength(1); + expect(completedTurns[0]!.durationMs).toBe(400); + + const spans = completed(snapshot()); + expect(byName(spans, "tool")).toHaveLength(1); + expect(byName(spans, "tool")[0]!.parentId).toBe(byName(spans, "turn")[0]!.id); + }); +}); diff --git a/src/perf/reactor-spans.ts b/src/perf/reactor-spans.ts new file mode 100644 index 000000000..fcefd0d3f --- /dev/null +++ b/src/perf/reactor-spans.ts @@ -0,0 +1,261 @@ +/** + * Map reactor stream events onto the shared PerfTrace span tree. + * + * Always-on, local-only. Nesting: + * turn + * inference + * inference.ttft (start → first content-bearing delta) + * inference.stream (first delta → inference.done) + * tool (per invocation) + */ + +import type { ReactorEmittedEvent } from "@intx/inference"; +import { end, start } from "./index.js"; + +/** Content-bearing events that end TTFT and open the stream phase. */ +const FIRST_TOKEN_TYPES: ReadonlySet = new Set([ + "inference.text.delta", + "inference.thinking.delta", + "inference.refusal.delta", + "inference.tool_call.start", + "inference.tool_call.delta", + "inference.tool_call.end", + "inference.citation", + "inference.image_output", + "inference.code_execution.start", + "inference.code_execution.delta", + "inference.code_execution.result", + "inference.thinking.signature", + "inference.thinking.redacted", +]); + +export type PerfReactorObserver = { + observe(event: ReactorEmittedEvent): void; + reset(): void; +}; + +type ObserverState = { + turnId: string | null; + inferenceId: string | null; + ttftId: string | null; + streamId: string | null; + /** Tool calls still expected before the current turn can close. */ + pendingTools: number; + /** Open tool spans keyed by callId. */ + openTools: Map; +}; + +function emptyState(): ObserverState { + return { + turnId: null, + inferenceId: null, + ttftId: null, + streamId: null, + pendingTools: 0, + openTools: new Map(), + }; +} + +function toolCallCount(event: ReactorEmittedEvent): number { + if (event.type !== "inference.done") return 0; + const data = event.data as { + turn?: { content?: ReadonlyArray<{ type: string }> }; + }; + const content = data.turn?.content; + if (content === undefined) return 0; + let n = 0; + for (const block of content) { + if (block.type === "tool_call") n += 1; + } + return n; +} + +function callIdFromToolStart(event: ReactorEmittedEvent): string | undefined { + if (event.type !== "tool.start") return undefined; + const data = event.data as { call?: { id?: unknown } }; + const id = data.call?.id; + return typeof id === "string" && id.length > 0 ? id : undefined; +} + +function callIdFromToolDone(event: ReactorEmittedEvent): string | undefined { + if (event.type !== "tool.done") return undefined; + const data = event.data as { result?: { callId?: unknown } }; + const id = data.result?.callId; + return typeof id === "string" && id.length > 0 ? id : undefined; +} + +function modelTags(event: ReactorEmittedEvent): Record | undefined { + if (event.type === "inference.start") { + const data = event.data as { model?: unknown }; + if (typeof data.model === "string" && data.model.length > 0) { + return { model_id: data.model }; + } + return undefined; + } + if (event.type === "inference.done") { + const data = event.data as { + source?: { provider?: unknown; model?: unknown }; + usage?: { input?: unknown; output?: unknown }; + }; + const tags: Record = {}; + if (typeof data.source?.provider === "string") tags.provider_id = data.source.provider; + if (typeof data.source?.model === "string") tags.model_id = data.source.model; + if (typeof data.usage?.input === "number") tags.input_tokens = data.usage.input; + if (typeof data.usage?.output === "number") tags.output_tokens = data.usage.output; + return Object.keys(tags).length > 0 ? tags : undefined; + } + return undefined; +} + +/** + * Observe reactor events and open/close nested PerfTrace spans. + * One observer per run-sink; call `reset` when the sink resets. + */ +export function createPerfReactorObserver(): PerfReactorObserver { + let state = emptyState(); + + function endIfOpen(id: string | null, tags?: Record): void { + if (id === null || id.length === 0) return; + end(id, tags); + } + + function closeInferenceTree(tags?: Record): void { + endIfOpen(state.streamId); + state.streamId = null; + // No first-token event: fold TTFT into the full inference window. + endIfOpen(state.ttftId); + state.ttftId = null; + endIfOpen(state.inferenceId, tags); + state.inferenceId = null; + } + + function closeOpenTools(): void { + for (const spanId of state.openTools.values()) { + endIfOpen(spanId); + } + state.openTools.clear(); + } + + /** + * Single exit for ending a turn: close orphan tool spans, then the turn. + * Inference tree must already be closed (or will be via abandonTurn). + */ + function closeTurn(): void { + closeOpenTools(); + endIfOpen(state.turnId); + state.turnId = null; + state.pendingTools = 0; + } + + /** Abandon the whole open tree (inference + tools + turn). */ + function abandonTurn(): void { + closeInferenceTree(); + closeTurn(); + } + + function ensureTurn(): string { + if (state.turnId === null) { + state.turnId = start("turn"); + } + return state.turnId; + } + + function onFirstToken(): void { + if (state.inferenceId === null || state.streamId !== null) return; + + endIfOpen(state.ttftId); + state.ttftId = null; + state.streamId = start("inference.stream", { parentId: state.inferenceId }); + } + + function observe(event: ReactorEmittedEvent): void { + const type = event.type; + + if (type === "inference.start") { + // Always abandon any prior turn before opening a new one. Interrupt mid- + // inference or mid-tool must not nest the next call under a stale turn or + // leave orphan tool spans in the process-wide open map. + abandonTurn(); + const turnId = ensureTurn(); + const tags = modelTags(event); + state.inferenceId = start("inference", { + parentId: turnId, + ...(tags !== undefined ? { tags } : {}), + }); + state.ttftId = start("inference.ttft", { parentId: state.inferenceId }); + state.streamId = null; + return; + } + + if (FIRST_TOKEN_TYPES.has(type)) { + onFirstToken(); + return; + } + + if (type === "inference.done") { + const tags = modelTags(event); + closeInferenceTree(tags); + state.pendingTools = toolCallCount(event); + if (state.pendingTools === 0) { + closeTurn(); + } + return; + } + + if (type === "inference.error") { + closeInferenceTree(); + // Drop the turn if nothing is waiting on tools; otherwise keep it open + // so in-flight tool spans can still close under it. + if (state.pendingTools === 0) { + closeTurn(); + } + return; + } + + if (type === "tool.start") { + const turnId = state.turnId; + if (turnId === null) return; + const callId = callIdFromToolStart(event); + if (callId === undefined) return; + if (state.openTools.has(callId)) return; + const spanId = start("tool", { + parentId: turnId, + tags: { tool_id: callId }, + }); + state.openTools.set(callId, spanId); + return; + } + + if (type === "tool.done") { + const callId = callIdFromToolDone(event); + if (callId !== undefined) { + const openId = state.openTools.get(callId); + if (openId !== undefined) { + end(openId); + state.openTools.delete(callId); + } else if (state.turnId !== null) { + // Blocked tools emit tool.done without tool.start. + const spanId = start("tool", { + parentId: state.turnId, + tags: { tool_id: callId }, + }); + end(spanId); + } + } + if (state.pendingTools > 0) { + state.pendingTools -= 1; + } + if (state.pendingTools === 0 && state.inferenceId === null && state.turnId !== null) { + closeTurn(); + } + return; + } + } + + function reset(): void { + abandonTurn(); + state = emptyState(); + } + + return { observe, reset }; +} diff --git a/src/session/run-sink.ts b/src/session/run-sink.ts index 170101c7c..89295ce12 100644 --- a/src/session/run-sink.ts +++ b/src/session/run-sink.ts @@ -1,6 +1,7 @@ import type { EventEmitter } from "node:events"; import type { ReactorEmittedEvent } from "@intx/inference"; import type { TokenUsage } from "@intx/types/runtime"; +import { createPerfReactorObserver } from "../perf/reactor-spans.js"; import { createTurnContextCollector, type LifecycleHookManager, @@ -89,9 +90,12 @@ export function createRunSink(args: RunSinkArgs): RunSink { let runCompleted = false; let runError: string | undefined; let turnCollector = createCollector(); + // Always-on local PerfTrace: not gated by lifecycle hooks. + let perfObserver = createPerfReactorObserver(); const sink = (event: ReactorEmittedEvent): void => { turnCollector.observe(event); + perfObserver.observe(event); if (event.type === "reactor.done") { runCompleted = true; // Terminal success clears any earlier transient inference error. @@ -126,6 +130,7 @@ export function createRunSink(args: RunSinkArgs): RunSink { runCompleted = false; runError = undefined; turnCollector = createCollector(); + perfObserver.reset(); }, }; }