diff --git a/changelog.d/features/9571-plugin-streaming-usage-timing.md b/changelog.d/features/9571-plugin-streaming-usage-timing.md new file mode 100644 index 0000000000..84720d2df5 --- /dev/null +++ b/changelog.d/features/9571-plugin-streaming-usage-timing.md @@ -0,0 +1,11 @@ +- **feat(plugins):** add onStreamComplete built-in event exposing streaming usage and timing (#9571) + + Adds a new `onStreamComplete` plugin event that fires after an SSE stream is fully + consumed, carrying usage token counts and timing metrics (latency, TTFT). Built-in + events now include `onStreamComplete` as a fire-and-forget lifecycle hook. + + Payload: `status`, `usage` (prompt_tokens, completion_tokens, reasoning_tokens, + cache_read_input_tokens, cache_creation_input_tokens), `timing` (latencyMs, ttft), + `model`, `provider`, `errorCode`. + + Non-breaking — existing `onResponse` hooks with `{ streamed: true }` remain unchanged. diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index e71de50bd8..4130bed738 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -251,7 +251,10 @@ import { recordCompressionCacheStats } from "./chatCore/compressionCacheStats.ts import { writeCavemanOutputAnalytics } from "./chatCore/cavemanOutputAnalytics.ts"; import { scheduleQuotaShareConsumption } from "./chatCore/quotaShareConsumption.ts"; import { emitRequestGamificationEvent } from "./chatCore/gamificationEvent.ts"; -import { runPluginOnResponseHook } from "./chatCore/pluginOnResponse.ts"; +import { + runPluginOnResponseHook, + runPluginOnStreamCompleteHook, +} from "./chatCore/pluginOnResponse.ts"; import { scheduleStreamingQuotaShareConsumption } from "./chatCore/streamingQuotaShare.ts"; import { recordStreamingUsageStats } from "./chatCore/streamingUsageStats.ts"; import { recordStreamingCost } from "./chatCore/streamingCost.ts"; @@ -4934,6 +4937,17 @@ export async function handleChatCore({ streamUsage, log, }); + + // Plugin onStreamComplete hook — fire-and-forget, fail-open (#9571) + runPluginOnStreamCompleteHook({ + status: normalizedStreamStatus, + usage: streamUsage as Record | undefined, + ttft, + model, + provider, + errorCode: streamErrorCode, + startTime, + }); }; const streamFailureFinalizers = streamFailure.createStreamFailureFinalizers({ diff --git a/open-sse/handlers/chatCore/pluginOnResponse.ts b/open-sse/handlers/chatCore/pluginOnResponse.ts index 63055e2e74..eb100c5b5e 100644 --- a/open-sse/handlers/chatCore/pluginOnResponse.ts +++ b/open-sse/handlers/chatCore/pluginOnResponse.ts @@ -45,3 +45,57 @@ export async function runPluginOnResponseHook(args: { /* plugin onResponse optional */ } } + +/** + * Payload passed to plugin onStreamComplete hooks after a streaming response is consumed. + * Carries usage token counts, timing metrics (latency, TTFT), model, provider, and error code. + */ +export type PluginOnStreamCompletePayload = { + status: number; + usage?: { + prompt_tokens?: number; + completion_tokens?: number; + reasoning_tokens?: number; + cache_read_input_tokens?: number; + cache_creation_input_tokens?: number; + }; + timing?: { + latencyMs: number; + ttft?: number; + }; + model?: string; + provider?: string; + errorCode?: string; +}; + +/** + * Run plugin onStreamComplete hooks — fire-and-forget and fail-open. + * Called inside the onStreamComplete callback (chatCore.ts) where usage and timing data + * converge after an SSE stream is fully consumed. + */ +export async function runPluginOnStreamCompleteHook(args: { + status: number; + usage?: Record; + ttft?: number; + model: string | null | undefined; + provider: string | null | undefined; + errorCode?: string | null | undefined; + startTime: number; +}): Promise { + try { + const { runOnStreamComplete } = await import("@/lib/plugins/hooks"); + runOnStreamComplete({ + status: args.status, + usage: args.usage as PluginOnStreamCompletePayload["usage"], + timing: { + latencyMs: Date.now() - args.startTime, + ttft: args.ttft, + }, + model: args.model ?? undefined, + provider: args.provider ?? undefined, + errorCode: args.errorCode ?? undefined, + }).catch(() => {}); + } catch (_) { + /* plugin onStreamComplete optional */ + } +} diff --git a/tests/unit/chatcore-plugin-onresponse.test.ts b/tests/unit/chatcore-plugin-onresponse.test.ts index ffd8161fd9..5125a07647 100644 --- a/tests/unit/chatcore-plugin-onresponse.test.ts +++ b/tests/unit/chatcore-plugin-onresponse.test.ts @@ -6,9 +6,7 @@ import { test, after } from "node:test"; import assert from "node:assert/strict"; -const { registerHook, unregisterHook } = await import("../../src/lib/plugins/hooks.ts"); -const { runPluginOnResponseHook } = - await import("../../open-sse/handlers/chatCore/pluginOnResponse.ts"); +const { registerHook, unregisterHook } = await import("../../src/lib/plugins/hooks.ts");const { runPluginOnResponseHook, runPluginOnStreamCompleteHook } = await import("../../open-sse/handlers/chatCore/pluginOnResponse.ts"); async function waitFor(pred: () => boolean, timeoutMs = 2000): Promise { const deadline = Date.now() + timeoutMs; @@ -19,6 +17,7 @@ async function waitFor(pred: () => boolean, timeoutMs = 2000): Promise { after(() => { unregisterHook("onResponse", "test-onresponse-plugin"); + unregisterHook("onStreamComplete", "test-onstreamcomplete-plugin"); }); test("no registered hooks → resolves without throwing (no-op)", async () => { @@ -143,3 +142,140 @@ test("a throwing hook never rejects the caller (fail-open)", async () => { ); await new Promise((r) => setTimeout(r, 30)); }); + +// ── onStreamComplete hook tests (#9571) ── + +test("onStreamComplete: no registered hooks resolves without throwing (no-op)", async () => { + const start = Date.now(); + await assert.doesNotReject( + runPluginOnStreamCompleteHook({ + status: 200, + usage: { prompt_tokens: 10, completion_tokens: 20 }, + ttft: 150, + model: "gpt-4", + provider: "openai", + errorCode: undefined, + startTime: start - 500, + }) + ); +}); + +test("onStreamComplete: registered hook receives usage + timing payload", async () => { + let captured: Record | undefined; + registerHook( + "onStreamComplete", + "test-onstreamcomplete-plugin", + async (payload: Record) => { + captured = payload; + } + ); + + const startTime = Date.now() - 500; + await runPluginOnStreamCompleteHook({ + status: 200, + usage: { prompt_tokens: 42, completion_tokens: 100, reasoning_tokens: 5 }, + ttft: 200, + model: "claude-3-opus", + provider: "anthropic", + errorCode: undefined, + startTime, + }); + + await waitFor(() => captured !== undefined); + assert.ok(captured, "expected onStreamComplete hook to be invoked"); + + // payload shape: status, usage, timing, model, provider + assert.equal(captured!.status, 200); + assert.ok(captured!.usage, "usage should be present"); + assert.equal((captured!.usage as Record).prompt_tokens, 42); + assert.equal((captured!.usage as Record).completion_tokens, 100); + assert.equal((captured!.usage as Record).reasoning_tokens, 5); + + assert.ok(captured!.timing, "timing should be present"); + const timing = captured!.timing as Record; + assert.equal(timing.ttft, 200); + assert.ok(timing.latencyMs > 450, "latencyMs should be near 500"); + + assert.equal(captured!.model, "claude-3-opus"); + assert.equal(captured!.provider, "anthropic"); + assert.equal(captured!.errorCode, undefined); +}); + +test("onStreamComplete: payload includes cache token fields when present", async () => { + let captured: Record | undefined; + registerHook( + "onStreamComplete", + "test-onstreamcomplete-plugin", + async (payload: Record) => { + captured = payload; + } + ); + + await runPluginOnStreamCompleteHook({ + status: 200, + usage: { + prompt_tokens: 50, + completion_tokens: 30, + cache_read_input_tokens: 20, + cache_creation_input_tokens: 10, + }, + ttft: 100, + model: "gpt-4", + provider: "openai", + errorCode: undefined, + startTime: Date.now(), + }); + + await waitFor(() => captured !== undefined); + assert.ok(captured); + const usage = captured!.usage as Record; + assert.equal(usage.cache_read_input_tokens, 20); + assert.equal(usage.cache_creation_input_tokens, 10); +}); + +test("onStreamComplete: throwing hook never rejects the caller (fail-open)", async () => { + registerHook("onStreamComplete", "test-onstreamcomplete-plugin", async () => { + throw new Error("stream-complete-boom"); + }); + + await assert.doesNotReject( + runPluginOnStreamCompleteHook({ + status: 500, + usage: undefined, + ttft: undefined, + model: "gpt-4", + provider: "openai", + errorCode: "upstream_error", + startTime: Date.now(), + }) + ); + await new Promise((r) => setTimeout(r, 30)); +}); + +test("onStreamComplete: errorCode is passed through when provided", async () => { + let captured: Record | undefined; + registerHook( + "onStreamComplete", + "test-onstreamcomplete-plugin", + async (payload: Record) => { + captured = payload; + } + ); + + await runPluginOnStreamCompleteHook({ + status: 502, + usage: undefined, + ttft: undefined, + model: "grok-3", + provider: "xai", + errorCode: "upstream_timeout", + startTime: Date.now(), + }); + + await waitFor(() => captured !== undefined); + assert.ok(captured); + assert.equal(captured!.status, 502); + assert.equal(captured!.errorCode, "upstream_timeout"); + assert.equal(captured!.model, "grok-3"); + assert.equal(captured!.provider, "xai"); +});