Skip to content

Commit 47fa183

Browse files
Merge pull request #694 from corbitsdev/cl-6816-cut-posthog-event-volume-aggregate-tool-spans-into-the-turn
Cut PostHog event volume with generation aggregates and subagent rollups
2 parents 1df7821 + 3ec8c44 commit 47fa183

24 files changed

Lines changed: 1579 additions & 169 deletions

docs/TELEMETRY.md

Lines changed: 85 additions & 38 deletions
Large diffs are not rendered by default.

src/exec/runner.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,7 @@ import type { ApprovalOutcome, PermissionRequest } from "../permission/types.js"
6262
import { createAgentToolset, type AgentToolset, type OperatorResult } from "../agent/tools.js";
6363
import { createAgentWithLiveToolDispatch } from "../agent/live-tool-dispatch.js";
6464
import { liveTelemetry } from "../telemetry/singleton.js";
65+
import { createTurnObserver } from "../telemetry/ai-observability.js";
6566
import { collectToolPlugins, resolveToolPlugins } from "../plugins/tool-plugins.js";
6667
import {
6768
expandExistingPluginMembers,
@@ -657,9 +658,15 @@ export async function runExec(config: Config): Promise<ExecResult> {
657658
const hookManager = createLifecycleHookManager({
658659
hooks: await discoverLifecycleHooks(hookDirectories(config.cwd)),
659660
});
661+
const turnObserver = createTurnObserver({
662+
telemetry: () => liveTelemetry,
663+
getSessionId: () => sessionId,
664+
getSource: () => liveSource,
665+
});
660666
const liveSink = createRunSink({
661667
emitter,
662668
hookManager,
669+
...turnObserver,
663670
onTurnBoundarySnapshot: () => {
664671
void persist("running");
665672
},

src/plugins/loader.ts

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,10 @@ import type { CommandPlugin } from "../tui/commands/registry.js";
1010
import { pathIsInsideOrEqual } from "../util/path-contain.js";
1111
import { parsePluginManifest, type PluginManifest } from "./manifest.js";
1212
import { NOOP_TELEMETRY, type Telemetry } from "../telemetry/index.js";
13+
import type { PluginLoadReporter } from "../telemetry/product-events.js";
14+
import { runtimePluginLoadReporter } from "../telemetry/singleton.js";
1315
import { loadDataOnlyPlugin } from "./data-only.js";
16+
1417
import {
1518
resolvePluginWarningHandler,
1619
stderrPluginWarning,
@@ -125,6 +128,7 @@ export async function loadPluginEntry(
125128
diagnostics?: PluginLoadDiagnostics;
126129
origin?: PluginOrigin;
127130
telemetry?: Telemetry;
131+
pluginLoadReporter?: PluginLoadReporter;
128132
} = {},
129133
): Promise<PluginModule | null> {
130134
const cwd = opts.cwd ?? process.cwd();
@@ -139,6 +143,7 @@ export async function loadPluginEntry(
139143
);
140144
const origin = opts.origin;
141145
const telemetry = opts.telemetry ?? NOOP_TELEMETRY;
146+
const reportPluginLoaded = opts.pluginLoadReporter ?? runtimePluginLoadReporter;
142147
let target = entryPath;
143148
let pluginDir = entryPath;
144149
try {
@@ -173,7 +178,9 @@ export async function loadPluginEntry(
173178
mod.origin = origin;
174179
mod.pluginPath = resolve(entryPath);
175180
}
176-
if (origin !== undefined) telemetry.capture("plugin_loaded", { origin });
181+
if (origin !== undefined) {
182+
reportPluginLoaded(telemetry, origin, resolve(entryPath));
183+
}
177184
return mod;
178185
}
179186
return null;
@@ -234,7 +241,9 @@ export async function loadPluginEntry(
234241
result.origin = origin;
235242
result.pluginPath = resolve(pluginDir);
236243
}
237-
if (origin !== undefined) telemetry.capture("plugin_loaded", { origin });
244+
if (origin !== undefined) {
245+
reportPluginLoaded(telemetry, origin, resolve(pluginDir));
246+
}
238247
return result;
239248
} catch (err) {
240249
// Route through the same sink as skill/load warnings so a diagnostics
@@ -848,7 +857,9 @@ export async function discoverClaudeInstalledPlugins(
848857
name: plugin.manifest.name === plugin.manifest.id ? idFromKey : plugin.manifest.name,
849858
};
850859
}
851-
opts.telemetry?.capture("plugin_loaded", { origin: "user" });
860+
if (opts.telemetry !== undefined) {
861+
runtimePluginLoadReporter(opts.telemetry, "user", resolve(d));
862+
}
852863
results.push(plugin);
853864
}
854865
}

src/session/run-sink.test.ts

Lines changed: 159 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@ import { describe, expect, test } from "bun:test";
33
import type { ReactorEmittedEvent } from "@intx/inference";
44
import { createRunSink } from "./run-sink.js";
55
import type { LifecycleHookStatus } from "./hooks.js";
6+
import { createTurnObserver } from "../telemetry/ai-observability.js";
7+
import type { Telemetry } from "../telemetry/index.js";
68

79
function event(type: string, data: unknown): ReactorEmittedEvent {
810
return { type, seq: 1, data } as ReactorEmittedEvent;
@@ -23,6 +25,42 @@ const enabledHook: LifecycleHookStatus = {
2325
enabled: true,
2426
};
2527

28+
function attributionHarness(selectedSource = { provider: "provider-a", model: "model-a" }) {
29+
const captured: { event: string; properties: Record<string, unknown> }[] = [];
30+
const telemetry: Telemetry = {
31+
enabled: true,
32+
installationId: "test",
33+
capture: (capturedEvent, properties = {}) => {
34+
captured.push({ event: capturedEvent, properties });
35+
},
36+
captureIntentional: () => false,
37+
flush: async () => {},
38+
discard: () => {},
39+
};
40+
const observer = createTurnObserver({
41+
telemetry: () => telemetry,
42+
getSessionId: () => "session-1",
43+
getSource: () => selectedSource,
44+
});
45+
const runSink = createRunSink({
46+
emitter: new EventEmitter(),
47+
hookManager: stubHookManager([]),
48+
...observer,
49+
});
50+
return { captured, runSink };
51+
}
52+
53+
function failMessageRun(runSink: ReturnType<typeof createRunSink>): void {
54+
runSink.sink(event("inference.error", { error: { message: "attempt failed" } }));
55+
runSink.sink(
56+
event("message.run.ended", {
57+
messageRunId: "run-1",
58+
messageId: "message-1",
59+
status: "failed",
60+
}),
61+
);
62+
}
63+
2664
describe("createRunSink", () => {
2765
test("allocates no turn collector when no lifecycle hooks are configured", () => {
2866
const runSink = createRunSink({
@@ -112,7 +150,7 @@ describe("createRunSink", () => {
112150
expect(runSink.getTurnCount()).toBe(1);
113151
});
114152

115-
test("reports the in-flight turn to onTurnFailed when a turn errors instead of completing", () => {
153+
test("settles a pending inference failure only when the message run fails", () => {
116154
const failures: { turnIndex: number; error: string }[] = [];
117155
const runSink = createRunSink({
118156
emitter: new EventEmitter(),
@@ -122,26 +160,121 @@ describe("createRunSink", () => {
122160

123161
runSink.sink(event("inference.start", {}));
124162
runSink.sink(event("inference.error", { error: { message: "429 rate limit" } }));
163+
expect(failures).toEqual([]);
164+
165+
runSink.sink(
166+
event("message.run.ended", {
167+
messageRunId: "run-1",
168+
messageId: "message-1",
169+
status: "failed",
170+
error: { message: "reactor gave up", kind: "inference_error" },
171+
}),
172+
);
125173

126174
expect(failures).toEqual([{ turnIndex: 0, error: "429 rate limit" }]);
127175
});
128176

129-
// Regression: one give-up reaches the sink twice — the director surfaces
130-
// the failed inference, then the reactor terminates the run — and reporting
131-
// both files two failed turns under a single turn's identity.
132-
test("reports one failure per turn across both error paths, not one per error event", () => {
177+
test("uses unknown attribution when a fallback model fails before usage", () => {
178+
const { captured, runSink } = attributionHarness();
179+
180+
runSink.sink(event("inference.start", { model: "model-b" }));
181+
failMessageRun(runSink);
182+
183+
expect(captured).toHaveLength(1);
184+
expect(captured[0]?.event).toBe("$ai_generation");
185+
expect(captured[0]?.properties).toMatchObject({
186+
$ai_provider: "unknown",
187+
$ai_model: "model-b",
188+
$ai_is_error: true,
189+
});
190+
});
191+
192+
test("uses the selected source when its model fails before usage", () => {
193+
const { captured, runSink } = attributionHarness();
194+
195+
runSink.sink(event("inference.start", { model: "model-a" }));
196+
failMessageRun(runSink);
197+
198+
expect(captured).toHaveLength(1);
199+
expect(captured[0]?.properties).toMatchObject({
200+
$ai_provider: "provider-a",
201+
$ai_model: "model-a",
202+
$ai_is_error: true,
203+
});
204+
});
205+
206+
test("uses authoritative usage attribution for a failed fallback", () => {
207+
const { captured, runSink } = attributionHarness();
208+
209+
runSink.sink(event("inference.start", { model: "model-b" }));
210+
runSink.sink(
211+
event("inference.usage", {
212+
usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, thinking: 0 },
213+
source: { sourceId: "fallback", provider: "provider-b", model: "model-b" },
214+
}),
215+
);
216+
failMessageRun(runSink);
217+
218+
expect(captured).toHaveLength(1);
219+
expect(captured[0]?.properties).toMatchObject({
220+
$ai_provider: "provider-b",
221+
$ai_model: "model-b",
222+
$ai_is_error: true,
223+
});
224+
});
225+
226+
test("does not leak authoritative source attribution across retry attempts", () => {
227+
const { captured, runSink } = attributionHarness();
228+
229+
runSink.sink(event("inference.start", { model: "model-a" }));
230+
runSink.sink(
231+
event("inference.usage", {
232+
usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, thinking: 0 },
233+
source: { sourceId: "selected", provider: "provider-a", model: "model-a" },
234+
}),
235+
);
236+
runSink.sink(event("inference.error", { error: { message: "retry" } }));
237+
runSink.sink(event("inference.start", { model: "model-b" }));
238+
failMessageRun(runSink);
239+
240+
expect(captured).toHaveLength(1);
241+
expect(captured[0]?.properties).toMatchObject({
242+
$ai_provider: "unknown",
243+
$ai_model: "model-b",
244+
$ai_is_error: true,
245+
});
246+
});
247+
248+
test("discards a recoverable inference failure after retry success", () => {
133249
const failures: { turnIndex: number; error: string }[] = [];
250+
const completions: number[] = [];
134251
const runSink = createRunSink({
135252
emitter: new EventEmitter(),
136253
hookManager: stubHookManager([]),
254+
onTurnComplete: (ctx) => completions.push(ctx.turnIndex),
137255
onTurnFailed: (info) => failures.push(info),
138256
});
139257

140258
runSink.sink(event("inference.start", {}));
141-
runSink.sink(event("inference.error", { error: { message: "429 rate limit" } }));
142-
runSink.sink(event("reactor.error", { error: "reactor gave up" }));
259+
runSink.sink(event("inference.error", { error: { message: "retry me" } }));
260+
runSink.sink(event("inference.start", {}));
261+
runSink.sink(
262+
event("inference.done", {
263+
turn: { role: "assistant", content: [], model: "test", timestamp: 0 },
264+
usage: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0, thinking: 0 },
265+
source: { provider: "test", model: "test" },
266+
}),
267+
);
268+
runSink.sink(
269+
event("message.run.ended", {
270+
messageRunId: "run-1",
271+
messageId: "message-1",
272+
status: "completed",
273+
}),
274+
);
143275

144-
expect(failures).toEqual([{ turnIndex: 0, error: "429 rate limit" }]);
276+
expect(completions).toEqual([0]);
277+
expect(failures).toEqual([]);
145278
});
146279

147280
test("reports no failure for a turn that already completed", () => {
@@ -168,9 +301,7 @@ describe("createRunSink", () => {
168301
expect(failures).toEqual([]);
169302
});
170303

171-
// A retry re-enters inference.start under the same turn index, so a second
172-
// report would land on the trace id the first one already claimed.
173-
test("reports one failure for a turn that fails, retries, and fails again", () => {
304+
test("settles the latest failed retry exactly once", () => {
174305
const failures: { turnIndex: number; error: string }[] = [];
175306
const runSink = createRunSink({
176307
emitter: new EventEmitter(),
@@ -182,11 +313,18 @@ describe("createRunSink", () => {
182313
runSink.sink(event("inference.error", { error: { message: "500 upstream" } }));
183314
runSink.sink(event("inference.start", {}));
184315
runSink.sink(event("inference.error", { error: { message: "500 upstream again" } }));
316+
runSink.sink(
317+
event("message.run.ended", {
318+
messageRunId: "run-1",
319+
messageId: "message-1",
320+
status: "failed",
321+
}),
322+
);
185323

186-
expect(failures).toEqual([{ turnIndex: 0, error: "500 upstream" }]);
324+
expect(failures).toEqual([{ turnIndex: 0, error: "500 upstream again" }]);
187325
});
188326

189-
test("still reports a failure after reset clears the latch", () => {
327+
test("reset discards an unresolved pending failure", () => {
190328
const failures: { turnIndex: number; error: string }[] = [];
191329
const runSink = createRunSink({
192330
emitter: new EventEmitter(),
@@ -199,11 +337,15 @@ describe("createRunSink", () => {
199337
runSink.reset();
200338
runSink.sink(event("inference.start", {}));
201339
runSink.sink(event("inference.error", { error: { message: "second session" } }));
340+
runSink.sink(
341+
event("message.run.ended", {
342+
messageRunId: "run-2",
343+
messageId: "message-2",
344+
status: "failed",
345+
}),
346+
);
202347

203-
expect(failures).toEqual([
204-
{ turnIndex: 0, error: "first session" },
205-
{ turnIndex: 0, error: "second session" },
206-
]);
348+
expect(failures).toEqual([{ turnIndex: 0, error: "second session" }]);
207349
});
208350

209351
test("seeds the turn count from a resumed session's prior turnsUsed", () => {

0 commit comments

Comments
 (0)