mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-09 00:32:13 +03:00
Compare commits
4 Commits
feat/9544-
...
feat/9571-
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4a2a9978bc | ||
|
|
867fbc8bad | ||
|
|
d557ea2222 | ||
|
|
1bdd1be8c6 |
11
changelog.d/features/9571-plugin-streaming-usage-timing.md
Normal file
11
changelog.d/features/9571-plugin-streaming-usage-timing.md
Normal file
@@ -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.
|
||||
@@ -245,7 +245,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";
|
||||
@@ -4895,6 +4898,17 @@ export async function handleChatCore({
|
||||
streamUsage,
|
||||
log,
|
||||
});
|
||||
|
||||
// Plugin onStreamComplete hook — fire-and-forget, fail-open (#9571)
|
||||
runPluginOnStreamCompleteHook({
|
||||
status: normalizedStreamStatus,
|
||||
usage: streamUsage as Record<string, unknown> | undefined,
|
||||
ttft,
|
||||
model,
|
||||
provider,
|
||||
errorCode: streamErrorCode,
|
||||
startTime,
|
||||
});
|
||||
};
|
||||
|
||||
const streamFailureFinalizers = streamFailure.createStreamFailureFinalizers({
|
||||
|
||||
@@ -43,3 +43,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<string, unknown>;
|
||||
ttft?: number;
|
||||
model: string | null | undefined;
|
||||
provider: string | null | undefined;
|
||||
errorCode?: string | null | undefined;
|
||||
startTime: number;
|
||||
}): Promise<void> {
|
||||
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 */
|
||||
}
|
||||
}
|
||||
|
||||
@@ -40,6 +40,7 @@ export const BUILTIN_EVENTS = [
|
||||
"onActivate",
|
||||
"onDeactivate",
|
||||
"onUninstall",
|
||||
"onStreamComplete",
|
||||
] as const;
|
||||
|
||||
export type BuiltinEvent = (typeof BUILTIN_EVENTS)[number];
|
||||
@@ -251,6 +252,35 @@ export interface Plugin {
|
||||
onActivate?: (payload: unknown) => Promise<void> | void;
|
||||
onDeactivate?: (payload: unknown) => Promise<void> | void;
|
||||
onUninstall?: (payload: unknown) => Promise<void> | void;
|
||||
onStreamComplete?: (payload: PluginOnStreamCompletePayload) => Promise<void> | void;
|
||||
}
|
||||
|
||||
// ── onStreamComplete event types ──
|
||||
|
||||
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 onStreamComplete hooks — fire-and-forget notification with usage/timing data.
|
||||
* Called when an SSE stream is fully consumed and usage/timing data is available.
|
||||
*/
|
||||
export async function runOnStreamComplete(payload: PluginOnStreamCompletePayload): Promise<void> {
|
||||
await emitHook("onStreamComplete", payload);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -7,9 +7,8 @@ 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 { runPluginOnResponseHook, runPluginOnStreamCompleteHook } =
|
||||
await import("../../open-sse/handlers/chatCore/pluginOnResponse.ts");
|
||||
|
||||
async function waitFor(pred: () => boolean, timeoutMs = 2000): Promise<void> {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
@@ -20,6 +19,7 @@ async function waitFor(pred: () => boolean, timeoutMs = 2000): Promise<void> {
|
||||
|
||||
after(() => {
|
||||
unregisterHook("onResponse", "test-onresponse-plugin");
|
||||
unregisterHook("onStreamComplete", "test-onstreamcomplete-plugin");
|
||||
});
|
||||
|
||||
test("no registered hooks → resolves without throwing (no-op)", async () => {
|
||||
@@ -101,3 +101,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<string, unknown> | undefined;
|
||||
registerHook(
|
||||
"onStreamComplete",
|
||||
"test-onstreamcomplete-plugin",
|
||||
async (payload: Record<string, unknown>) => {
|
||||
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<string, number>).prompt_tokens, 42);
|
||||
assert.equal((captured!.usage as Record<string, number>).completion_tokens, 100);
|
||||
assert.equal((captured!.usage as Record<string, number>).reasoning_tokens, 5);
|
||||
|
||||
assert.ok(captured!.timing, "timing should be present");
|
||||
const timing = captured!.timing as Record<string, number>;
|
||||
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<string, unknown> | undefined;
|
||||
registerHook(
|
||||
"onStreamComplete",
|
||||
"test-onstreamcomplete-plugin",
|
||||
async (payload: Record<string, unknown>) => {
|
||||
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<string, number>;
|
||||
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<string, unknown> | undefined;
|
||||
registerHook(
|
||||
"onStreamComplete",
|
||||
"test-onstreamcomplete-plugin",
|
||||
async (payload: Record<string, unknown>) => {
|
||||
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");
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user