Skip to content
Merged
20 changes: 12 additions & 8 deletions crates/aisix-core/src/models/guardrail.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
24 changes: 11 additions & 13 deletions crates/aisix-proxy/src/audio.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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<S, F>(
upstream: S,
// Set when `upstream` ended on a read timeout rather than its own end.
Expand All @@ -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).
Expand Down
31 changes: 16 additions & 15 deletions crates/aisix-proxy/src/chat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
Expand All @@ -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,
Expand Down Expand Up @@ -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"))
Expand All @@ -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),
);
}
Expand Down Expand Up @@ -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)
Expand Down
68 changes: 25 additions & 43 deletions crates/aisix-proxy/src/completions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand All @@ -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<CompletionUsage>,
/// Whether the request reached the provider. False only for the 501
/// provider-unsupported branch.
Expand Down Expand Up @@ -209,18 +208,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,
Expand All @@ -231,13 +223,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,
Expand Down Expand Up @@ -766,8 +758,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(),
Expand All @@ -785,9 +776,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
Expand Down Expand Up @@ -1553,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;
Expand Down Expand Up @@ -1619,15 +1608,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;

Expand Down Expand Up @@ -1670,14 +1658,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",
Expand All @@ -1698,14 +1685,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;
Expand Down
Loading
Loading