diff --git a/crates/aether-data/adapters/postgres/src/usage/queries/find_by_id_sql.sql b/crates/aether-data/adapters/postgres/src/usage/queries/find_by_id_sql.sql index a5252794e..298d4909e 100644 --- a/crates/aether-data/adapters/postgres/src/usage/queries/find_by_id_sql.sql +++ b/crates/aether-data/adapters/postgres/src/usage/queries/find_by_id_sql.sql @@ -185,6 +185,10 @@ SELECT usage_http_audits.provider_request_body_ref AS http_provider_request_body_ref, usage_http_audits.response_body_ref AS http_response_body_ref, usage_http_audits.client_response_body_ref AS http_client_response_body_ref, + usage_http_audits.request_body_state AS http_request_body_state, + usage_http_audits.provider_request_body_state AS http_provider_request_body_state, + usage_http_audits.response_body_state AS http_response_body_state, + usage_http_audits.client_response_body_state AS http_client_response_body_state, usage_routing_snapshots.candidate_id AS routing_candidate_id, usage_routing_snapshots.candidate_index AS routing_candidate_index, usage_routing_snapshots.key_name AS routing_key_name, diff --git a/crates/aether-data/adapters/postgres/src/usage/tests.rs b/crates/aether-data/adapters/postgres/src/usage/tests.rs index 082765ae7..a1cef6095 100644 --- a/crates/aether-data/adapters/postgres/src/usage/tests.rs +++ b/crates/aether-data/adapters/postgres/src/usage/tests.rs @@ -3467,6 +3467,18 @@ fn usage_sql_reads_http_audits_for_single_record_fetches() { assert!(super::FIND_BY_ID_SQL.contains("LEFT JOIN usage_http_audits")); assert!(super::FIND_BY_REQUEST_ID_SQL.contains("http_request_body_ref")); assert!(super::FIND_BY_ID_SQL.contains("http_client_response_body_ref")); + for sql in [super::FIND_BY_REQUEST_ID_SQL, super::FIND_BY_ID_SQL] { + for field in [ + "request_body", + "provider_request_body", + "response_body", + "client_response_body", + ] { + assert!(sql.contains(&format!( + "usage_http_audits.{field}_state AS http_{field}_state" + ))); + } + } } #[test] diff --git a/crates/aether-usage/runtime/src/queue.rs b/crates/aether-usage/runtime/src/queue.rs index 23b6dadb6..04d8fe6fe 100644 --- a/crates/aether-usage/runtime/src/queue.rs +++ b/crates/aether-usage/runtime/src/queue.rs @@ -104,6 +104,13 @@ impl UsageQueue { pub async fn enqueue(&self, event: &UsageEvent) -> Result { let encoded = self.encode_event(event)?; + self.enqueue_encoded(encoded).await + } + + pub(crate) async fn enqueue_encoded( + &self, + encoded: EncodedUsageEvent, + ) -> Result { 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 { + pub(crate) fn encode_event( + &self, + event: &UsageEvent, + ) -> Result { let encoded = match event.to_bounded_stream_fields(self.config.queue_payload_max_bytes) { Ok(encoded) => encoded, Err(error) => { diff --git a/crates/aether-usage/runtime/src/runtime.rs b/crates/aether-usage/runtime/src/runtime.rs index 94044cafe..99706f638 100644 --- a/crates/aether-usage/runtime/src/runtime.rs +++ b/crates/aether-usage/runtime/src/runtime.rs @@ -4919,7 +4919,7 @@ impl UsageRuntime { &self, data: &T, queue: UsageQueue, - event: UsageEvent, + mut event: UsageEvent, ) -> TerminalPersistenceOutcome where T: UsageRuntimeAccess, @@ -4952,7 +4952,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 @@ -12725,6 +12738,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 = + 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 = + 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 { diff --git a/docs/operations/concurrency-design-audit-2026-09-09.md b/docs/operations/concurrency-design-audit-2026-09-09.md index a81109632..9c68c95af 100644 --- a/docs/operations/concurrency-design-audit-2026-09-09.md +++ b/docs/operations/concurrency-design-audit-2026-09-09.md @@ -184,6 +184,7 @@ ## 第九轮修复状态 - **新增队列消息的完整字节上限:** `AETHER_GATEWAY_USAGE_QUEUE_PAYLOAD_MAX_BYTES` 默认 1 MiB,启用 usage runtime 时显式 `0` 非法。`UsageQueue::enqueue` 在发送 Redis 命令前,以有界 writer 编码完整 v1 JSON envelope,包含 UTF-8、转义、metadata、正文、headers 及其他字段;不会先生成无限制的完整 JSON 字符串再检查长度。普通消息保持原格式,原公开编码接口及历史消息解码继续兼容。 +- **FULL 终态正文保留补充:** 正常终态入队发现需要剥离诊断正文时,先用原事件尝试已有的受限数据库直写,成功后不再入队降级副本。该路径仍遵守数据库压力检查、共享写入并发门限和直写门限,不提高队列字节上限;仅队列节点、写入失败或门限饱和时仍允许有界降级,保留计费事实及 `Truncated` 状态。按使用记录 ID 读取详情同时返回四类正文的采集状态,避免将已知截断误报为 `legacy_unknown`。已经丢弃的历史正文无法由此恢复。 - **诊断降级保留计费语义:** 完整消息超限后,先借用检查去掉四份正文和四份 headers 的核心字段大小,核心可容纳才克隆 metadata 并生成诊断投影;复用同一字节缓冲。保留 token、费用、显式零、错误存在性、身份、时间、终态、正文引用、预留 token 和计费维度。按完整 v1 消费者规则保留请求档位、推理参数、响应实际档位及缓存 TTL,并标记被移除的正文为 `Truncated`。显式 `None`、`Disabled`、`Unavailable` 不改为可回退的状态;JSON null 与非对象正文分别按旧解码和权威规则处理。 - **无法安全编码时的失败路径:** 核心仍超限,或去掉正文无法保留原缓存 TTL 计费语义时,返回 `InvalidInput`。终态沿现有有并发限制的数据库路径使用原事件回退;数据库不可用、受压或写入失败时明确返回 `Failed`,保留 first-byte 状态,不误报已入队或已缓冲。这类输入错误不打开 Redis 熔断,也不会无限重试。重试接收前再次校验,覆盖主路径已熔断或入队槽耗尽的旁路;重试 worker 也会终止单条永久失败并继续处理后续条目。 - **指标与临时分配:** 导出队列 payload 上限、诊断降级、编码拒绝及永久重试失败计数。payload 计数是进程级编码尝试,包含入队与重试预校验,不是唯一事件数。超长 tier/reasoning 字符串先检查已有的 64 字节限制,再执行大小写规范化,避免明知非法仍复制整个字符串。 diff --git a/frontend/src/features/usage/utils/__tests__/body-document-engine.spec.ts b/frontend/src/features/usage/utils/__tests__/body-document-engine.spec.ts index ed5e9f413..b04d86c07 100644 --- a/frontend/src/features/usage/utils/__tests__/body-document-engine.spec.ts +++ b/frontend/src/features/usage/utils/__tests__/body-document-engine.spec.ts @@ -1,12 +1,31 @@ -import { describe, expect, it } from 'vitest' +import { describe, expect, it, vi } from 'vitest' import { gzipSync } from 'node:zlib' import { BodyDocumentEngine, decodeBody } from '../body-document-engine' import { JSON_PAGE_SIZE, JSON_TEXT_CHUNK_SIZE } from '../json-viewer' +import type { BodyWorkerRequest } from '../body-document-protocol' function bytes(value: string) { return new TextEncoder().encode(value).buffer } function gzip(value: string) { return Uint8Array.from(gzipSync(value)).buffer } describe('body document decoding', () => { + it('loads and copies complete captured bodies through the worker entry point', async () => { + const postMessage = vi.fn() + vi.stubGlobal('postMessage', postMessage) + vi.stubGlobal('onmessage', undefined) + try { + await import('../body-document.worker') + const dispatch = globalThis.onmessage as unknown as (event: { data: BodyWorkerRequest }) => Promise + const value = { messages: [{ role: 'user', content: `${'x'.repeat(100_000)}BODY-END` }] } + const text = JSON.stringify(value) + await dispatch({ data: { id: 1, action: 'load', bytes: gzip(text), encoding: 'gzip' } }) + expect(postMessage).toHaveBeenLastCalledWith({ id: 1, ok: true, result: { byteLength: bytes(text).byteLength } }) + await dispatch({ data: { id: 2, action: 'copy' } }) + expect(postMessage).toHaveBeenLastCalledWith({ id: 2, ok: true, result: JSON.stringify(value, null, 2) }) + } finally { + vi.unstubAllGlobals() + } + }) + it.each(['gzip', 'json'] as const)('decodes %s off the UI protocol with a byte count', async encoding => { const text = JSON.stringify({ text: '你好🙂', count: 0, enabled: false }) const decoded = await decodeBody(encoding === 'gzip' ? gzip(text) : bytes(text), encoding)