Merge origin/main into main

Integrate upstream updates while preserving the local analytics dashboards and schema-only migration changes.

Combine user account analysis with upstream user/group usage statistics in separate tabs, retain all migration versions, and keep the deleted audit document removed.

Validation: gateway all-target cargo check, frontend type check and 57 focused tests, 48 migration tests, schema composition checks, and diff whitespace checks.
This commit is contained in:
elky
2026-10-02 11:57:18 +08:00
343 changed files with 27929 additions and 2549 deletions
+12 -3
View File
@@ -15,9 +15,9 @@ use super::{BorrowedUsageEventEnvelope, UsageEvent, UsageEventData, USAGE_EVENT_
use crate::body_capture::mark_usage_event_capture_truncated;
use crate::request_metadata::{
attach_client_request_body_metadata, attach_provider_request_body_metadata,
attach_provider_response_body_metadata, clear_client_request_body_metadata,
clear_provider_request_body_metadata, request_body_derived_facts_action,
RequestBodyDerivedFactsAction,
attach_provider_response_body_metadata, attach_provider_response_model_metadata,
clear_client_request_body_metadata, clear_provider_request_body_metadata,
request_body_derived_facts_action, RequestBodyDerivedFactsAction,
};
const DIAGNOSTIC_FIELDS: [&str; 8] = [
@@ -238,6 +238,15 @@ impl WireOverrides {
RequestBodyDerivedFactsAction::Clear | RequestBodyDerivedFactsAction::Preserve => {}
}
metadata = attach_provider_response_body_metadata(metadata, data.response_body.as_ref());
metadata = attach_provider_response_model_metadata(
metadata,
request_body,
data.request_body_state,
data.api_format.as_deref(),
data.response_body.as_ref(),
data.response_body_state,
data.endpoint_api_format.as_deref(),
);
// Billing reads raw-body TTL before metadata regardless of capture state.
// Preserve that precedence independently of reasoning and tier authority.
if let Some(cache_ttl) = body_cache_ttl {
+11 -1
View File
@@ -104,6 +104,13 @@ impl UsageQueue {
pub async fn enqueue(&self, event: &UsageEvent) -> Result<String, DataLayerError> {
let encoded = self.encode_event(event)?;
self.enqueue_encoded(encoded).await
}
pub(crate) async fn enqueue_encoded(
&self,
encoded: EncodedUsageEvent,
) -> Result<String, DataLayerError> {
self.runner
.append_fields_with_maxlen(
&self.stream,
@@ -117,7 +124,10 @@ impl UsageQueue {
self.encode_event(event).map(|_| ())
}
fn encode_event(&self, event: &UsageEvent) -> Result<EncodedUsageEvent, DataLayerError> {
pub(crate) fn encode_event(
&self,
event: &UsageEvent,
) -> Result<EncodedUsageEvent, DataLayerError> {
let encoded = match event.to_bounded_stream_fields(self.config.queue_payload_max_bytes) {
Ok(encoded) => encoded,
Err(error) => {
+13 -3
View File
@@ -5,9 +5,9 @@ use aether_data_contracts::DataLayerError;
use crate::request_metadata::{
attach_client_request_body_metadata, attach_provider_request_body_metadata,
clear_client_request_body_metadata, clear_provider_request_body_metadata,
request_body_derived_facts_action, sanitize_usage_request_metadata,
RequestBodyDerivedFactsAction,
attach_provider_response_model_metadata, clear_client_request_body_metadata,
clear_provider_request_body_metadata, request_body_derived_facts_action,
sanitize_usage_request_metadata, RequestBodyDerivedFactsAction,
};
use crate::{UsageEvent, UsageEventType};
@@ -85,6 +85,16 @@ pub fn build_upsert_usage_record_from_event(
}
RequestBodyDerivedFactsAction::Preserve => {}
}
// 响应模型必须在 body 被裁剪前从客户端请求体和上游响应体共同派生;缺少权威 body 时保留已派生事实。
data.request_metadata = attach_provider_response_model_metadata(
data.request_metadata,
data.request_body.as_ref(),
data.request_body_state,
data.api_format.as_deref(),
data.response_body.as_ref(),
data.response_body_state,
data.endpoint_api_format.as_deref(),
);
let now_unix_secs = event.timestamp_ms / 1_000;
Ok(UpsertUsageRecord {
@@ -1,13 +1,15 @@
use aether_contracts::ExecutionPlan;
use aether_data_contracts::repository::usage::{
extract_provider_actual_service_tier_from_response,
extract_provider_reasoning_effort_from_body, extract_provider_service_tier_from_body,
normalize_provider_service_tier, resolve_provider_cache_ttl_minutes,
extract_provider_reasoning_effort_from_body, extract_provider_response_model_from_bodies,
extract_provider_service_tier_from_body, normalize_provider_service_tier,
resolve_provider_cache_ttl_minutes,
sanitize_usage_request_metadata as project_usage_request_metadata,
sanitize_usage_request_metadata_object as project_usage_request_metadata_object,
sanitize_usage_request_metadata_ref as project_usage_request_metadata_ref,
UsageBodyCaptureState, PROVIDER_ACTUAL_SERVICE_TIER_METADATA_KEY,
PROVIDER_CACHE_TTL_MINUTES_METADATA_KEY, PROVIDER_REASONING_EFFORT_METADATA_KEY,
usage_body_capture_is_authoritative, UsageBodyCaptureState,
PROVIDER_ACTUAL_SERVICE_TIER_METADATA_KEY, PROVIDER_CACHE_TTL_MINUTES_METADATA_KEY,
PROVIDER_REASONING_EFFORT_METADATA_KEY, PROVIDER_RESPONSE_MODEL_METADATA_KEY,
PROVIDER_SERVICE_TIER_METADATA_KEY, REQUESTED_REASONING_EFFORT_METADATA_KEY,
};
use serde_json::{Map, Value};
@@ -249,6 +251,83 @@ pub(crate) fn attach_provider_response_body_metadata(
attach_provider_actual_service_tier_metadata(metadata, actual_service_tier.as_deref())
}
pub(crate) fn attach_provider_response_model_metadata(
metadata: Option<Value>,
request_body: Option<&Value>,
request_body_state: Option<UsageBodyCaptureState>,
request_api_format: Option<&str>,
response_body: Option<&Value>,
response_body_state: Option<UsageBodyCaptureState>,
provider_api_format: Option<&str>,
) -> Option<Value> {
let both_bodies_are_authoritative =
usage_body_capture_is_authoritative(request_body, request_body_state)
&& usage_body_capture_is_authoritative(response_body, response_body_state);
let response_model = extract_provider_response_model_from_bodies(
request_body,
request_body_state,
request_api_format,
response_body,
response_body_state,
provider_api_format,
);
if !both_bodies_are_authoritative && response_model.is_none() {
return metadata;
}
let mut object = match metadata {
Some(Value::Object(object)) => object,
_ => Map::new(),
};
// 完整终态 body 是最终候选的权威事实;相同、无效或缺失模型都要清除旧候选值。
if both_bodies_are_authoritative {
object.remove(PROVIDER_RESPONSE_MODEL_METADATA_KEY);
}
if let Some(response_model) = response_model {
object.insert(
PROVIDER_RESPONSE_MODEL_METADATA_KEY.to_string(),
Value::String(response_model),
);
}
(!object.is_empty()).then_some(Value::Object(object))
}
/// 终态候选无法完成比较时,显式清除旧响应模型,避免重试/故障转移残留。
pub(crate) fn refresh_provider_response_model_metadata(
metadata: Option<Value>,
request_body: Option<&Value>,
request_body_state: Option<UsageBodyCaptureState>,
request_api_format: Option<&str>,
response_body: Option<&Value>,
response_body_state: Option<UsageBodyCaptureState>,
provider_api_format: Option<&str>,
) -> Option<Value> {
let mut object = match metadata {
Some(Value::Object(object)) => object,
_ => Map::new(),
};
object.remove(PROVIDER_RESPONSE_MODEL_METADATA_KEY);
if usage_body_capture_is_authoritative(request_body, request_body_state)
&& usage_body_capture_is_authoritative(response_body, response_body_state)
{
if let Some(response_model) = extract_provider_response_model_from_bodies(
request_body,
request_body_state,
request_api_format,
response_body,
response_body_state,
provider_api_format,
) {
object.insert(
PROVIDER_RESPONSE_MODEL_METADATA_KEY.to_string(),
Value::String(response_model),
);
}
}
(!object.is_empty()).then_some(Value::Object(object))
}
/// Refreshes the response-derived tier for a terminal snapshot. Complete response objects are
/// authoritative even when they contain no tier (which clears a stale candidate value). Capture
/// placeholders/absent bodies are not authoritative, so a terminal summary already present in
@@ -312,6 +391,7 @@ pub(crate) fn attach_provider_actual_service_tier_metadata(
#[cfg(test)]
mod tests {
use aether_contracts::{ExecutionPlan, RequestBody};
use aether_data_contracts::repository::usage::UsageBodyCaptureState;
use serde_json::{json, Value};
use std::collections::BTreeMap;
@@ -323,8 +403,9 @@ mod tests {
use super::{
attach_client_request_body_metadata, attach_provider_actual_service_tier_metadata,
attach_provider_request_body_metadata, attach_provider_response_body_metadata,
build_usage_request_metadata_seed, merge_usage_request_metadata,
merge_usage_request_metadata_owned, refresh_provider_response_body_metadata,
attach_provider_response_model_metadata, build_usage_request_metadata_seed,
merge_usage_request_metadata, merge_usage_request_metadata_owned,
refresh_provider_response_body_metadata, refresh_provider_response_model_metadata,
retain_first_byte_request_metadata, sanitize_usage_request_metadata,
sanitize_usage_request_metadata_ref,
};
@@ -802,6 +883,44 @@ mod tests {
);
}
#[test]
fn response_model_metadata_is_independent_from_mapping_and_clears_stale_values() {
let metadata = attach_provider_response_model_metadata(
Some(json!({"provider_response_model": "old-model", "trace_id": "trace-1"})),
Some(&json!({"model": "gpt-5"})),
Some(UsageBodyCaptureState::Inline),
Some("openai:chat"),
Some(&json!({"model": "gpt-5.1"})),
Some(UsageBodyCaptureState::Inline),
Some("openai:chat"),
)
.expect("response model should be attached");
assert_eq!(metadata["provider_response_model"], "gpt-5.1");
assert_eq!(metadata["trace_id"], "trace-1");
let metadata = refresh_provider_response_model_metadata(
Some(json!({"provider_response_model": "gpt-5.1"})),
Some(&json!({"model": "gpt-5"})),
Some(UsageBodyCaptureState::Inline),
Some("openai:chat"),
Some(&json!({"model": "gpt-5"})),
Some(UsageBodyCaptureState::Inline),
Some("openai:chat"),
);
assert!(metadata.is_none());
let metadata = refresh_provider_response_model_metadata(
Some(json!({"provider_response_model": "gpt-5.1"})),
None,
Some(UsageBodyCaptureState::Disabled),
Some("openai:chat"),
Some(&json!({"model": "gpt-5.2"})),
Some(UsageBodyCaptureState::Inline),
Some("openai:chat"),
);
assert!(metadata.is_none());
}
#[test]
fn terminal_response_refresh_replaces_stale_actual_tier() {
let metadata = refresh_provider_response_body_metadata(
+317 -67
View File
@@ -21,9 +21,10 @@ use crate::executor::spawn_on_usage_background_runtime;
use crate::queue::is_permanent_enqueue_error;
use crate::request_metadata::{
attach_client_request_body_metadata, attach_provider_request_body_metadata,
attach_provider_response_body_metadata, clear_client_request_body_metadata,
clear_provider_request_body_metadata, request_body_derived_facts_action,
retain_first_byte_request_metadata, RequestBodyDerivedFactsAction,
attach_provider_response_body_metadata, attach_provider_response_model_metadata,
clear_client_request_body_metadata, clear_provider_request_body_metadata,
request_body_derived_facts_action, retain_first_byte_request_metadata,
RequestBodyDerivedFactsAction,
};
use crate::settlement::{
reconcile_usage_policy_cost_for_event_with_result, settle_usage_with_reconciled_cost,
@@ -4919,7 +4920,7 @@ impl UsageRuntime {
&self,
data: &T,
queue: UsageQueue,
event: UsageEvent,
mut event: UsageEvent,
) -> TerminalPersistenceOutcome
where
T: UsageRuntimeAccess,
@@ -4952,7 +4953,20 @@ impl UsageRuntime {
.await;
};
if let Err(err) = queue.enqueue(&event).await {
let enqueue_result = match queue.encode_event(&event) {
Ok(encoded) => {
if encoded.diagnostics_omitted
&& self
.try_write_terminal_direct_fallback(data, &mut event, "queue_wire_limit")
.await
{
return TerminalPersistenceOutcome::PersistedDirectly;
}
queue.enqueue_encoded(encoded).await
}
Err(err) => Err(err),
};
if let Err(err) = enqueue_result {
drop(_guard);
if is_permanent_enqueue_error(&err) {
return self
@@ -5284,8 +5298,17 @@ fn preserve_request_facts_with_legacy_missing(
fn preserve_provider_response_facts(event: &mut UsageEvent) {
let metadata = event.data.request_metadata.take();
event.data.request_metadata =
let metadata =
attach_provider_response_body_metadata(metadata, event.data.response_body.as_ref());
event.data.request_metadata = attach_provider_response_model_metadata(
metadata,
event.data.request_body.as_ref(),
event.data.request_body_state,
event.data.api_format.as_deref(),
event.data.response_body.as_ref(),
event.data.response_body_state,
event.data.endpoint_api_format.as_deref(),
);
}
impl UsageQueueHealthSnapshot {
@@ -6649,7 +6672,7 @@ mod tests {
}
use std::collections::BTreeMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Instant;
@@ -7375,9 +7398,27 @@ mod tests {
queue: Arc<dyn RuntimeQueueStore>,
policy_started: Arc<tokio::sync::Notify>,
release_policy: Arc<tokio::sync::Notify>,
policy_released: Arc<AtomicBool>,
policy_reads: Arc<AtomicUsize>,
}
impl BlockingPolicyQueueConfiguredUsageStore {
fn new(queue: Arc<dyn RuntimeQueueStore>) -> Self {
Self {
queue,
policy_started: Arc::new(tokio::sync::Notify::new()),
release_policy: Arc::new(tokio::sync::Notify::new()),
policy_released: Arc::new(AtomicBool::new(false)),
policy_reads: Arc::new(AtomicUsize::new(0)),
}
}
fn release_blocked_policy(&self) {
self.policy_released.store(true, Ordering::Release);
self.release_policy.notify_waiters();
}
}
#[derive(Default)]
struct FailingPolicyUsageStore {
inner: NoRedisUsageStore,
@@ -8409,7 +8450,18 @@ mod tests {
async fn body_capture_policy(&self) -> Result<UsageBodyCapturePolicy, DataLayerError> {
self.policy_reads.fetch_add(1, Ordering::AcqRel);
self.policy_started.notify_one();
self.release_policy.notified().await;
// Latch the gate: Notify is edge-triggered, and later policy reads
// (or a waiter that subscribed after a single notify) must not hang.
loop {
if self.policy_released.load(Ordering::Acquire) {
break;
}
let notified = self.release_policy.notified();
if self.policy_released.load(Ordering::Acquire) {
break;
}
notified.await;
}
Ok(UsageBodyCapturePolicy::default())
}
}
@@ -9043,15 +9095,39 @@ mod tests {
.await
.expect("a duplicate first-byte marker must release the terminal barrier");
let records = store.records.lock().expect("records lock");
assert_eq!(
records.len(),
2,
"the duplicate first byte must be coalesced"
);
assert_eq!(records[0].status, "streaming");
assert_eq!(records[1].status, "completed");
drop(records);
{
let records = store.records.lock().expect("records lock");
assert_eq!(
records.len(),
2,
"the duplicate first byte must be coalesced"
);
assert_eq!(records[0].status, "streaming");
assert_eq!(records[1].status, "completed");
}
// The terminal persistence notification can arrive before the submission
// dispatcher accounts for its completed task and releases admission.
timeout(Duration::from_secs(1), async {
loop {
let snapshot = runtime.metrics_snapshot();
if snapshot.lifecycle_submission_pending == 0
&& snapshot.first_byte_persistence_pending == 0
&& snapshot.ordered_lifecycle_pending == 0
&& runtime
.lifecycle_submission
.state
.admission
.available_permits()
== CAPACITY
{
break;
}
sleep(Duration::from_millis(1)).await;
}
})
.await
.expect("duplicate first-byte submission accounting should drain");
let snapshot = runtime.metrics_snapshot();
assert_eq!(snapshot.lifecycle_submission_pending, 0);
@@ -9979,6 +10055,24 @@ mod tests {
2,
"the later direct caller should receive its own bounded write attempt"
);
// Direct persistence can finish before the submission worker joins the
// barrier handoff and accounts for its completed slot.
timeout(Duration::from_secs(1), async {
loop {
let snapshot = runtime.metrics_snapshot();
let submission = &runtime.lifecycle_submission.state;
if snapshot.terminal_submission_pending == 0
&& snapshot.ordered_lifecycle_pending == 0
&& snapshot.lifecycle_submission_pending == 0
&& submission.admission.available_permits() == submission.capacity
{
break;
}
sleep(Duration::from_millis(1)).await;
}
})
.await
.expect("failed terminal submission accounting and admission should drain");
let snapshot = runtime.metrics_snapshot();
assert_eq!(snapshot.terminal_submission_pending, 0);
assert_eq!(snapshot.ordered_lifecycle_pending, 0);
@@ -10086,24 +10180,43 @@ mod tests {
.await;
assert_eq!(remaining_policy_panics.load(Ordering::Acquire), 0);
let records = records.lock().expect("records lock");
assert_eq!(
records
.iter()
.filter(|record| record.request_id == healthy_request_id)
.count(),
2,
"the same terminal shard should continue processing healthy requests"
);
assert!(
records
.iter()
.filter(|record| record.request_id == failed_request_id)
.count()
== 1,
"only the later healthy attempt should persist for the panicked request"
);
drop(records);
{
let records = records.lock().expect("records lock");
assert_eq!(
records
.iter()
.filter(|record| record.request_id == healthy_request_id)
.count(),
2,
"the same terminal shard should continue processing healthy requests"
);
assert!(
records
.iter()
.filter(|record| record.request_id == failed_request_id)
.count()
== 1,
"only the later healthy attempt should persist for the panicked request"
);
}
// The final direct attempt also submits a barrier whose worker may
// account for completion after the persistence call has returned.
timeout(Duration::from_secs(1), async {
loop {
let snapshot = runtime.metrics_snapshot();
let submission = &runtime.lifecycle_submission.state;
if snapshot.terminal_submission_pending == 0
&& snapshot.ordered_lifecycle_pending == 0
&& snapshot.lifecycle_submission_pending == 0
&& submission.admission.available_permits() == submission.capacity
{
break;
}
sleep(Duration::from_millis(1)).await;
}
})
.await
.expect("panicked terminal submission accounting and admission should drain");
let snapshot = runtime.metrics_snapshot();
assert_eq!(snapshot.terminal_submission_pending, 0);
assert_eq!(snapshot.ordered_lifecycle_pending, 0);
@@ -12325,12 +12438,9 @@ mod tests {
async fn event_capture_budget_bounds_blocked_policy_waiters_and_releases_on_cancel_or_basic() {
for limit in [0, 64 * 1024] {
let runtime = UsageRuntime::new(UsageRuntimeConfig::default()).expect("runtime");
let store = BlockingPolicyQueueConfiguredUsageStore {
queue: Arc::new(RuntimeState::memory(MemoryRuntimeStateConfig::default())),
policy_started: Arc::new(tokio::sync::Notify::new()),
release_policy: Arc::new(tokio::sync::Notify::new()),
policy_reads: Arc::new(AtomicUsize::new(0)),
};
let store = BlockingPolicyQueueConfiguredUsageStore::new(Arc::new(
RuntimeState::memory(MemoryRuntimeStateConfig::default()),
));
let budget = Arc::new(crate::event_capture_budget::EventCaptureMemoryBudget::new(
limit,
));
@@ -12391,7 +12501,7 @@ mod tests {
.await
.expect("replacement policy read starts");
assert_eq!(budget.retained_bytes(), retained);
store.release_policy.notify_one();
store.release_blocked_policy();
let event = timeout(Duration::from_secs(2), completing)
.await
.expect("Basic policy completes")
@@ -12675,6 +12785,159 @@ mod tests {
assert_eq!(queued.data.total_cost_usd, None);
}
fn oversized_full_terminal_event(max_bytes: usize) -> (serde_json::Value, UsageEvent) {
let body = json!({"content": "full-body".repeat(max_bytes / 16)});
let mut event = UsageEvent::new(
UsageEventType::Completed,
"oversized-full-capture",
UsageEventData {
provider_name: "openai".to_string(),
model: "gpt-5".to_string(),
status_code: Some(200),
total_tokens: Some(12),
request_body: Some(body.clone()),
provider_request_body: Some(body.clone()),
response_body: Some(body.clone()),
client_response_body: Some(body.clone()),
..UsageEventData::default()
},
);
apply_usage_body_capture_policy_to_event(
UsageBodyCapturePolicy {
record_level: UsageRequestRecordLevel::Full,
},
&mut event,
);
(body, event)
}
#[tokio::test]
async fn oversized_full_terminal_capture_is_persisted_without_queue_truncation() {
let config = UsageRuntimeConfig {
enabled: true,
queue_terminal_events: true,
consumer_block_ms: 1,
..UsageRuntimeConfig::default()
};
let queue_runner: Arc<dyn RuntimeQueueStore> =
Arc::new(RuntimeState::memory(MemoryRuntimeStateConfig::default()));
let queue = UsageQueue::new(Arc::clone(&queue_runner), config.clone()).unwrap();
let store = EnrichmentCountingQueueStore {
records: Mutex::new(Vec::new()),
queue: queue_runner,
enrich_calls: AtomicUsize::new(0),
};
let runtime = UsageRuntime::new(config.clone()).unwrap();
let (body, event) = oversized_full_terminal_event(config.queue_payload_max_bytes);
assert!(
event
.to_bounded_stream_fields(config.queue_payload_max_bytes)
.unwrap()
.diagnostics_omitted
);
let outcome = runtime.enqueue_or_write_terminal(&store, event).await;
assert_eq!(
outcome,
super::TerminalPersistenceOutcome::PersistedDirectly
);
assert_eq!(store.enrich_calls.load(Ordering::Acquire), 1);
assert_eq!(queue.stats().await.unwrap().stream_length, 0);
let records = store.records.lock().unwrap();
assert_eq!(records.len(), 1);
for captured in [
&records[0].request_body,
&records[0].provider_request_body,
&records[0].response_body,
&records[0].client_response_body,
] {
assert_eq!(captured.as_ref(), Some(&body));
}
}
#[tokio::test]
async fn oversized_full_terminal_capture_keeps_bounded_queue_fallback() {
for unavailable in ["writer", "write_failure", "worker_gate", "fallback_gate"] {
let config = UsageRuntimeConfig {
enabled: true,
queue_terminal_events: true,
consumer_block_ms: 1,
worker_record_concurrency_limit: Some(1),
..UsageRuntimeConfig::default()
};
let queue_runner: Arc<dyn RuntimeQueueStore> =
Arc::new(RuntimeState::memory(MemoryRuntimeStateConfig::default()));
let queue = UsageQueue::new(Arc::clone(&queue_runner), config.clone()).unwrap();
queue.ensure_consumer_group().await.unwrap();
let store = FailingWriteQueueConfiguredUsageStore {
queue: Arc::clone(&queue_runner),
upsert_attempts: Arc::new(AtomicUsize::new(0)),
};
let queue_only = QueueOnlyUsageStore {
queue: queue_runner,
upsert_attempts: Arc::clone(&store.upsert_attempts),
};
let runtime = UsageRuntime::new(config.clone()).unwrap();
let worker_permit = (unavailable == "worker_gate").then(|| {
runtime
.worker_record_gate
.as_ref()
.unwrap()
.try_acquire()
.unwrap()
});
let fallback_permit = (unavailable == "fallback_gate").then(|| {
runtime
.terminal_direct_fallback_state
.try_acquire()
.unwrap()
});
let (_, event) = oversized_full_terminal_event(config.queue_payload_max_bytes);
let outcome = if unavailable == "writer" {
runtime.enqueue_or_write_terminal(&queue_only, event).await
} else {
runtime.enqueue_or_write_terminal(&store, event).await
};
assert_eq!(
outcome,
super::TerminalPersistenceOutcome::Queued,
"{unavailable}"
);
assert_eq!(
store.upsert_attempts.load(Ordering::Acquire),
usize::from(unavailable == "write_failure")
);
let entries = queue
.read_group("oversized-capture-consumer")
.await
.unwrap();
assert_eq!(entries.len(), 1);
assert!(entries[0].fields["payload"].len() <= config.queue_payload_max_bytes);
let queued = UsageEvent::from_stream_fields(&entries[0].fields).unwrap();
assert_eq!(queued.data.total_tokens, Some(12));
for (body, state) in [
(&queued.data.request_body, queued.data.request_body_state),
(
&queued.data.provider_request_body,
queued.data.provider_request_body_state,
),
(&queued.data.response_body, queued.data.response_body_state),
(
&queued.data.client_response_body,
queued.data.client_response_body_state,
),
] {
assert!(body.is_none());
assert_eq!(state, Some(UsageBodyCaptureState::Truncated));
}
assert_eq!(runtime.metrics_snapshot().terminal_enqueue_failed_total, 0);
drop((worker_permit, fallback_permit));
}
}
#[tokio::test]
async fn terminal_enqueue_failure_uses_bounded_direct_database_fallback() {
let config = UsageRuntimeConfig {
@@ -13398,12 +13661,7 @@ mod tests {
Arc::new(RuntimeState::memory(MemoryRuntimeStateConfig::default()));
let tracked_queue = Arc::new(FlakyAppendQueueStore::new(inner_queue, 0));
let queue: Arc<dyn RuntimeQueueStore> = tracked_queue.clone();
let store = BlockingPolicyQueueConfiguredUsageStore {
queue,
policy_started: Arc::new(tokio::sync::Notify::new()),
release_policy: Arc::new(tokio::sync::Notify::new()),
policy_reads: Arc::new(AtomicUsize::new(0)),
};
let store = BlockingPolicyQueueConfiguredUsageStore::new(queue);
let runtime = UsageRuntime::new(config).expect("usage runtime should build");
let request_id = "req-terminal-seed-waits-for-turn";
let plan = terminal_test_plan(request_id);
@@ -13425,9 +13683,10 @@ mod tests {
assert_eq!(blocked_snapshot.terminal_submission_in_flight, 0);
assert!(blocked_snapshot.lifecycle_submission_pending >= 2);
store.release_policy.notify_waiters();
store.release_blocked_policy();
timeout(Duration::from_secs(2), async {
loop {
store.release_blocked_policy();
let snapshot = runtime.metrics_snapshot();
if tracked_queue.successful_appends.load(Ordering::Acquire) == 1
&& snapshot.lifecycle_submission_pending == 0
@@ -13435,7 +13694,7 @@ mod tests {
{
break;
}
tokio::task::yield_now().await;
sleep(Duration::from_millis(1)).await;
}
})
.await
@@ -13468,12 +13727,7 @@ mod tests {
Arc::new(RuntimeState::memory(MemoryRuntimeStateConfig::default()));
let tracked_queue = Arc::new(FlakyAppendQueueStore::new(inner_queue, 0));
let queue: Arc<dyn RuntimeQueueStore> = tracked_queue.clone();
let store = BlockingPolicyQueueConfiguredUsageStore {
queue,
policy_started: Arc::new(tokio::sync::Notify::new()),
release_policy: Arc::new(tokio::sync::Notify::new()),
policy_reads: Arc::new(AtomicUsize::new(0)),
};
let store = BlockingPolicyQueueConfiguredUsageStore::new(queue);
let runtime = UsageRuntime::new(config).expect("usage runtime should build");
let policy_started = store.policy_started.notified();
@@ -13520,9 +13774,10 @@ mod tests {
assert_eq!(blocked_snapshot.terminal_submission_in_flight, 1);
assert!(blocked_snapshot.lifecycle_submission_pending <= BACKLOG + 1);
store.release_policy.notify_waiters();
store.release_blocked_policy();
timeout(Duration::from_secs(5), async {
loop {
store.release_blocked_policy();
let snapshot = runtime.metrics_snapshot();
if tracked_queue.successful_appends.load(Ordering::Acquire) == BACKLOG + 1
&& snapshot.lifecycle_submission_pending == 0
@@ -13530,7 +13785,7 @@ mod tests {
{
break;
}
tokio::task::yield_now().await;
sleep(Duration::from_millis(1)).await;
}
})
.await
@@ -13701,7 +13956,7 @@ mod tests {
.iter()
.any(|record| record.request_id == blocked_request_id));
release_build.notify_waiters();
release_build.notify_one();
timeout(Duration::from_secs(2), async {
loop {
let snapshot = runtime.metrics_snapshot();
@@ -13828,12 +14083,7 @@ mod tests {
Arc::new(RuntimeState::memory(MemoryRuntimeStateConfig::default()));
let tracked_queue = Arc::new(FlakyAppendQueueStore::new(inner_queue, 0));
let queue: Arc<dyn RuntimeQueueStore> = tracked_queue.clone();
let store = BlockingPolicyQueueConfiguredUsageStore {
queue,
policy_started: Arc::new(tokio::sync::Notify::new()),
release_policy: Arc::new(tokio::sync::Notify::new()),
policy_reads: Arc::new(AtomicUsize::new(0)),
};
let store = BlockingPolicyQueueConfiguredUsageStore::new(queue);
let runtime = UsageRuntime::new(config).expect("usage runtime should build");
let policy_started = store.policy_started.notified();
runtime
@@ -13892,10 +14142,10 @@ mod tests {
.expect("terminal submissions should reach the execution backlog");
let saturated_snapshot = runtime.metrics_snapshot();
store.release_policy.notify_waiters();
store.release_blocked_policy();
let all_completed = timeout(Duration::from_secs(2), async {
loop {
store.release_policy.notify_waiters();
store.release_blocked_policy();
if tracked_queue.successful_appends.load(Ordering::Acquire)
== EXCESS_SUBMISSIONS + 1
&& runtime.metrics_snapshot().terminal_submission_in_flight == 0
+64 -5
View File
@@ -17,9 +17,10 @@ use crate::body_capture::{
};
use crate::request_metadata::{
attach_client_request_body_metadata, attach_provider_actual_service_tier_metadata,
attach_provider_request_body_metadata, build_usage_request_metadata_seed,
merge_usage_request_metadata, merge_usage_request_metadata_owned,
refresh_provider_response_body_metadata, sanitize_usage_request_metadata,
attach_provider_request_body_metadata, attach_provider_response_model_metadata,
build_usage_request_metadata_seed, merge_usage_request_metadata,
merge_usage_request_metadata_owned, refresh_provider_response_body_metadata,
refresh_provider_response_model_metadata, sanitize_usage_request_metadata,
sanitize_usage_request_metadata_ref,
};
use crate::{
@@ -727,6 +728,15 @@ fn build_terminal_usage_event_from_seed_impl(
Some(model.as_str()),
provider_request.as_ref(),
);
let request_metadata = attach_provider_response_model_metadata(
request_metadata,
request_body.as_ref(),
body_states.request_body_state,
Some(client_contract.as_str()),
provider_response.as_ref(),
body_states.response_body_state,
Some(provider_contract.as_str()),
);
let mut data = UsageEventData {
user_id,
@@ -1073,6 +1083,15 @@ pub fn build_sync_terminal_usage_seed(
context_seed.request_metadata,
provider_response_full.as_ref(),
);
let request_metadata = refresh_provider_response_model_metadata(
request_metadata,
context_seed.request_body.as_ref(),
context_seed.body_states.request_body_state,
Some(context_seed.client_contract.as_str()),
provider_response_full.as_ref(),
provider_response_body_state,
Some(context_seed.provider_contract.as_str()),
);
TerminalUsageSeed {
token_measurement_source,
@@ -1252,6 +1271,15 @@ pub fn build_stream_terminal_usage_seed(
context_seed.request_metadata,
provider_response_full.as_ref(),
);
let request_metadata = refresh_provider_response_model_metadata(
request_metadata,
context_seed.request_body.as_ref(),
context_seed.body_states.request_body_state,
Some(context_seed.client_contract.as_str()),
provider_response_full.as_ref(),
provider_response_body_state,
Some(context_seed.provider_contract.as_str()),
);
// The parser's terminal summary is authoritative when a response body is truncated or the
// body and summary disagree; attach it after the body refresh so it wins.
let request_metadata = attach_provider_actual_service_tier_metadata(
@@ -3080,14 +3108,26 @@ fn parse_sse_body_for_storage(text: &str) -> Option<Value> {
let mut chunks = Vec::new();
let mut total_chunks = 0_u64;
let mut saw_done = false;
let mut first_parse_error = None;
for_each_sse_payload(text, |payload| {
if payload == "[DONE]" {
saw_done = true;
return;
}
total_chunks += 1;
if let Ok(json_body) = serde_json::from_str::<Value>(payload) {
chunks.push(json_body);
match serde_json::from_str::<Value>(payload) {
Ok(json_body) => chunks.push(json_body),
Err(error) if first_parse_error.is_none() => {
// A later valid event must not hide an earlier broken one.
// Store diagnostics only, without duplicating raw user content.
first_parse_error = Some(json!({
"chunk_index": total_chunks - 1,
"line": error.line(),
"column": error.column(),
"message": error.to_string(),
}));
}
Err(_) => {}
}
});
if total_chunks == 0 && !saw_done {
@@ -3104,6 +3144,11 @@ fn parse_sse_body_for_storage(text: &str) -> Option<Value> {
if saw_done {
metadata.insert("has_completion".to_string(), Value::Bool(true));
}
if let Some(error) = first_parse_error {
// Capture truncation can also cause a parse error; this describes the
// captured payload, not an assertion that the provider sent bad JSON.
metadata.insert("first_parse_error".to_string(), error);
}
if stored_chunks < total_chunks {
metadata.insert(
"dropped_chunks".to_string(),
@@ -7221,6 +7266,20 @@ mod tests {
);
}
#[test]
fn parse_sse_body_for_storage_reports_bad_event_before_valid_terminal() {
let body = concat!(
"data: {\"tools\":[}\n\n",
"data: {\"type\":\"response.completed\"}\n\n",
);
let parsed = parse_sse_body_for_storage(body).unwrap();
assert_eq!(parsed["metadata"]["dropped_chunks"], 1);
assert_eq!(parsed["metadata"]["first_parse_error"]["chunk_index"], 0);
assert!(parsed["metadata"]["first_parse_error"]["message"].is_string());
assert_eq!(parsed["chunks"][0]["type"], "response.completed");
assert!(parsed.get("raw_response").is_none());
}
#[test]
fn extract_token_counts_from_value_handles_crlf_and_cr_sse_text() {
let sse_body = concat!(