From 8317bf78742f5b9858bf795e153d62a1eb9ca45a Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 01:44:21 +0000 Subject: [PATCH 1/8] fix(guardrails): raise the raw hold-cap factor from 64 to 128 OpenAI streams of CJK output frame at 74-75x their generated content on both Chat Completions and Responses. On the routes that count the upstream's own bytes (native /v1/responses, passthrough routes, native /v1/messages) the 64x raw bound therefore tripped at about 86% of the configured max_buffer_bytes, contradicting the rule that the cap counts generated content. At 128x the content cap binds again. Worst-case memory per held stream at the default 256 KiB goes from 16 MiB to 32 MiB. --- crates/aisix-proxy/src/held_content.rs | 14 +- .../stream-output-raw-hold-cap-e2e.test.ts | 159 +++++++++++++++++- 2 files changed, 165 insertions(+), 8 deletions(-) diff --git a/crates/aisix-proxy/src/held_content.rs b/crates/aisix-proxy/src/held_content.rs index 8a1f6e84c..d64a6a045 100644 --- a/crates/aisix-proxy/src/held_content.rs +++ b/crates/aisix-proxy/src/held_content.rs @@ -22,10 +22,11 @@ use aisix_gateway::ChatDelta; use serde_json::Value; /// Raw bytes a hold-back may keep, as a multiple of `max_buffer_bytes` -/// (16 MiB at the 256 KiB default). Far above the framing an ordinary -/// token stream wraps around its content, so a normal response still trips -/// on content first. -pub(crate) const RAW_HOLD_FACTOR: usize = 64; +/// (32 MiB at the 256 KiB default). Above the framing an ordinary token +/// stream wraps around its content — OpenAI streams of CJK text measure +/// about 75× their content, on Chat Completions and Responses alike — so a +/// normal response still trips on content first. +pub(crate) const RAW_HOLD_FACTOR: usize = 128; /// What one hold-back buffer holds: generated content (the cap /// `max_buffer_bytes` names) and the raw bytes kept to hold it. @@ -73,6 +74,11 @@ impl BoundedValues { self.values.push(v.clone()); } + /// Whether a value was refused for lack of room. + pub(crate) fn is_full(&self) -> bool { + self.full + } + /// The kept values, or `None` when there are none. pub(crate) fn take(&mut self) -> Option> { self.bytes = 0; diff --git a/tests/e2e/src/cases/stream-output-raw-hold-cap-e2e.test.ts b/tests/e2e/src/cases/stream-output-raw-hold-cap-e2e.test.ts index df1f7729b..99a14fcb2 100644 --- a/tests/e2e/src/cases/stream-output-raw-hold-cap-e2e.test.ts +++ b/tests/e2e/src/cases/stream-output-raw-hold-cap-e2e.test.ts @@ -13,7 +13,7 @@ import { // E2E for the raw-byte bound on a held-back stream. `max_buffer_bytes` // counts generated content only (#513), but the hold-back keeps whole frames, -// so each hold-back also bounds the raw bytes it keeps at 64 times the cap. +// so each hold-back also bounds the raw bytes it keeps at 128 times the cap. // // - A stream of frames that carry no content (pings, empty deltas, // keep-alives, base64 image previews) past that bound is a buffer-exceeded @@ -21,6 +21,9 @@ import { // releases unscanned. // - An ordinary token-by-token text stream still trips on the content cap: // content exactly at the cap is scanned and released, one byte more trips. +// - So does one whose framing runs at ~100× its content, as OpenAI streams +// of CJK text do (~75×): on the routes that measure the upstream's own +// frames, the raw bound must not cut it before the content cap. // // Every stream ends with an email the guardrail masks whenever it scans, so a // masked email proves nothing tripped and a raw one proves the stream went out @@ -29,8 +32,8 @@ import { const CALLER = "sk-stream-raw-hold-caller"; const CALLER_HASH = createHash("sha256").update(CALLER).digest("hex"); const CAP = 1_000; -// Comfortably past the raw bound (64 × CAP = 64 000 bytes). -const RAW_TARGET = 3 * 64 * CAP; +// Comfortably past the raw bound (128 × CAP = 128 000 bytes). +const RAW_TARGET = 3 * 128 * CAP; const EMAIL = "raw-probe@example.com"; const MASKED = "[EMAIL_REDACTED]"; const TAIL = `reach me at ${EMAIL}`; @@ -76,7 +79,7 @@ const CHAT_EMPTY_TOOL_CALLS = [ // raw bound; the terminal event the bridge adds echoes them a third time and // carries the stream past it. const CHAT_NO_FINISH = [chatChunk({ role: "assistant" }), chatChunk({ content: TAIL })]; -const LONG_INSTRUCTIONS = `Be brief. ${"i".repeat(25 * CAP)}`; +const LONG_INSTRUCTIONS = `Be brief. ${"i".repeat(50 * CAP)}`; // Token-by-token text whose content totals `bytes`, ending with TAIL. const tokens = (bytes: number) => { @@ -160,6 +163,112 @@ const RESPONSES_PARTIAL_IMAGES = [ }), ]; +// Frames whose envelope is ~HIGH_RATIO times the content they carry: each +// token's frame is padded with an `obfuscation` field (which OpenAI streams +// really carry) or, on Anthropic, followed by `ping` events, until the raw +// bytes so far reach HIGH_RATIO × the content so far. +const HIGH_RATIO = 100; +const sseData = (payload: string) => `data: ${payload}\n\n`; +const obfuscated = (build: (pad: string) => string, content: number) => { + const bare = sseData(build("")).length; + return build("o".repeat(Math.max(0, HIGH_RATIO * content - bare))); +}; + +const responsesHighRatio = (bytes: number) => [ + sseData( + JSON.stringify({ + type: "response.created", + response: { id: "resp_ratio", object: "response", status: "in_progress", model: "gpt-4o-mini", output: [] }, + }), + ), + ...tokens(bytes).map((t, i) => + sseData( + obfuscated( + (pad) => + JSON.stringify({ + type: "response.output_text.delta", + item_id: "msg_ratio", + output_index: 0, + content_index: 0, + delta: t, + sequence_number: i, + obfuscation: pad, + }), + t.length, + ), + ), + ), + sseData( + JSON.stringify({ + type: "response.completed", + response: { + id: "resp_ratio", + object: "response", + status: "completed", + model: "gpt-4o-mini", + output: [], + usage: { input_tokens: 5, output_tokens: 40, total_tokens: 45 }, + }, + }), + ), +]; + +const PING = anthropicFrame("ping", {}); +const anthropicHighRatio = (bytes: number) => { + const text = tokens(bytes); + const frames: string[] = []; + let raw = 0; + let content = 0; + for (const t of text) { + const f = anthropicFrame("content_block_delta", { index: 0, delta: { type: "text_delta", text: t } }); + frames.push(f); + raw += f.length; + content += t.length; + while (raw < HIGH_RATIO * content) { + frames.push(PING); + raw += PING.length; + } + } + return [ + anthropicFrame("message_start", { + message: { + id: "msg_ratio", + type: "message", + role: "assistant", + content: [], + model: "claude-3-5-haiku-20241022", + stop_reason: null, + usage: { input_tokens: 5, output_tokens: 1 }, + }, + }), + anthropicFrame("content_block_start", { index: 0, content_block: { type: "text", text: "" } }), + ...frames, + anthropicFrame("content_block_stop", { index: 0 }), + anthropicFrame("message_delta", { delta: { stop_reason: "end_turn" }, usage: { output_tokens: 40 } }), + anthropicFrame("message_stop", {}), + ]; +}; + +const chatHighRatio = (bytes: number) => [ + ...tokens(bytes).map((t) => + sseData( + obfuscated( + (pad) => + JSON.stringify({ + id: "chatcmpl-ratio", + object: "chat.completion.chunk", + model: "gpt-4o-mini", + choices: [{ index: 0, delta: { content: t }, finish_reason: null }], + obfuscation: pad, + }), + t.length, + ), + ), + ), + sseData(chatChunk({}, "stop")), + sseData("[DONE]"), +]; + // Passthrough: keep-alive chunks with no choices. const PASSTHROUGH_KEEPALIVES = [ ...repeatToBytes(JSON.stringify({ id: "chatcmpl-raw", object: "chat.completion.chunk", choices: [] }), RAW_TARGET), @@ -260,6 +369,23 @@ describe("held-back stream raw-byte bound", () => { await model("raw-chat-over-cap", "openai", chatText(CAP + 1), policies.closed); await model("raw-msg-at-cap", "anthropic", anthropicText(CAP), policies.closed, { raw: true }); await model("raw-msg-over-cap", "anthropic", anthropicText(CAP + 1), policies.closed, { raw: true }); + await model("raw-resp-ratio", "openai", responsesHighRatio(CAP), policies.closed, { raw: true }); + await model("raw-msg-ratio", "anthropic", anthropicHighRatio(CAP), policies.closed, { raw: true }); + const ratioBacking = await model("raw-route-ratio-backing", "openai", chatHighRatio(CAP), policies.closed, { + raw: true, + }); + const ratioRoute = await seed.createPassthroughRoute({ + name: "raw-route-ratio", + path_prefix: "/passthrough/raw-ratio", + target_url: String(ratioBacking.value.api_base), + provider_key_id: ratioBacking.id, + }); + await seed.update("guardrail_attachments", randomUUID(), { + guardrail_id: routePolicies.closed!.id, + scope_type: "passthrough_route", + scope_id: ratioRoute.id, + priority: 100, + }); // Seeded last: this key authenticating implies the whole seed set landed. await seed.createApiKey({ key_hash: CALLER_HASH, allowed_models: ["*"], allowed_routes: ["*"] }); @@ -353,4 +479,29 @@ describe("held-back stream raw-byte bound", () => { expect(overCap).not.toContain(EMAIL); }); } + + // Raw frames between 64× and 128× the content: every route that measures + // the upstream's own bytes must hold them until the content cap binds. + const ratioFrames = (frames: string[]) => frames.join("").length; + test("high-ratio fixtures sit between 64× and 128× the content cap", () => { + for (const frames of [responsesHighRatio(CAP), anthropicHighRatio(CAP), chatHighRatio(CAP)]) { + expect(ratioFrames(frames)).toBeGreaterThan(64 * CAP); + expect(ratioFrames(frames)).toBeLessThan(128 * CAP); + } + }); + + const ratioCases: Array<[string, () => Promise, string]> = [ + ["/v1/responses (native)", () => responses("raw-resp-ratio"), MASKED], + ["/v1/messages (native)", () => messages("raw-msg-ratio"), MASKED], + ["passthrough route", () => route("ratio"), "content_filter"], + ]; + for (const [surface, send, scanned] of ratioCases) { + test(`${surface}: content at the cap framed at ~100× is scanned, not cut by the raw bound`, async (ctx) => { + if (!ready(ctx)) return; + const body = await send(); + expect(body).not.toContain("output_buffer_exceeded"); + expect(body, "the guardrail scanned the whole response").toContain(scanned); + expect(body).not.toContain(EMAIL); + }); + } }); From e1e1a3324d8ab99544c03fb1c568ed831bfde4b8 Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 01:44:41 +0000 Subject: [PATCH 2/8] fix(guardrails): monitor-only output chains scan all of the generated text A monitor-only output chain (EndOfStreamCheck) holds nothing back, yet several streamed routes silently cut its end-of-stream scan at 256 KiB: native /v1/responses and streamed audio transcription (EosOutputScan and the accumulators feeding it), the /v1/responses chat bridge's live-mode scan, and the chat tool-call text on /v1/chat/completions. A monitor rule whose trigger arrived past that point was never recorded. Every route now scans the whole generated text, as chat body text and native /v1/messages already did. Only generated text is accumulated; the chat route's raw tool-call deltas stay bounded, and past that bound the local kinds judge the tool-call text instead. --- crates/aisix-proxy/src/audio.rs | 24 +- crates/aisix-proxy/src/chat.rs | 31 +- crates/aisix-proxy/src/guardrail_stream.rs | 17 +- crates/aisix-proxy/src/responses.rs | 20 +- crates/aisix-proxy/src/responses_bridge.rs | 18 -- ...drail-monitor-full-output-scan-e2e.test.ts | 302 ++++++++++++++++++ 6 files changed, 342 insertions(+), 70 deletions(-) create mode 100644 tests/e2e/src/cases/guardrail-monitor-full-output-scan-e2e.test.ts diff --git a/crates/aisix-proxy/src/audio.rs b/crates/aisix-proxy/src/audio.rs index d0e2fecdb..2b86dc5ac 100644 --- a/crates/aisix-proxy/src/audio.rs +++ b/crates/aisix-proxy/src/audio.rs @@ -871,8 +871,8 @@ struct StreamedTranscript { impl StreamedTranscript { /// The transcript the caller received, for the end-of-stream scan and - /// the content capture. Both are capped at their own limits on top of - /// this one. + /// the content capture. The capture is capped at its own limit on top + /// of this one. fn text(&self) -> &str { self.terminal.as_deref().unwrap_or(&self.deltas) } @@ -978,9 +978,8 @@ where /// once the upstream ends or the caller disconnects, whichever comes /// first; it owns the UsageEvent for this request. /// -/// `capture_cap` bounds the assembled transcript; an output chain that -/// needs a scan raises the floor to its own scan bound, so neither -/// consumer sees past its own limit. +/// `content_cap` bounds the assembled transcript, unless an output chain +/// needs the end-of-stream scan, which reads all of it. fn transcription_relay( upstream: S, // Set when `upstream` ended on a read timeout rather than its own end. @@ -997,14 +996,13 @@ where // Anchors a timed-out read's reported elapsed time. let started = std::time::Instant::now(); - let text_cap = content_cap - .map(|cap| cap as usize) - .unwrap_or(0) - .max(if eos_scan.is_some() { - aisix_guardrails::DEFAULT_STREAM_OUTPUT_BUFFER_BYTES - } else { - 0 - }); + // The end-of-stream scan reads the whole transcript; the capture + // re-truncates to its own cap. + let text_cap = if eos_scan.is_some() { + usize::MAX + } else { + content_cap.map_or(0, |cap| cap as usize) + }; // Re-attach the request span: the body is polled after the request-id // middleware returns, so anything logged from here would otherwise lose // its `request_id` correlation (AISIX-Cloud#1060). diff --git a/crates/aisix-proxy/src/chat.rs b/crates/aisix-proxy/src/chat.rs index 870ca502c..1238e9dde 100644 --- a/crates/aisix-proxy/src/chat.rs +++ b/crates/aisix-proxy/src/chat.rs @@ -5615,7 +5615,8 @@ where // collected). Under Window every window re-scans the whole buffer, so // the chain's folded hold cap bounds it, and outgrowing that cap is a // buffer trip like any other (#513, #1029). The live-forward monitor - // path bounds it at the default cap. Allocated only with a guardrail. + // path holds nothing back, so it collects all of it, like the content + // buffer above. Allocated only with a guardrail. let mut tool_calls_buf = if output_guardrail.is_some() { Some(String::new()) } else { @@ -5625,7 +5626,7 @@ where // the channels the buffered branch's mask walker reads (#1027). // Only the live-forward branch reads them; bounded by their own size, // since a delta with no name or arguments adds nothing to the buffer - // above. + // above. Past that bound the check judges the tool-call text instead. let mut eos_tool_calls = crate::held_content::BoundedValues::default(); // P2 (#379) / #466: streamed-output policy folded over the output-hook // guardrails. EndOfStreamCheck (reached only when no output-hook @@ -5637,12 +5638,6 @@ where .map(|ctx| ctx.chain.stream_output_policy()) .unwrap_or_default(); let hold_back = stream_policy.holds_back(); - let tool_calls_cap = match stream_policy { - aisix_guardrails::StreamOutputPolicy::EndOfStreamCheck => { - aisix_guardrails::DEFAULT_STREAM_OUTPUT_BUFFER_BYTES - } - _ => usize::MAX, - }; let mut window_tool_calls_cap = match stream_policy { aisix_guardrails::StreamOutputPolicy::Window { .. } => stream_policy.hold_cap(), _ => None, @@ -5802,9 +5797,6 @@ where tool_calls_overflowed, ) { for tc in tcs { - if buf.len() >= tool_calls_cap { - break; - } let function = tc.get("function"); let name = function .and_then(|f| f.get("name")) @@ -5823,10 +5815,10 @@ where if !hold_back { // Serialized deltas carry their envelope, so // they get the raw guard's headroom over the - // text cap the buffer above keeps. + // default text cap. eos_tool_calls.push( tc, - tool_calls_cap + aisix_guardrails::DEFAULT_STREAM_OUTPUT_BUFFER_BYTES .saturating_mul(crate::held_content::RAW_HOLD_FACTOR), ); } @@ -6332,20 +6324,29 @@ where let verdict = if verdict.is_block() { verdict } else { + // Deltas past their bound: judge the tool-call text whole. + let tool_calls_whole = eos_tool_calls.is_full(); let mut chunks = vec![aisix_gateway::ChatChunk { id: String::new(), model: String::new(), delta: aisix_gateway::ChatDelta { content: Some(content.clone()), - tool_calls: eos_tool_calls.take(), + tool_calls: eos_tool_calls.take().filter(|_| !tool_calls_whole), ..Default::default() }, finish_reason: None, usage: None, }]; - let segments = crate::redact::collect_segments(|g| { + let mut segments = crate::redact::collect_segments(|g| { let _ = crate::redact::redact_chat_chunks(g, &mut chunks); }); + if tool_calls_whole { + segments.push(aisix_guardrails::ScanSegment { + text: tc_part.to_string(), + role: aisix_guardrails::SegmentRole::ScanOnly, + in_latest_turn: true, + }); + } let (local, hits) = ctx.chain.check_local_segments(&segments, false); guard.comp().monitor_hits.extend(hits); verdict.merged_with(local) diff --git a/crates/aisix-proxy/src/guardrail_stream.rs b/crates/aisix-proxy/src/guardrail_stream.rs index 6f7f12f3f..fab0478e1 100644 --- a/crates/aisix-proxy/src/guardrail_stream.rs +++ b/crates/aisix-proxy/src/guardrail_stream.rs @@ -63,19 +63,12 @@ impl EosOutputScan { text: &str, local_segments: Option>, ) -> Vec { - // Bound the provider calls the same way the buffered branch's byte - // cap does — scan at most the cap's worth of text. - let mut end = text - .len() - .min(aisix_guardrails::DEFAULT_STREAM_OUTPUT_BUFFER_BYTES); - while end > 0 && !text.is_char_boundary(end) { - end -= 1; - } - let scan_text = &text[..end]; - if scan_text.is_empty() { + // The whole generated text, as chat and `/v1/messages` scan it: a + // monitor-only chain holds nothing back, so no hold cap applies. + if text.is_empty() { return Vec::new(); } - let synth = synth_chat_response(&self.upstream_model, scan_text.to_string()); + let synth = synth_chat_response(&self.upstream_model, text.to_string()); let (verdict, mut hits) = aisix_guardrails::Guardrail::check_output_non_segment_observed( self.chain.as_ref(), &synth, @@ -95,7 +88,7 @@ impl EosOutputScan { &mut hits, local_segments, |g| { - let _ = g.redact_output_text(scan_text); + let _ = g.redact_output_text(text); crate::redact::RedactionCounts::new() }, ) diff --git a/crates/aisix-proxy/src/responses.rs b/crates/aisix-proxy/src/responses.rs index 8c39efbd9..4a08d82c3 100644 --- a/crates/aisix-proxy/src/responses.rs +++ b/crates/aisix-proxy/src/responses.rs @@ -3281,21 +3281,17 @@ where // at end-of-stream — with the estimation cap as the floor. The cap // bounds delta accumulation; a terminal event's full output text is // instead bounded by MAX_SSE_FRAME_BUF_BYTES (an oversized frame - // never parses) and re-truncated by each consumer — - // `CapturedContent::new` at the exporter cap and - // `EosOutputScan::observe` at the scan bound — so none sees beyond - // its own limit. - let capture_cap = Some( + // never parses). The end-of-stream scan reads all of it, so with one + // the accumulator is unbounded; `CapturedContent::new` re-truncates + // to the exporter cap either way. + let capture_cap = Some(if eos_scan.is_some() { + usize::MAX + } else { content_cap .map(|cap| cap as usize) .unwrap_or(0) - .max(if eos_scan.is_some() { - aisix_guardrails::DEFAULT_STREAM_OUTPUT_BUFFER_BYTES - } else { - 0 - }) - .max(crate::token_estimate::OUTPUT_ACCUMULATION_CAP), - ); + .max(crate::token_estimate::OUTPUT_ACCUMULATION_CAP) + }); // Re-attach the request span: the body is polled after the request-id // middleware returns, so the end-of-stream output-guardrail scan // (`EosOutputScan::observe`) would otherwise log without a diff --git a/crates/aisix-proxy/src/responses_bridge.rs b/crates/aisix-proxy/src/responses_bridge.rs index 95b68a6f3..bc8c1185a 100644 --- a/crates/aisix-proxy/src/responses_bridge.rs +++ b/crates/aisix-proxy/src/responses_bridge.rs @@ -2695,24 +2695,6 @@ pub fn build_responses_bridge_stream( // EndOfStreamCheck behavior. let (text, tool_calls) = encoder.assembled_assistant_message(); if !text.is_empty() || !tool_calls.is_empty() { - // Live mode releases oversized streams (that's the point of - // AISIX-Cloud#1010), so the assembled text is unbounded here — - // cap the scan input like the verbatim path's EosOutputScan - // does, keeping the observation provider calls bounded. Held - // (buffering) text is already capped by the hold-back budget. - let text = if buffering { - text - } else { - let mut text = text; - let mut end = text - .len() - .min(aisix_guardrails::DEFAULT_STREAM_OUTPUT_BUFFER_BYTES); - while end > 0 && !text.is_char_boundary(end) { - end -= 1; - } - text.truncate(end); - text - }; // Live mode has no held frames for the segment pass to walk — // offer the flattened text as one segment so monitor-mode // segment moderators still record their observations. diff --git a/tests/e2e/src/cases/guardrail-monitor-full-output-scan-e2e.test.ts b/tests/e2e/src/cases/guardrail-monitor-full-output-scan-e2e.test.ts new file mode 100644 index 000000000..a54ef6459 --- /dev/null +++ b/tests/e2e/src/cases/guardrail-monitor-full-output-scan-e2e.test.ts @@ -0,0 +1,302 @@ +import { createHash } from "node:crypto"; +import { afterAll, beforeAll, describe, expect, test } from "vitest"; +import { + EtcdClient, + ProxyClient, + SeedClient, + spawnApp, + startMockSls, + startOpenAiUpstream, + waitConfigPropagation, + waitForSlsLog, + type MockSls, + type OpenAiUpstream, + type SpawnedApp, +} from "../harness/index.js"; + +// E2E: a monitor-mode output guardrail on a streamed response judges ALL of +// the generated text, on every streamed route. A monitor-only chain never +// holds the stream back, so no hold cap applies to what it scans: a phrase +// that arrives after 256 KiB of output is still a would-be block on the +// request's usage event. +// +// The row is `kind: custom`, which judges the flattened text, so the chat +// route's tool-call arguments reach it through the same text the other +// routes scan. + +const CALLER = "sk-monitor-full-output-scan-e2e"; +const CREDENTIAL_REF = "mock"; +const LOGSTORE = "monitor-full-output-scan-events"; +const GUARD_NAME = "monitor-full-output-scan"; +const hash = (s: string) => createHash("sha256").update(s).digest("hex"); + +const MARKER = "latemonitormarker"; +const SCRIPT = ` +export function checkOutput(ctx) { + return ctx.text.includes("${MARKER}") ? { action: "block" } : { action: "none" }; +}`; + +// 220 × 1250 = 275 000 bytes of clean output ahead of the marker: past the +// 256 KiB (262 144-byte) default hold cap. +const PIECE = "z".repeat(1250); +const pieces = [...Array.from({ length: 220 }, () => PIECE), ` then ${MARKER}`]; + +const chatChunk = (delta: Record, finish: string | null = null) => + JSON.stringify({ + id: "chatcmpl-late", + object: "chat.completion.chunk", + created: 1, + model: "gpt-4o-mini", + choices: [{ index: 0, delta, finish_reason: finish }], + }); +const chatDone = [ + JSON.stringify({ + id: "chatcmpl-late", + object: "chat.completion.chunk", + created: 1, + model: "gpt-4o-mini", + choices: [{ index: 0, delta: {}, finish_reason: "stop" }], + usage: { prompt_tokens: 5, completion_tokens: 12, total_tokens: 17 }, + }), + "[DONE]", +]; + +const chatContent = [ + chatChunk({ role: "assistant" }), + ...pieces.map((content) => chatChunk({ content })), + ...chatDone, +]; + +const chatToolCall = [ + chatChunk({ + role: "assistant", + tool_calls: [{ index: 0, id: "call_late", type: "function", function: { name: "lookup", arguments: "" } }], + }), + ...pieces.map((args) => chatChunk({ tool_calls: [{ index: 0, function: { arguments: args } }] })), + ...chatDone, +]; + +const responsesText = [ + JSON.stringify({ type: "response.created", response: { id: "resp_late" } }), + ...pieces.map((delta) => JSON.stringify({ type: "response.output_text.delta", delta })), + JSON.stringify({ + type: "response.completed", + response: { id: "resp_late", status: "completed", usage: { input_tokens: 5, output_tokens: 12 } }, + }), + "[DONE]", +]; + +const anthropicText = [ + JSON.stringify({ + type: "message_start", + message: { + id: "msg_late", + role: "assistant", + content: [], + model: "claude-3-5-haiku-20241022", + stop_reason: null, + usage: { input_tokens: 5, output_tokens: 1 }, + }, + }), + JSON.stringify({ type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }), + ...pieces.map((text) => JSON.stringify({ type: "content_block_delta", index: 0, delta: { type: "text_delta", text } })), + JSON.stringify({ type: "content_block_stop", index: 0 }), + JSON.stringify({ type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 12 } }), + JSON.stringify({ type: "message_stop" }), +]; + +// Streamed transcription deltas; the terminal event carries usage only, so +// the transcript is the assembled deltas. +const transcriptFrames = [ + ...pieces.map((delta) => `data: ${JSON.stringify({ type: "transcript.text.delta", delta })}\n\n`), + `data: ${JSON.stringify({ + type: "transcript.text.done", + usage: { type: "tokens", input_tokens: 26, output_tokens: 12, total_tokens: 38 }, + })}\n\n`, + "data: [DONE]\n\n", +]; + +interface Route { + name: string; + provider: "anthropic" | "openai"; + upstream: { streamEvents?: string[]; rawStreamFrames?: string[] }; + // `apis: {}` declares an OpenAI-compatible endpoint with no + // `/v1/responses`, which puts `/v1/responses` on the Chat bridge. + apis?: Record; + send: (proxyUrl: string, model: string) => Promise; +} + +const json = (proxyUrl: string, path: string, body: Record) => + fetch(`${proxyUrl}${path}`, { + method: "POST", + headers: { + "content-type": "application/json", + authorization: `Bearer ${CALLER}`, + "x-api-key": CALLER, + "anthropic-version": "2023-06-01", + }, + body: JSON.stringify({ ...body, stream: true }), + }); + +const ROUTES: Route[] = [ + { + name: "responses-native", + provider: "openai", + upstream: { streamEvents: responsesText }, + send: (u, model) => json(u, "/v1/responses", { model, input: "go" }), + }, + { + name: "responses-bridge", + provider: "openai", + upstream: { streamEvents: chatContent }, + apis: {}, + send: (u, model) => json(u, "/v1/responses", { model, input: "go" }), + }, + { + name: "chat-tool-call", + provider: "openai", + upstream: { streamEvents: chatToolCall }, + send: (u, model) => json(u, "/v1/chat/completions", { model, messages: [{ role: "user", content: "go" }] }), + }, + { + name: "chat-content", + provider: "openai", + upstream: { streamEvents: chatContent }, + send: (u, model) => json(u, "/v1/chat/completions", { model, messages: [{ role: "user", content: "go" }] }), + }, + { + name: "messages-native", + provider: "anthropic", + upstream: { streamEvents: anthropicText }, + send: (u, model) => json(u, "/v1/messages", { model, max_tokens: 64, messages: [{ role: "user", content: "go" }] }), + }, + { + name: "transcription", + provider: "openai", + upstream: { rawStreamFrames: transcriptFrames }, + send: (u, model) => { + const form = new FormData(); + form.set("model", model); + form.set("stream", "true"); + form.set("file", new Blob([new Uint8Array([0x49, 0x44, 0x33])], { type: "audio/mpeg" }), "a.mp3"); + return fetch(`${u}/v1/audio/transcriptions`, { + method: "POST", + headers: { authorization: `Bearer ${CALLER}` }, + body: form, + }); + }, + }, +]; + +describe("monitor-mode output guardrail scans all of a streamed response", () => { + let app: SpawnedApp | undefined; + let sls: MockSls | undefined; + const upstreams: OpenAiUpstream[] = []; + let etcdReachable = false; + + beforeAll(async () => { + const etcd = new EtcdClient(); + etcdReachable = await etcd.ping(); + if (!etcdReachable) return; + + sls = await startMockSls(); + app = await spawnApp({ + extraEnv: { + [`SLS_CRED_${CREDENTIAL_REF.toUpperCase()}_AK_ID`]: "mock-akid", + [`SLS_CRED_${CREDENTIAL_REF.toUpperCase()}_AK_SECRET`]: "mock-secret", + }, + }); + const seed = new SeedClient(etcd, app.etcdPrefix); + await seed.createObservabilityExporter({ + name: "sls-monitor-full-output-scan", + enabled: true, + kind: "aliyun_sls", + endpoint: sls.url, + project: "aisix-e2e-obs", + logstore: LOGSTORE, + credential_ref: CREDENTIAL_REF, + content_mode: "metadata_only", + }); + const monitor = await seed.createGuardrail( + { + name: GUARD_NAME, + enabled: true, + kind: "custom", + hook_point: "output", + enforcement_mode: "monitor", + output_fail_open: false, + timeout_ms: 5000, + script: SCRIPT, + }, + { attach: false }, + ); + + const models: string[] = []; + for (const route of ROUTES) { + const upstream = await startOpenAiUpstream(route.upstream); + upstreams.push(upstream); + const display = `late-${route.name}`; + const pk = await seed.createProviderKey({ + display_name: `${display}-pk`, + secret: "sk-mock", + api_base: route.provider === "anthropic" ? upstream.baseUrl : `${upstream.baseUrl}/v1`, + ...(route.apis ? { apis: route.apis } : {}), + }); + const model = await seed.createModel({ + display_name: display, + provider: route.provider, + model_name: + route.provider === "anthropic" + ? "claude-3-5-haiku-20241022" + : route.name === "transcription" + ? "gpt-4o-transcribe" + : "gpt-4o-mini", + provider_key_id: pk.id, + }); + await seed.attachGuardrailToModel(monitor.id as string, model.id as string); + models.push(display); + } + + // Seeded last: its key authenticating implies the whole seed is live. + await seed.createApiKey({ key_hash: hash(CALLER), allowed_models: models }); + await waitConfigPropagation( + async () => (await new ProxyClient(app!.proxyUrl, CALLER).listModels()).status === 200, + ); + }, 120_000); + + afterAll(async () => { + await app?.exit(); + await sls?.close(); + await Promise.all(upstreams.map((u) => u.close())); + }); + + for (const route of ROUTES) { + test(`${route.name}: a trigger past 256 KiB of output is a monitor hit`, async (ctx) => { + if (!etcdReachable || !app || !sls) { + ctx.skip(); + return; + } + const model = `late-${route.name}`; + const res = await route.send(app.proxyUrl, model); + expect(res.status).toBe(200); + const body = await res.text(); + expect(body, "monitor mode releases the whole stream").toContain(MARKER); + expect(body).not.toContain("blocked by content policy"); + expect(body).not.toContain("output_buffer_exceeded"); + const event = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("requested_model") === model, + `usage event for ${model}`, + ); + const hits = JSON.parse(event.get("guardrail_monitor_hits") ?? "[]") as Array<{ + action: string; + hook: string; + guardrail_name: string; + }>; + expect(hits).toContainEqual( + expect.objectContaining({ action: "would_block", hook: "output", guardrail_name: GUARD_NAME }), + ); + }); + } +}); From 8403bb9bda282202cb4763eec76493e01c18aa1b Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 01:44:49 +0000 Subject: [PATCH 3/8] docs(guardrails): max_buffer_bytes also caps chat streaming tool-call arguments In window mode the chat route holds streamed tool-call arguments whole and caps them with the chain's folded max_buffer_bytes; the field description named only the whole-stream holds on /v1/messages and /v1/responses. Schemas regenerated. --- crates/aisix-core/src/models/guardrail.rs | 20 +++++++++++-------- .../resources-lenient/guardrail.schema.json | 8 ++++---- schemas/resources/guardrail.schema.json | 8 ++++---- 3 files changed, 20 insertions(+), 16 deletions(-) diff --git a/crates/aisix-core/src/models/guardrail.rs b/crates/aisix-core/src/models/guardrail.rs index 029aaf453..09e8f6321 100644 --- a/crates/aisix-core/src/models/guardrail.rs +++ b/crates/aisix-core/src/models/guardrail.rs @@ -245,8 +245,9 @@ pub struct AzureContentSafetyTextModerationConfig { pub window_overlap_size: u32, /// Max bytes of model-generated content held back from a streamed response /// before `on_buffer_exceeded` applies: in `buffer_full` mode, and in - /// `window` mode where the stream is held whole (`/v1/messages`, - /// `/v1/responses`) and no output guardrail in the chain uses + /// `window` mode wherever output is held whole (the whole stream on + /// `/v1/messages` and `/v1/responses`, and tool-call arguments on + /// `/v1/chat/completions`) and no output guardrail in the chain uses /// `buffer_full`, whose rows alone then set the cap. Counts assistant /// text, reasoning, and tool-call arguments; SSE and JSON framing is not /// counted. @@ -361,8 +362,9 @@ pub struct AliyunTextModerationConfig { pub window_overlap_size: u32, /// Max bytes of model-generated content held back from a streamed response /// before `on_buffer_exceeded` applies: in `buffer_full` mode, and in - /// `window` mode where the stream is held whole (`/v1/messages`, - /// `/v1/responses`) and no output guardrail in the chain uses + /// `window` mode wherever output is held whole (the whole stream on + /// `/v1/messages` and `/v1/responses`, and tool-call arguments on + /// `/v1/chat/completions`) and no output guardrail in the chain uses /// `buffer_full`, whose rows alone then set the cap. Counts assistant /// text, reasoning, and tool-call arguments; SSE and JSON framing is not /// counted. @@ -447,8 +449,9 @@ pub struct AliyunAiGuardrailConfig { pub window_overlap_size: u32, /// Max bytes of model-generated content held back from a streamed response /// before `on_buffer_exceeded` applies: in `buffer_full` mode, and in - /// `window` mode where the stream is held whole (`/v1/messages`, - /// `/v1/responses`) and no output guardrail in the chain uses + /// `window` mode wherever output is held whole (the whole stream on + /// `/v1/messages` and `/v1/responses`, and tool-call arguments on + /// `/v1/chat/completions`) and no output guardrail in the chain uses /// `buffer_full`, whose rows alone then set the cap. Counts assistant /// text, reasoning, and tool-call arguments; SSE and JSON framing is not /// counted. @@ -1016,8 +1019,9 @@ pub struct CustomConfig { pub window_overlap_size: u32, /// Max bytes of model-generated content held back from a streamed response /// before `on_buffer_exceeded` applies: in `buffer_full` mode, and in - /// `window` mode where the stream is held whole (`/v1/messages`, - /// `/v1/responses`) and no output guardrail in the chain uses + /// `window` mode wherever output is held whole (the whole stream on + /// `/v1/messages` and `/v1/responses`, and tool-call arguments on + /// `/v1/chat/completions`) and no output guardrail in the chain uses /// `buffer_full`, whose rows alone then set the cap. Counts assistant /// text, reasoning, and tool-call arguments; SSE and JSON framing is not /// counted. diff --git a/schemas/resources-lenient/guardrail.schema.json b/schemas/resources-lenient/guardrail.schema.json index 45fc43425..69a4175c8 100644 --- a/schemas/resources-lenient/guardrail.schema.json +++ b/schemas/resources-lenient/guardrail.schema.json @@ -629,7 +629,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" @@ -797,7 +797,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" @@ -944,7 +944,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" @@ -1677,7 +1677,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" diff --git a/schemas/resources/guardrail.schema.json b/schemas/resources/guardrail.schema.json index 571f10070..6d31921f2 100644 --- a/schemas/resources/guardrail.schema.json +++ b/schemas/resources/guardrail.schema.json @@ -641,7 +641,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" @@ -810,7 +810,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" @@ -958,7 +958,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" @@ -1755,7 +1755,7 @@ }, "max_buffer_bytes": { "default": 262144, - "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode where the stream is held whole (`/v1/messages`, `/v1/responses`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", + "description": "Max bytes of model-generated content held back from a streamed response before `on_buffer_exceeded` applies: in `buffer_full` mode, and in `window` mode wherever output is held whole (the whole stream on `/v1/messages` and `/v1/responses`, and tool-call arguments on `/v1/chat/completions`) and no output guardrail in the chain uses `buffer_full`, whose rows alone then set the cap. Counts assistant text, reasoning, and tool-call arguments; SSE and JSON framing is not counted.", "format": "uint64", "minimum": 1.0, "type": "integer" From 3910aa0543593221f23ab5e2c9e3abcfe2c011bd Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 01:44:56 +0000 Subject: [PATCH 4/8] fix(usage): every model-serving request emits its UsageEvent A usage event is the observability record of a request, and the console Logs and budgets read nothing else. Several endpoints still skipped it when there was nothing to bill: - /v1/rerank emitted only when the upstream reported a token count. A live Cohere rerank-v3.5 call reports meta.billed_units.search_units and no input_tokens, so every real Cohere rerank was missing from Logs and budgets. - /v1/completions, /v1/embeddings, /v1/images/generations and POST /v1/videos skipped the event when the gateway answered 501 itself because the provider lacks the capability. Each now emits, at zero tokens when the upstream reported none. The rule is stated on usage_attr::emit_usage, the emission chokepoint. No search-unit pricing is added. --- crates/aisix-proxy/src/completions.rs | 59 +++---- crates/aisix-proxy/src/embeddings.rs | 46 ++---- crates/aisix-proxy/src/images.rs | 39 ++--- crates/aisix-proxy/src/rerank.rs | 68 ++++---- crates/aisix-proxy/src/responses.rs | 6 +- crates/aisix-proxy/src/usage_attr.rs | 90 ++--------- crates/aisix-proxy/src/videos.rs | 40 +++-- .../usage-event-every-request-e2e.test.ts | 153 ++++++++++++++++++ 8 files changed, 279 insertions(+), 222 deletions(-) create mode 100644 tests/e2e/src/cases/usage-event-every-request-e2e.test.ts diff --git a/crates/aisix-proxy/src/completions.rs b/crates/aisix-proxy/src/completions.rs index f828524eb..9cd548481 100644 --- a/crates/aisix-proxy/src/completions.rs +++ b/crates/aisix-proxy/src/completions.rs @@ -209,18 +209,11 @@ pub async fn completions( elapsed, ); // Issue #403: emit UsageEvent so cp-api's budget ledger - // and customer-facing /logs see /v1/completions spend. - // Pre-#403 the legacy completions handler dropped the - // event entirely. A 501 or malformed 200 normally remains - // suppressed, but a guardrail decision is an audit fact rather - // than token-accounting noise and gets a zero-token event. - let guardrail_attributed = - crate::usage_attr::has_guardrail_attribution(&audit, &success.monitor_hits); + // and customer-facing /logs see /v1/completions traffic. The + // route's own 501 records zero tokens. // One zero-token event per attempt that failed first (#655). - // The route's own 501 keeps its gated event below; when that - // stays silent the last failed attempt is the terminal one. - let (superseded, refused) = crate::usage_attr::split_route_refusal(&routing); - let answered = success.usage.is_some() || guardrail_attributed; + // The route's own 501 is the terminal event below. + let superseded = crate::usage_attr::split_route_refusal(&routing); crate::usage_attr::emit_failed_attempts( &state, &snapshot, @@ -231,13 +224,13 @@ pub async fn completions( &client, &success.applied_guardrails, superseded, - refused && !answered, + false, success.guardrail_blocked, success.monitor_hits.clone(), success.redactions.clone(), &audit, ); - if answered { + { let usage = success.usage.as_ref().unwrap_or(&CompletionUsage { prompt_tokens: 0, completion_tokens: 0, @@ -766,8 +759,7 @@ async fn dispatch( provider_key_id: pk_id.clone(), upstream_model: model.upstream_model().unwrap_or("unknown").to_string(), applied_guardrails, - // No upstream call → no token usage. The handler emits only - // if screening already produced guardrail attribution. + // No upstream call → no token usage; the event records zero. usage: None, upstream_called: false, provider_request_id: String::new(), @@ -785,9 +777,7 @@ async fn dispatch( /// - The `usage` block is missing entirely (non-conformant edge), or /// - `usage.prompt_tokens` is missing / non-numeric (malformed) /// -/// Those cases normally skip UsageEvent emission rather than attributing a -/// zero-everything noise row to the api_key. A guardrail decision overrides -/// that suppression so its audit fields are not lost. +/// The caller then estimates the counts locally (AISIX-Cloud#1074). /// /// `completion_tokens`, by contrast, defaults to 0 when absent: a 200 /// that reports a prompt side but omits the completion side is still a @@ -1619,15 +1609,14 @@ mod tests { } } - /// A 501 without a guardrail decision stays out of usage, while a 501 - /// reached after a mask must preserve that attribution in a zero-token - /// event (#1083). Triggers the path + /// A 501 emits a zero-token event, and one reached after a mask carries + /// that attribution (#1083). Triggers the path /// by routing /v1/completions at an Anthropic-backed model; /// `AnthropicBridge` doesn't override `Bridge::complete()`, so the trait /// default returns `UnsupportedCapability(TextCompletions)`, which maps to /// 501. #[tokio::test] - async fn provider_lacking_complete_emits_only_for_guardrail_attribution() { + async fn provider_lacking_complete_emits_a_zero_token_event() { use aisix_obs::UsageSink; use aisix_provider_anthropic::AnthropicBridge; @@ -1670,14 +1659,13 @@ mod tests { (default Bridge::complete returns BridgeError::Config)", ); - let recv = tokio::time::timeout(std::time::Duration::from_millis(200), rx.recv()).await; - if let Ok(Some(ev)) = recv { - panic!( - "501 NotImplemented must not emit UsageEvent, \ - got prompt_tokens={}, status_code={}", - ev.prompt_tokens, ev.status_code, - ); - } + let ev = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) + .await + .expect("the 501 must emit its zero-token UsageEvent") + .expect("usage sink remains open"); + assert_eq!(ev.status_code, 501); + assert_eq!((ev.prompt_tokens, ev.completion_tokens), (0, 0)); + assert!(ev.guardrail_enforced_hits.is_empty(), "{ev:?}"); let body = serde_json::json!({ "model": "claude-instruct", @@ -1698,14 +1686,9 @@ mod tests { assert_eq!(ev.applied_guardrails.len(), 1); } - /// The same 501 path, but the guardrail FAILS OPEN instead of masking. - /// - /// A bypass leaves no enforced hit and no score, so before the gate - /// learned about it this event was suppressed outright — the reason was - /// written onto a row nobody received, which is the same silence this - /// field exists to break, one layer further out. The unbilled paths are - /// where it bites: `success.usage` is `None`, so the guardrail - /// attribution is the only thing that can keep the row alive. + /// The same 501 path, but the guardrail FAILS OPEN instead of masking: + /// the bypass reason, which leaves no enforced hit and no score, still + /// lands on the unbilled event. #[tokio::test] async fn a_fail_open_bypass_alone_keeps_the_unbilled_event_alive() { use aisix_obs::UsageSink; diff --git a/crates/aisix-proxy/src/embeddings.rs b/crates/aisix-proxy/src/embeddings.rs index bf9b84d59..1b29cfbc7 100644 --- a/crates/aisix-proxy/src/embeddings.rs +++ b/crates/aisix-proxy/src/embeddings.rs @@ -177,11 +177,8 @@ pub async fn embeddings( elapsed, ); // One zero-token event per attempt that failed first (#655). - // The route's own 501 keeps its gated event below; when that - // stays silent the last failed attempt is the terminal one. - let (superseded, refused) = crate::usage_attr::split_route_refusal(&routing); - let answered = success.upstream_called - || crate::usage_attr::has_guardrail_attribution(&audit, &success.monitor_hits); + // The route's own 501 is the terminal event below. + let superseded = crate::usage_attr::split_route_refusal(&routing); crate::usage_attr::emit_failed_attempts( &state, &snapshot, @@ -192,7 +189,7 @@ pub async fn embeddings( &client, &success.applied_guardrails, superseded, - refused && !answered, + false, false, success.monitor_hits.clone(), success.redactions.clone(), @@ -203,13 +200,9 @@ pub async fn embeddings( // spend. Pre-#226 the embedding handler dropped the // event entirely, so any /v1/embeddings traffic was // invisible to budget enforcement and billing - // reconciliation. A 501 normally remains suppressed because no - // upstream call happened; if screening recorded a guardrail - // decision, preserve it in a zero-token event. Distinguished from - // `prompt_tokens == 0` so a 200 with legitimately zero - // tokens (empty input, provider-specific billing - // convention) still emits. - if answered { + // reconciliation. The route's own 501, which made no upstream + // call, records zero tokens. + { emit_usage_event( &state, &snapshot, @@ -1743,17 +1736,13 @@ mod tests { // The `.expect(1)` on both mocks asserts exactly two upstream calls. } - /// A 501 without a guardrail decision stays out of usage, while a 501 - /// reached after a mask must preserve that attribution in a zero-token - /// event (#1083). + /// A 501 emits a zero-token event, and one reached after a mask carries + /// that attribution (#1083). /// Triggers the path by routing /v1/embeddings at an Anthropic-backed /// model; `AnthropicBridge` doesn't override `Bridge::embed()` so the - /// trait default returns `BridgeError::Config(...)` → 501. Without this - /// test, a regression flipping `upstream_called: false` → `true` (or - /// `usage: None` → `Some(zero)`) on the 501 branch would silently emit - /// a bogus zero event. + /// trait default returns `BridgeError::Config(...)` → 501. #[tokio::test] - async fn provider_lacking_embed_emits_only_for_guardrail_attribution() { + async fn provider_lacking_embed_emits_a_zero_token_event() { use aisix_obs::UsageSink; use aisix_provider_anthropic::AnthropicBridge; @@ -1796,14 +1785,13 @@ mod tests { (default Bridge::embed returns BridgeError::Config)", ); - let recv = tokio::time::timeout(std::time::Duration::from_millis(200), rx.recv()).await; - if let Ok(Some(ev)) = recv { - panic!( - "501 NotImplemented must not emit UsageEvent, \ - got prompt_tokens={}, status_code={}", - ev.prompt_tokens, ev.status_code, - ); - } + let ev = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) + .await + .expect("the 501 must emit its zero-token UsageEvent") + .expect("usage sink remains open"); + assert_eq!(ev.status_code, 501); + assert_eq!((ev.prompt_tokens, ev.completion_tokens), (0, 0)); + assert!(ev.guardrail_enforced_hits.is_empty(), "{ev:?}"); let body = serde_json::json!({ "model": "claude-embed", diff --git a/crates/aisix-proxy/src/images.rs b/crates/aisix-proxy/src/images.rs index 7c2eb8076..dd3bdf44b 100644 --- a/crates/aisix-proxy/src/images.rs +++ b/crates/aisix-proxy/src/images.rs @@ -156,11 +156,8 @@ pub async fn image_generations( elapsed, ); // One zero-token event per attempt that failed first (#655). - // The route's own 501 keeps its gated event below; when that - // stays silent the last failed attempt is the terminal one. - let (superseded, refused) = crate::usage_attr::split_route_refusal(&routing); - let answered = success.upstream_called - || crate::usage_attr::has_guardrail_attribution(&audit, &success.monitor_hits); + // The route's own 501 is the terminal event below. + let superseded = crate::usage_attr::split_route_refusal(&routing); crate::usage_attr::emit_failed_attempts( &state, &snapshot, @@ -171,7 +168,7 @@ pub async fn image_generations( &client, &success.applied_guardrails, superseded, - refused && !answered, + false, false, success.monitor_hits.clone(), success.redactions.clone(), @@ -179,14 +176,14 @@ pub async fn image_generations( ); // Issue #407: emit UsageEvent so cp-api's budget ledger + // /logs see image-generation traffic. Pre-#407 the handler - // dropped the event entirely. Emit on a real upstream call - // (even zero tokens — request visible/attributed). A 501 emits - // only to preserve a guardrail decision. Tokens come from the upstream + // dropped the event entirely. Every request emits, the route's + // own 501 included, at zero tokens when there are none to + // report. Tokens come from the upstream // `usage` block when present (gpt-image-1); dall-e-3 has no // usage block → zero tokens (precise per-image cost is a // documented cross-repo follow-up — needs image-count / // size / quality on the wire + cp-api pricing). - if answered { + { let (prompt_tokens, completion_tokens) = success.usage.unwrap_or((0, 0)); emit_usage_event( &state, @@ -1190,9 +1187,8 @@ mod tests { assert_eq!(event.inbound_protocol, "openai"); } - /// A 501 without a guardrail decision stays out of usage, while a 501 - /// reached after a mask must preserve that attribution in a zero-token - /// event (#1083). Unlike completions and embeddings, /v1/images/generations + /// A 501 emits a zero-token event, and one reached after a mask carries + /// that attribution (#1083). Unlike completions and embeddings, /v1/images/generations /// rejects non-OpenAI providers with 400 *before* dispatch (see /// `non_openai_provider_returns_400_invalid_request`) and the real /// `OpenAiBridge` overrides `generate_image`, so the only way to reach @@ -1200,7 +1196,7 @@ mod tests { /// leaves `Bridge::generate_image` at the trait default. We register a /// minimal stub under the "openai" key to exercise exactly that. #[tokio::test] - async fn image_501_emits_only_for_guardrail_attribution() { + async fn image_501_emits_a_zero_token_event() { use aisix_gateway::{ Bridge, BridgeContext, BridgeError, ChatChunkStream, ChatFormat, ChatMessage, ChatResponse, FinishReason, UsageStats, @@ -1265,14 +1261,13 @@ mod tests { (default Bridge::generate_image returns BridgeError::Config)", ); - let recv = tokio::time::timeout(std::time::Duration::from_millis(200), rx.recv()).await; - if let Ok(Some(ev)) = recv { - panic!( - "501 NotImplemented must not emit UsageEvent, \ - got prompt_tokens={}, status_code={}", - ev.prompt_tokens, ev.status_code, - ); - } + let ev = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) + .await + .expect("the 501 must emit its zero-token UsageEvent") + .expect("usage sink remains open"); + assert_eq!(ev.status_code, 501); + assert_eq!((ev.prompt_tokens, ev.completion_tokens), (0, 0)); + assert!(ev.guardrail_enforced_hits.is_empty(), "{ev:?}"); let body = serde_json::json!({ "model": "stub-image", diff --git a/crates/aisix-proxy/src/rerank.rs b/crates/aisix-proxy/src/rerank.rs index 209ce4a42..b571695ca 100644 --- a/crates/aisix-proxy/src/rerank.rs +++ b/crates/aisix-proxy/src/rerank.rs @@ -182,14 +182,10 @@ pub async fn rerank( &audit, ); // Issue #405: emit UsageEvent so cp-api's budget ledger - // and customer-facing /logs see /v1/rerank spend. - // Pre-#405 the rerank handler dropped the event entirely. - // Skip on 200 without a recognisable usage field — avoids - // attributing zero-everything noise rows when an - // upstream returns a malformed / unsupported shape. - let guardrail_attributed = - crate::usage_attr::has_guardrail_attribution(&audit, &success.monitor_hits); - if success.usage.is_some() || guardrail_attributed { + // and customer-facing /logs see /v1/rerank traffic. An upstream + // that reports no token count (Cohere bills search units) still + // gets its event, at zero tokens. + { let usage = success .usage .as_ref() @@ -680,11 +676,9 @@ async fn dispatch( // forward raw bytes downstream — preserves any provider-specific // fields (Cohere `meta.api_version`, Jina-specific fields, etc.) // that the JSON round-trip would otherwise re-format. A parse - // failure here is non-fatal: the request still succeeds, and it - // still leaves a usage event whenever a guardrail attributed the - // request — only an unattributed one goes unrecorded, which is what - // keeps a zero-everything noise row off the ledger. Audit HIGH: log - // the parse failure so a silent billing gap is visible in operator + // failure here is non-fatal: the request still succeeds and still + // leaves a usage event, at zero tokens. Audit HIGH: log the parse + // failure so a silent billing gap is visible in operator // dashboards (the upstream returned 200 + claimed JSON but the body // was unparseable — this is upstream-malformed, not gateway-bug, but // operators need to see it). @@ -776,8 +770,9 @@ async fn dispatch( /// /// Three known wire shapes (per #213): /// - **OpenAI-compat** — `usage.prompt_tokens` (or `usage.input_tokens`) -/// - **Cohere** — `meta.billed_units.input_tokens` -/// () +/// - **Cohere** — `meta.billed_units.input_tokens`, when present +/// (); its rerank models +/// report only `search_units`, which yields `None` /// - **Jina** — `usage.total_tokens` /// () /// @@ -1791,23 +1786,21 @@ mod tests { assert_eq!(event.inbound_protocol, "openai"); } - /// Issue #405: Cohere's wire shape puts the token counter at - /// `meta.billed_units.input_tokens` instead of `usage.prompt_tokens`. - /// The extractor must handle this — without coverage, customers - /// running Cohere-backed rerank would see zero spend in cp-api - /// even though billing is happening. + /// Issue #405: Cohere reports `meta.billed_units.search_units` and no + /// token count. The request must still reach the usage-event table, + /// which the console Logs and budgets read, at zero tokens. #[tokio::test] async fn emits_usage_event_on_cohere_wire_shape_issue_405() { use aisix_obs::UsageSink; let upstream = MockServer::start().await; - // Cohere wire shape: `meta.billed_units.input_tokens`. + // Cohere's live wire shape: search units only, no token count. let upstream_body = serde_json::json!({ "id": "rerank-cohere", "results": [{"index": 0, "relevance_score": 0.95}], "meta": { - "api_version": {"version": "1"}, - "billed_units": {"input_tokens": 47, "search_units": 1} + "api_version": {"version": "2"}, + "billed_units": {"search_units": 1} } }); Mock::given(method("POST")) @@ -1845,18 +1838,18 @@ mod tests { .expect("usage_sink sender dropped"); assert_eq!( - event.prompt_tokens, 47, - "Cohere meta.billed_units.input_tokens must be surfaced as prompt_tokens", + event.prompt_tokens, 0, + "Cohere reports search units, not tokens: the event records zero tokens", ); + assert_eq!(event.status_code, 200); assert_eq!(event.inbound_protocol, "openai"); } - /// Issue #405: an upstream 200 with no recognisable usage field - /// (neither `usage` nor `meta.billed_units`) must NOT emit a - /// zero-everything noise row. Same edge-case discipline as - /// PR #425 audit MEDIUM-1. + /// An upstream 200 with no recognisable usage field (neither `usage` + /// nor `meta.billed_units`) still emits its usage event, at zero + /// tokens: every request is recorded. #[tokio::test] - async fn skips_usage_event_when_upstream_lacks_usage_fields() { + async fn emits_zero_token_usage_event_when_upstream_lacks_usage_fields() { use aisix_obs::UsageSink; let upstream = MockServer::start().await; @@ -1893,14 +1886,13 @@ mod tests { .unwrap(); assert_eq!(resp.status(), StatusCode::OK); - let recv = tokio::time::timeout(std::time::Duration::from_millis(200), rx.recv()).await; - if let Ok(Some(ev)) = recv { - panic!( - "no UsageEvent should be emitted when upstream lacks usage fields, \ - got prompt_tokens={}", - ev.prompt_tokens, - ); - } + let event = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) + .await + .expect("UsageEvent must be emitted when upstream lacks usage fields") + .expect("usage_sink sender dropped"); + assert_eq!(event.prompt_tokens, 0); + assert_eq!(event.completion_tokens, 0); + assert_eq!(event.status_code, 200); } #[tokio::test] diff --git a/crates/aisix-proxy/src/responses.rs b/crates/aisix-proxy/src/responses.rs index 4a08d82c3..2e2522030 100644 --- a/crates/aisix-proxy/src/responses.rs +++ b/crates/aisix-proxy/src/responses.rs @@ -2819,9 +2819,9 @@ async fn responses_cross_provider_to_target( /// - The `usage` block is missing entirely, OR /// - `usage.input_tokens` is missing / non-numeric /// -/// Those cases skip UsageEvent emission rather than attributing a -/// zero-everything noise row to the api_key. The `input_tokens` gate -/// distinguishes "no upstream usage at all" from a legitimate reply. +/// The caller then estimates the counts locally (AISIX-Cloud#1074). The +/// `input_tokens` gate distinguishes "no upstream usage at all" from a +/// legitimate reply. /// /// `output_tokens`, by contrast, defaults to 0 when absent: a 200 that /// reports an input side but omits the output side is still a real diff --git a/crates/aisix-proxy/src/usage_attr.rs b/crates/aisix-proxy/src/usage_attr.rs index c81047234..b83be69cc 100644 --- a/crates/aisix-proxy/src/usage_attr.rs +++ b/crates/aisix-proxy/src/usage_attr.rs @@ -102,32 +102,6 @@ pub(crate) fn guardrail_scores(audit: &GuardrailAudit) -> Vec bool { - !monitor_hits.is_empty() - || audit.as_ref().is_some_and(|log| { - !log.snapshot().is_empty() - || !log.score_snapshot().is_empty() - || log.bypass_reason().is_some() - }) -} - /// [`guardrail_scores`] for the retrying families — see /// [`terminal_enforced_hits`]. pub(crate) fn terminal_guardrail_scores( @@ -762,21 +736,18 @@ pub(crate) fn failed_attempts_are_terminal(routing: &crate::attempt::RoutingTele !routing.attempts.is_empty() && routing.winner().is_none() } -/// A single-shot success branch's attempts, split for emission: the ones -/// [`emit_failed_attempts`] owns, and whether the request was answered by -/// the route's own 501 for a provider lacking the capability — the LAST -/// record, failed and never dispatched. That refusal keeps its own -/// handler-emitted event, gated as it always was (no event unless a -/// guardrail decision needs recording), so it is left out of the slice; -/// when that gate stays shut the last superseded attempt, if any, is the -/// terminal event instead. +/// A single-shot success branch's attempts that [`emit_failed_attempts`] +/// owns. When the request was answered by the route's own 501 for a +/// provider lacking the capability — the LAST record, failed and never +/// dispatched — that refusal is the handler's own terminal event, so it is +/// left out of the slice. pub(crate) fn split_route_refusal( routing: &crate::attempt::RoutingTelemetry, -) -> (&[crate::attempt::AttemptRecord], bool) { +) -> &[crate::attempt::AttemptRecord] { if failed_attempts_are_terminal(routing) { - (&routing.attempts[..routing.attempts.len() - 1], true) + &routing.attempts[..routing.attempts.len() - 1] } else { - (&routing.attempts, false) + &routing.attempts } } @@ -916,6 +887,11 @@ pub(crate) fn emit_prepared_usage_event( /// through here, so the CP telemetry leg and the exporter fan-out cannot /// drift, and the trace snapshot is taken in exactly one place. /// +/// A usage event is the observability record of a request, not a billing +/// line: every request that reaches dispatch emits one, at zero tokens when +/// the upstream reported none or nothing was billable. Never skip it for +/// lack of usage. +/// /// `trace` is the request's bundle (`ClientContext::trace`); `terminal` /// says whether this event ends the request — the terminal event carries /// the SERVER + logical spans (ending the SERVER span at NOW, the real @@ -1002,46 +978,6 @@ pub(crate) fn emit_usage( mod tests { use super::*; - #[test] - fn similarity_score_alone_requires_a_zero_token_event() { - let log = Arc::new(aisix_guardrails::GuardrailAuditLog::new()); - let audit = Some(Arc::clone(&log)); - assert!(!has_guardrail_attribution(&audit, &[])); - - log.record_score(aisix_core::GuardrailScore { - guardrail_name: "semantic-policy".into(), - hook: "input".into(), - direction: "deny".into(), - score: 0.7, - threshold: 0.8, - matched: false, - top_example_index: 0, - embedding_model: "embedder".into(), - }); - assert!(has_guardrail_attribution(&audit, &[])); - } - - /// A bypass leaves no enforced hit and no score — it is the outcome - /// where no policy acted — so it has to be named here explicitly. - /// Without it, the unbilled paths that consult this gate - /// (`/v1/completions`, `/v1/embeddings`, `/v1/images/*`, `/v1/rerank`) - /// suppress the whole event, and the reason lands on a row nobody - /// receives. - #[test] - fn a_bypass_alone_requires_a_zero_token_event() { - let log = Arc::new(aisix_guardrails::GuardrailAuditLog::new()); - let audit = Some(Arc::clone(&log)); - assert!(!has_guardrail_attribution(&audit, &[])); - - log.record_bypass("lakera_timeout"); - assert!( - log.snapshot().is_empty() && log.score_snapshot().is_empty(), - "premise: a bypass is not an enforced hit and not a score, so the \ - other two arms of this gate cannot be what carries it", - ); - assert!(has_guardrail_attribution(&audit, &[])); - } - #[test] fn metric_model_label_three_outcomes() { use aisix_core::resource::ResourceEntry; diff --git a/crates/aisix-proxy/src/videos.rs b/crates/aisix-proxy/src/videos.rs index ec989e0dc..42a3ec480 100644 --- a/crates/aisix-proxy/src/videos.rs +++ b/crates/aisix-proxy/src/videos.rs @@ -1631,10 +1631,8 @@ pub async fn create_video( .unwrap_or_else(|| crate::usage_attr::UNRESOLVED_MODEL_LABEL.to_string()); telemetry.finish_routed(status, &success.provider, &model_label, None, &routing); // One zero-token event per attempt that failed first (#655). - // The route's own 501 keeps its gated event below; when that - // stays silent the last failed attempt is the terminal one. - let (superseded, refused) = crate::usage_attr::split_route_refusal(&routing); - let answered = success.upstream_called; + // The route's own 501 is the terminal event below. + let superseded = crate::usage_attr::split_route_refusal(&routing); crate::usage_attr::emit_failed_attempts( &state, &snapshot, @@ -1645,7 +1643,7 @@ pub async fn create_video( &client, &success.applied_guardrails, superseded, - refused && !answered, + false, false, success.monitor_hits.clone(), crate::redact::RedactionCounts::new(), @@ -1655,9 +1653,9 @@ pub async fn create_video( // /logs and the budget ledger like every other endpoint. // Per-second cost is computed control-plane-side once the // per-second cost schema lands (AISIX-Cloud#1118 decision 2); - // token fields stay zero. Skipped when no upstream call - // happened (the 501 unsupported-provider branch). - if answered { + // token fields stay zero. The 501 unsupported-provider branch, + // which made no upstream call, emits too. + { emit_submit_usage_event( &state, &snapshot, @@ -1740,9 +1738,6 @@ struct CreateSuccess { provider_key_id: String, applied_guardrails: Vec, monitor_hits: Vec, - /// `false` on the 501 unsupported-provider branch — no upstream call - /// happened, so no UsageEvent is attributed (embeddings convention). - upstream_called: bool, } async fn dispatch_create( @@ -1786,7 +1781,6 @@ async fn dispatch_create( provider_key_id: String::new(), applied_guardrails: Vec::new(), monitor_hits: Vec::new(), - upstream_called: false, }); } } else { @@ -1933,7 +1927,6 @@ async fn dispatch_create( provider_key_id: String::new(), applied_guardrails, monitor_hits, - upstream_called: false, }) } }; @@ -1972,7 +1965,6 @@ async fn dispatch_create( provider_key_id: target.pk_id.clone(), applied_guardrails, monitor_hits, - upstream_called: true, }) } @@ -2998,9 +2990,21 @@ mod tests { assert_eq!(resp.status(), StatusCode::NOT_FOUND); } + /// The route's own 501 still records the request, at zero tokens. #[tokio::test] async fn unsupported_provider_returns_501_not_implemented() { - let app = build_app(new_snap("http://unused", "minimax", "")); + use aisix_obs::UsageSink; + + let (tx, mut rx) = tokio::sync::mpsc::channel(8); + let app = crate::build_router( + crate::ProxyState::new( + SnapshotHandle::new(new_snap("http://unused", "minimax", "")), + Arc::new(Hub::new()), + &cfg(), + ) + .without_cache() + .with_usage_sink(UsageSink::new(tx)), + ); let resp = tower::ServiceExt::oneshot( app, post_videos(serde_json::json!({"model": "my-video", "prompt": "hi"})), @@ -3010,6 +3014,12 @@ mod tests { assert_eq!(resp.status(), StatusCode::NOT_IMPLEMENTED); let v = body_json(resp).await; assert_eq!(v["error"]["type"], "not_implemented"); + let ev = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) + .await + .expect("the 501 must emit its zero-token UsageEvent") + .expect("usage_sink sender dropped"); + assert_eq!(ev.status_code, 501); + assert_eq!((ev.prompt_tokens, ev.completion_tokens), (0, 0)); } #[tokio::test] diff --git a/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts b/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts new file mode 100644 index 000000000..a846e3427 --- /dev/null +++ b/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts @@ -0,0 +1,153 @@ +import { createHash } from "node:crypto"; +import { afterAll, beforeAll, describe, expect, test } from "vitest"; +import { + EtcdClient, + SeedClient, + spawnApp, + startMockSls, + startOpenAiUpstream, + waitConfigPropagation, + waitForSlsLog, + type MockSls, + type OpenAiUpstream, + type SpawnedApp, +} from "../harness/index.js"; + +// E2E: every model-serving request leaves a usage event — the record the +// console Logs and budgets read — even when there is nothing to bill: +// - a Cohere rerank, whose upstream reports search units and no tokens; +// - a request the gateway answers 501 itself because the provider lacks +// the capability, so no upstream call is made. +// Each records zero tokens. + +const CALLER = "sk-usage-event-every-request"; +const CALLER_HASH = createHash("sha256").update(CALLER).digest("hex"); +const CREDENTIAL_REF = "mock"; +const LOGSTORE = "usage-event-every-request"; +const auth = { authorization: `Bearer ${CALLER}`, "content-type": "application/json" }; + +// Cohere `rerank-v3.5`'s live response shape: no `input_tokens`. +const COHERE_RERANK = { + id: "rerank-cohere-live", + results: [ + { index: 1, relevance_score: 0.91 }, + { index: 0, relevance_score: 0.12 }, + ], + meta: { api_version: { version: "2" }, billed_units: { search_units: 1 } }, +}; + +describe("a usage event for every request", () => { + let app: SpawnedApp | undefined; + let sls: MockSls | undefined; + let cohere: OpenAiUpstream | undefined; + let etcdReachable = false; + + beforeAll(async () => { + const etcd = new EtcdClient(); + etcdReachable = await etcd.ping(); + if (!etcdReachable) return; + + sls = await startMockSls(); + cohere = await startOpenAiUpstream({ nonStreamBody: COHERE_RERANK }); + app = await spawnApp({ + extraEnv: { + [`SLS_CRED_${CREDENTIAL_REF.toUpperCase()}_AK_ID`]: "mock-akid", + [`SLS_CRED_${CREDENTIAL_REF.toUpperCase()}_AK_SECRET`]: "mock-secret", + }, + }); + const seed = new SeedClient(etcd, app.etcdPrefix); + await seed.createObservabilityExporter({ + name: "sls-usage-event-every-request", + enabled: true, + kind: "aliyun_sls", + endpoint: sls.url, + project: "aisix-e2e-obs", + logstore: LOGSTORE, + credential_ref: CREDENTIAL_REF, + content_mode: "metadata_only", + }); + + const coherePk = await seed.createProviderKey({ + display_name: "every-request-cohere", + secret: "cohere-mock-key", + provider: "cohere", + adapter: "openai", + api_base: cohere.baseUrl, + }); + await seed.createModel({ + display_name: "every-request-rerank", + provider: "cohere", + model_name: "rerank-v3.5", + provider_key_id: coherePk.id, + }); + + // Anthropic serves neither legacy completions nor embeddings; the + // gateway answers those itself, without an upstream call. + const anthropicPk = await seed.createProviderKey({ + display_name: "every-request-anthropic", + secret: "sk-ant-mock", + provider: "anthropic", + api_base: "http://127.0.0.1:9", + }); + await seed.createModel({ + display_name: "every-request-claude", + provider: "anthropic", + model_name: "claude-3-5-haiku-20241022", + provider_key_id: anthropicPk.id, + }); + + // Seeded last: its key authenticating implies the whole seed is live. + await seed.createApiKey({ key_hash: CALLER_HASH, allowed_models: ["*"] }); + await waitConfigPropagation(async () => { + const res = await fetch(`${app!.proxyUrl}/v1/models`, { headers: auth }); + await res.arrayBuffer(); + return res.status === 200; + }); + }, 60_000); + + afterAll(async () => { + await app?.exit(); + await cohere?.close(); + await sls?.close(); + }); + + const post = async (path: string, body: Record) => { + const res = await fetch(`${app!.proxyUrl}${path}`, { method: "POST", headers: auth, body: JSON.stringify(body) }); + return { status: res.status, body: await res.text() }; + }; + + const eventFor = (model: string, operation: string) => + waitForSlsLog( + sls!, + LOGSTORE, + (log) => log.get("requested_model") === model && log.get("operation") === operation, + `${operation} usage event for ${model}`, + ); + + test("a Cohere rerank that reports only search units records a zero-token event", async (ctx) => { + if (!etcdReachable || !app || !sls) return ctx.skip(); + const res = await post("/v1/rerank", { model: "every-request-rerank", query: "q", documents: ["a", "b"] }); + expect(res.status, res.body).toBe(200); + expect(JSON.parse(res.body).results).toHaveLength(2); + + const event = await eventFor("every-request-rerank", "rerank"); + expect(event.get("status_code")).toBe("200"); + expect(event.get("prompt_tokens") ?? "0").toBe("0"); + expect(event.get("completion_tokens") ?? "0").toBe("0"); + }); + + for (const [path, operation, body] of [ + ["/v1/completions", "completions", { model: "every-request-claude", prompt: "hi" }], + ["/v1/embeddings", "embeddings", { model: "every-request-claude", input: "hi" }], + ] as const) { + test(`${path}: the gateway's own 501 records a zero-token event`, async (ctx) => { + if (!etcdReachable || !app || !sls) return ctx.skip(); + const res = await post(path, body); + expect(res.status, res.body).toBe(501); + + const event = await eventFor("every-request-claude", operation); + expect(event.get("status_code")).toBe("501"); + expect(event.get("prompt_tokens") ?? "0").toBe("0"); + }); + } +}); From b9e11f2217db7b815c3606410f3999f250362f54 Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 01:53:08 +0000 Subject: [PATCH 5/8] docs(usage): drop comments describing the old 501 emission gate --- crates/aisix-proxy/src/completions.rs | 9 ++++----- crates/aisix-proxy/src/images.rs | 4 ++-- 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/crates/aisix-proxy/src/completions.rs b/crates/aisix-proxy/src/completions.rs index 9cd548481..7d77df007 100644 --- a/crates/aisix-proxy/src/completions.rs +++ b/crates/aisix-proxy/src/completions.rs @@ -43,8 +43,7 @@ struct CompletionDispatchSuccess { /// UUID of the resolved Model row — required for UsageEvent /// `model_id`. Always populated on every success arm (including /// the 501 NotImplemented branch where no upstream call - /// happened); emission depends on usage or a recorded guardrail - /// decision, not this field. + /// happened); every arm emits a usage event. model_id: String, /// Resolved ProviderKey UUID — feeds per-PK telemetry attribution /// (AISIX-Cloud#867 parity). @@ -61,8 +60,8 @@ struct CompletionDispatchSuccess { provider_request_id: String, /// Upstream-reported token counts. `None` on the 501 /// NotImplemented path (provider doesn't support completions) - /// or on a 200 with no `usage` block (rare edge). Those paths still - /// emit a zero-token event when a guardrail recorded a decision. + /// or on a 200 with no `usage` block (rare edge). Those paths emit a + /// zero-token event. usage: Option, /// Whether the request reached the provider. False only for the 501 /// provider-unsupported branch. @@ -1543,7 +1542,7 @@ mod tests { /// ONE zero-token UsageEvent so the failed request is visible in Logs /// (status + error class) and attributed to the api_key — instead of being /// dropped, as the non-chat handlers used to do. The 501 NotImplemented - /// path still emits nothing (no upstream call); see the test below. + /// path emits a zero-token event too; see the test below. #[tokio::test] async fn upstream_5xx_emits_zero_token_error_event() { use aisix_obs::UsageSink; diff --git a/crates/aisix-proxy/src/images.rs b/crates/aisix-proxy/src/images.rs index dd3bdf44b..9c8914a15 100644 --- a/crates/aisix-proxy/src/images.rs +++ b/crates/aisix-proxy/src/images.rs @@ -49,8 +49,8 @@ struct ImageDispatchSuccess { /// event so the request is visible + attributed. usage: Option<(u32, u32)>, /// `false` on the 501 NotImplemented branch (provider lacks image - /// generation → no upstream call). That path emits only when screening - /// already produced guardrail attribution. + /// generation → no upstream call). That path still emits a zero-token + /// event. upstream_called: bool, /// Per-detector PII mask counts (#932/#696) applied to the prompt. /// Attached to the emitted UsageEvent. Empty = no redaction. From 4f1c5ceb575506235da9d93b3cea3021366dd5e3 Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 02:03:09 +0000 Subject: [PATCH 6/8] fix(mcp): record a tool call whose upstream result cannot be read back When a guardrail or content capture makes /mcp read the tool result back and that read fails (the result outgrows the body cap, or the body is broken), the gateway answered 502 without a usage event, although the call had reached the upstream. It now emits the same tool-call event every other exit does, with status 502. Video job polling (GET /v1/videos/:id) and content retrieval deliberately emit no usage event; the emission rule now says so. --- crates/aisix-proxy/src/mcp.rs | 23 +++++- crates/aisix-proxy/src/usage_attr.rs | 3 +- crates/aisix-proxy/src/videos.rs | 4 ++ .../usage-event-every-request-e2e.test.ts | 70 ++++++++++++++++++- 4 files changed, 95 insertions(+), 5 deletions(-) diff --git a/crates/aisix-proxy/src/mcp.rs b/crates/aisix-proxy/src/mcp.rs index a594a0a78..366b021e6 100644 --- a/crates/aisix-proxy/src/mcp.rs +++ b/crates/aisix-proxy/src/mcp.rs @@ -826,7 +826,28 @@ async fn dispatch( { Ok(bytes) => bytes, Err(_) => { - return (StatusCode::BAD_GATEWAY, "invalid upstream response").into_response() + // The tool call reached the upstream: record it like every + // other exit, with the status the caller receives. + if is_tool_call { + emit_tool_call_usage( + state, + &snapshot, + &auth, + request_id, + &mcp_server, + &mcp_tool, + StatusCode::BAD_GATEWAY.as_u16(), + latency, + false, + monitor_hits, + redaction_counts, + guardrail_chain.as_ref(), + None, + trace, + /* dispatched */ true, + ); + } + return (StatusCode::BAD_GATEWAY, "invalid upstream response").into_response(); } }; if let Some(chain) = &guardrail_chain { diff --git a/crates/aisix-proxy/src/usage_attr.rs b/crates/aisix-proxy/src/usage_attr.rs index b83be69cc..75b4908d0 100644 --- a/crates/aisix-proxy/src/usage_attr.rs +++ b/crates/aisix-proxy/src/usage_attr.rs @@ -890,7 +890,8 @@ pub(crate) fn emit_prepared_usage_event( /// A usage event is the observability record of a request, not a billing /// line: every request that reaches dispatch emits one, at zero tokens when /// the upstream reported none or nothing was billable. Never skip it for -/// lack of usage. +/// lack of usage. The one exception is by design: polling a video job +/// (`GET /v1/videos/:id`) and retrieving its content emit none. /// /// `trace` is the request's bundle (`ClientContext::trace`); `terminal` /// says whether this event ends the request — the terminal event carries diff --git a/crates/aisix-proxy/src/videos.rs b/crates/aisix-proxy/src/videos.rs index 42a3ec480..a664cfafd 100644 --- a/crates/aisix-proxy/src/videos.rs +++ b/crates/aisix-proxy/src/videos.rs @@ -2024,6 +2024,8 @@ fn group_routes_to(snapshot: &aisix_core::AisixSnapshot, group: &str, target: &s }) } +/// Polls a video job. Deliberately emits no usage event: polling a job is +/// not a recorded request (see `usage_attr::emit_usage`). pub async fn get_video( State(state): State, auth: AuthenticatedKey, @@ -2096,6 +2098,8 @@ pub async fn get_video( } } +/// Retrieves a finished video. Deliberately emits no usage event, like +/// [`get_video`]. pub async fn video_content( State(state): State, auth: AuthenticatedKey, diff --git a/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts b/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts index a846e3427..e8f2741a6 100644 --- a/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts +++ b/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts @@ -1,13 +1,15 @@ -import { createHash } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import { afterAll, beforeAll, describe, expect, test } from "vitest"; import { EtcdClient, SeedClient, spawnApp, + startMcpUpstream, startMockSls, startOpenAiUpstream, waitConfigPropagation, waitForSlsLog, + type McpUpstream, type MockSls, type OpenAiUpstream, type SpawnedApp, @@ -17,8 +19,13 @@ import { // console Logs and budgets read — even when there is nothing to bill: // - a Cohere rerank, whose upstream reports search units and no tokens; // - a request the gateway answers 501 itself because the provider lacks -// the capability, so no upstream call is made. +// the capability, so no upstream call is made; +// - an MCP tool call whose upstream result the gateway cannot read back +// (it outgrows the body cap), answered 502 after the call went out. // Each records zero tokens. +// +// The body cap is lowered so the MCP tool result outgrows it. +const BODY_LIMIT = 65_536; const CALLER = "sk-usage-event-every-request"; const CALLER_HASH = createHash("sha256").update(CALLER).digest("hex"); @@ -40,6 +47,7 @@ describe("a usage event for every request", () => { let app: SpawnedApp | undefined; let sls: MockSls | undefined; let cohere: OpenAiUpstream | undefined; + let mcp: McpUpstream | undefined; let etcdReachable = false; beforeAll(async () => { @@ -49,7 +57,11 @@ describe("a usage event for every request", () => { sls = await startMockSls(); cohere = await startOpenAiUpstream({ nonStreamBody: COHERE_RERANK }); + mcp = await startMcpUpstream("big", { + reportContent: { summary: "done", log: "x".repeat(4 * BODY_LIMIT), structuredLog: "ok" }, + }); app = await spawnApp({ + requestBodyLimitBytes: BODY_LIMIT, extraEnv: { [`SLS_CRED_${CREDENTIAL_REF.toUpperCase()}_AK_ID`]: "mock-akid", [`SLS_CRED_${CREDENTIAL_REF.toUpperCase()}_AK_SECRET`]: "mock-secret", @@ -96,8 +108,30 @@ describe("a usage event for every request", () => { provider_key_id: anthropicPk.id, }); + // A guardrail on the MCP server makes the gateway read the tool result + // back before relaying it. + const mcpServerId = randomUUID(); + await seed.update("mcp_servers", mcpServerId, { display_name: "big", url: mcp.url, enabled: true }); + const monitor = await seed.createGuardrail( + { + name: "every-request-mcp-monitor", + enabled: true, + hook_point: "output", + enforcement_mode: "monitor", + kind: "keyword", + patterns: [{ kind: "literal", value: "never-present-literal" }], + }, + { attach: false }, + ); + await seed.update("guardrail_attachments", randomUUID(), { + guardrail_id: monitor.id, + scope_type: "mcp_server", + scope_id: mcpServerId, + priority: 100, + }); + // Seeded last: its key authenticating implies the whole seed is live. - await seed.createApiKey({ key_hash: CALLER_HASH, allowed_models: ["*"] }); + await seed.createApiKey({ key_hash: CALLER_HASH, allowed_models: ["*"], mcp_access: { allow: ["*"] } }); await waitConfigPropagation(async () => { const res = await fetch(`${app!.proxyUrl}/v1/models`, { headers: auth }); await res.arrayBuffer(); @@ -108,6 +142,7 @@ describe("a usage event for every request", () => { afterAll(async () => { await app?.exit(); await cohere?.close(); + await mcp?.close(); await sls?.close(); }); @@ -150,4 +185,33 @@ describe("a usage event for every request", () => { expect(event.get("prompt_tokens") ?? "0").toBe("0"); }); } + + test("an MCP tool result the gateway cannot read back records a 502 event", async (ctx) => { + if (!etcdReachable || !app || !sls) return ctx.skip(); + const rpc = (id: number, method: string, params: Record) => + fetch(`${app!.proxyUrl}/mcp`, { + method: "POST", + headers: { ...auth, accept: "application/json, text/event-stream" }, + body: JSON.stringify({ jsonrpc: "2.0", id, method, params }), + }); + await ( + await rpc(1, "initialize", { + protocolVersion: "2025-11-25", + capabilities: {}, + clientInfo: { name: "usage-event-every-request", version: "0.1" }, + }) + ).text(); + const res = await rpc(2, "tools/call", { name: "big__report", arguments: {} }); + const body = await res.text(); + expect(res.status, body).toBe(502); + + const event = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("operation") === "mcp" && log.get("mcp_tool_name") === "report", + "mcp usage event for big__report", + ); + expect(event.get("status_code")).toBe("502"); + expect(event.get("mcp_server_name")).toBe("big"); + }); }); From a729bc2b00475065923422c4bb3b69a59a61d988 Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 02:08:25 +0000 Subject: [PATCH 7/8] test(e2e): size the unterminated-tail fixture past the 128x raw bound The #1029 passthrough fixture carried a 100 KB content-free tail, past the old 64x raw bound of its 1 000-byte cap but under the new 128x one, so it no longer tripped. It now carries 200 KB. --- .../src/cases/guardrail-buffer-cap-enforced-hit-e2e.test.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/e2e/src/cases/guardrail-buffer-cap-enforced-hit-e2e.test.ts b/tests/e2e/src/cases/guardrail-buffer-cap-enforced-hit-e2e.test.ts index 1effd4e92..621897dbb 100644 --- a/tests/e2e/src/cases/guardrail-buffer-cap-enforced-hit-e2e.test.ts +++ b/tests/e2e/src/cases/guardrail-buffer-cap-enforced-hit-e2e.test.ts @@ -123,12 +123,12 @@ const RESPONSES_STREAM = [ // A passthrough stream whose last frame never ends: a short answer, then a // keep-alive the upstream leaves unterminated, carrying no content but past -// the raw bound (64 × the tight cap) the hold-back also keeps. +// the raw bound (128 × the tight cap) the hold-back also keeps. const TAIL_MARKER = "tail-marker"; const UNTERMINATED_TAIL_STREAM = [ `data: ${chatChunk({ role: "assistant" })}\n\n`, `data: ${chatChunk({ content: PIECES[0] })}\n\n`, - `data: ${JSON.stringify({ id: `${TAIL_MARKER}-${"k".repeat(100_000)}`, object: "chat.completion.chunk", choices: [] })}`, + `data: ${JSON.stringify({ id: `${TAIL_MARKER}-${"k".repeat(200_000)}`, object: "chat.completion.chunk", choices: [] })}`, ]; interface EnforcedHit { From ab36d6da92218b5fb61e15d629d9129f75ecf6e0 Mon Sep 17 00:00:00 2001 From: Jarvis Date: Fri, 25 Sep 2026 02:33:48 +0000 Subject: [PATCH 8/8] fix(videos): the route's own 501 exports no upstream span The unsupported-provider 501 on POST /v1/videos now emits a usage event, but emit_submit_usage_event hard-coded dispatched=true, so the 501 exported an upstream CLIENT span for a call that never happened. Restore CreateSuccess.upstream_called and pass it through, as completions, embeddings and images already do. --- crates/aisix-proxy/src/videos.rs | 11 +++- .../usage-event-every-request-e2e.test.ts | 51 +++++++++++++++++++ 2 files changed, 61 insertions(+), 1 deletion(-) diff --git a/crates/aisix-proxy/src/videos.rs b/crates/aisix-proxy/src/videos.rs index a664cfafd..e0e6045bc 100644 --- a/crates/aisix-proxy/src/videos.rs +++ b/crates/aisix-proxy/src/videos.rs @@ -1670,6 +1670,7 @@ pub async fn create_video( started.elapsed(), &audit, routing.attempts.last(), + success.upstream_called, ); } success.response @@ -1738,6 +1739,9 @@ struct CreateSuccess { provider_key_id: String, applied_guardrails: Vec, monitor_hits: Vec, + /// `false` on the 501 unsupported-provider branch: no upstream call + /// happened, so its usage event exports no upstream span. + upstream_called: bool, } async fn dispatch_create( @@ -1781,6 +1785,7 @@ async fn dispatch_create( provider_key_id: String::new(), applied_guardrails: Vec::new(), monitor_hits: Vec::new(), + upstream_called: false, }); } } else { @@ -1927,6 +1932,7 @@ async fn dispatch_create( provider_key_id: String::new(), applied_guardrails, monitor_hits, + upstream_called: false, }) } }; @@ -1965,6 +1971,7 @@ async fn dispatch_create( provider_key_id: target.pk_id.clone(), applied_guardrails, monitor_hits, + upstream_called: true, }) } @@ -2244,6 +2251,8 @@ fn emit_submit_usage_event( audit: &crate::usage_attr::GuardrailAudit, // The attempt that answered (#655). winner: Option<&crate::attempt::AttemptRecord>, + // Whether the submit reached an upstream; the route's own 501 did not. + dispatched: bool, ) { let mut event = UsageEvent { request_id: client.request_id.clone(), @@ -2290,7 +2299,7 @@ fn emit_submit_usage_event( None, client.trace.as_ref(), /* terminal */ true, - /* dispatched */ true, + dispatched, ); } diff --git a/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts b/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts index e8f2741a6..7a972a4fd 100644 --- a/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts +++ b/tests/e2e/src/cases/usage-event-every-request-e2e.test.ts @@ -14,6 +14,7 @@ import { type OpenAiUpstream, type SpawnedApp, } from "../harness/index.js"; +import { startMockOtlp, type MockOtlp } from "../harness/otlp-mock.js"; // E2E: every model-serving request leaves a usage event — the record the // console Logs and budgets read — even when there is nothing to bill: @@ -48,6 +49,7 @@ describe("a usage event for every request", () => { let sls: MockSls | undefined; let cohere: OpenAiUpstream | undefined; let mcp: McpUpstream | undefined; + let otlp: MockOtlp | undefined; let etcdReachable = false; beforeAll(async () => { @@ -56,6 +58,7 @@ describe("a usage event for every request", () => { if (!etcdReachable) return; sls = await startMockSls(); + otlp = await startMockOtlp(); cohere = await startOpenAiUpstream({ nonStreamBody: COHERE_RERANK }); mcp = await startMcpUpstream("big", { reportContent: { summary: "done", log: "x".repeat(4 * BODY_LIMIT), structuredLog: "ok" }, @@ -79,6 +82,28 @@ describe("a usage event for every request", () => { content_mode: "metadata_only", }); + await seed.createObservabilityExporter({ + name: "otlp-usage-event-every-request", + enabled: true, + kind: "otlp_http", + endpoint: otlp.url, + }); + + // A video provider the gateway does not implement: it answers 501 + // itself, with no upstream call. + const videoPk = await seed.createProviderKey({ + display_name: "every-request-video", + secret: "sk-mock", + provider: "minimax", + api_base: "http://127.0.0.1:9", + }); + await seed.createModel({ + display_name: "every-request-video", + provider: "minimax", + model_name: "video-01", + provider_key_id: videoPk.id, + }); + const coherePk = await seed.createProviderKey({ display_name: "every-request-cohere", secret: "cohere-mock-key", @@ -144,6 +169,7 @@ describe("a usage event for every request", () => { await cohere?.close(); await mcp?.close(); await sls?.close(); + await otlp?.close(); }); const post = async (path: string, body: Record) => { @@ -214,4 +240,29 @@ describe("a usage event for every request", () => { expect(event.get("status_code")).toBe("502"); expect(event.get("mcp_server_name")).toBe("big"); }); + + test("POST /v1/videos: the gateway's own 501 records an event with no upstream span", async (ctx) => { + if (!etcdReachable || !app || !sls || !otlp) return ctx.skip(); + const res = await fetch(`${app.proxyUrl}/v1/videos`, { + method: "POST", + headers: auth, + body: JSON.stringify({ model: "every-request-video", prompt: "a cat" }), + }); + expect(res.status, await res.text()).toBe(501); + const requestId = res.headers.get("x-aisix-request-id"); + expect(requestId).toBeTruthy(); + + const event = await eventFor("every-request-video", "video_generation"); + expect(event.get("status_code")).toBe("501"); + + // The request's SERVER span arrives; no upstream CLIENT span beside it, + // since nothing was dispatched. + const deadline = Date.now() + 10_000; + const spans = () => otlp!.spans.filter((s) => s.attributes["aisix.request_id"] === requestId); + while (Date.now() < deadline && !spans().some((s) => s.kind === 2)) { + await new Promise((r) => setTimeout(r, 50)); + } + expect(spans().some((s) => s.kind === 2), "the SERVER span was exported").toBe(true); + expect(spans().filter((s) => s.kind === 3)).toEqual([]); + }); });