From 125cd40aa5e15d057b529fee0dc33b4293c9cd04 Mon Sep 17 00:00:00 2001 From: dalamudx Date: Tue, 29 Sep 2026 22:57:39 +0800 Subject: [PATCH] =?UTF-8?q?feat(claude-code):=20=E8=A2=AB=E5=8A=A8?= =?UTF-8?q?=E9=87=87=E6=A0=B7=20anthropic-ratelimit-unified=20=E5=93=8D?= =?UTF-8?q?=E5=BA=94=E5=A4=B4=E6=9B=B4=E6=96=B0=E9=A2=9D=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/orchestration/report_effects.rs | 164 ++++++++++++++++++ crates/aether-admin/src/provider/quota.rs | 106 ++++++++++- 2 files changed, 269 insertions(+), 1 deletion(-) diff --git a/apps/aether-gateway/src/orchestration/report_effects.rs b/apps/aether-gateway/src/orchestration/report_effects.rs index 7e28a9e78..b86414637 100644 --- a/apps/aether-gateway/src/orchestration/report_effects.rs +++ b/apps/aether-gateway/src/orchestration/report_effects.rs @@ -551,6 +551,24 @@ async fn sync_grok_quota_from_report_context( async fn apply_local_sync_report_effect(state: &AppState, payload: &GatewaySyncReportRequest) { apply_local_gemini_file_mapping_report_effect(state, payload).await; + if claude_code_quota_headers_reportable(payload.status_code) { + if let Err(err) = sync_claude_code_quota_from_response_headers( + state, + payload.report_context.as_ref(), + &payload.headers, + ) + .await + { + warn!( + event_name = "claude_code_realtime_quota_sync_failed", + log_type = "ops", + report_kind = %payload.report_kind, + report_request_id = %short_request_id(report_request_id(payload.report_context.as_ref())), + error = ?err, + "gateway failed to persist claude_code realtime quota from sync response headers" + ); + } + } if (200..300).contains(&payload.status_code) { if let Err(err) = sync_codex_quota_from_response_headers( state, @@ -641,6 +659,24 @@ async fn apply_local_stream_report_effect(state: &AppState, payload: &GatewayStr ); } } + if claude_code_quota_headers_reportable(payload.status_code) { + if let Err(err) = sync_claude_code_quota_from_response_headers( + state, + payload.report_context.as_ref(), + &payload.headers, + ) + .await + { + warn!( + event_name = "claude_code_realtime_quota_sync_failed", + log_type = "ops", + report_kind = %payload.report_kind, + report_request_id = %short_request_id(report_request_id(payload.report_context.as_ref())), + error = ?err, + "gateway failed to persist claude_code realtime quota from stream response headers" + ); + } + } if let Err(err) = sync_grok_quota_from_report_context( state, payload.report_context.as_ref(), @@ -896,6 +932,134 @@ async fn sync_codex_quota_from_response_headers( .await } +fn claude_code_quota_headers_reportable(status_code: u16) -> bool { + // Real limit 429s carry the unified headers (fingerprint-rejection 429s do not, and then + // the parser finds nothing to record). + (200..300).contains(&status_code) || status_code == 429 +} + +/// Passive sampling of Anthropic's `anthropic-ratelimit-unified-*` response headers into the +/// `claude_code` quota metadata, so the 5H / weekly windows stay fresh between active refreshes. +async fn sync_claude_code_quota_from_response_headers( + state: &AppState, + report_context: Option<&Value>, + headers: &BTreeMap, +) -> Result { + let observed_at_unix_secs = report_context_u64( + report_context, + "provider_response_headers_observed_at_unix_ms", + ) + .map(|value| value / 1_000) + .filter(|value| *value > 0) + .unwrap_or_else(current_unix_secs); + let parsed = report_context_provider_response_headers(report_context) + .and_then(|headers| { + admin_provider_quota_pure::parse_claude_code_usage_headers( + &headers, + observed_at_unix_secs, + ) + }) + .or_else(|| { + admin_provider_quota_pure::parse_claude_code_usage_headers( + headers, + observed_at_unix_secs, + ) + }); + let Some(parsed) = parsed else { + return Ok(false); + }; + let Some(key_id) = report_context_key_id(report_context) else { + return Ok(false); + }; + + for attempt in 0..RUNTIME_METADATA_CAS_MAX_ATTEMPTS { + let Some(key) = state + .read_provider_catalog_keys_by_ids(std::slice::from_ref(&key_id)) + .await? + .into_iter() + .next() + else { + return Ok(false); + }; + let Some(provider) = state + .read_provider_catalog_providers_by_ids(std::slice::from_ref(&key.provider_id)) + .await? + .into_iter() + .next() + else { + return Ok(false); + }; + if !provider + .provider_type + .trim() + .eq_ignore_ascii_case("claude_code") + { + return Ok(false); + } + + let expected_namespace_value = + upstream_metadata_namespace_value(key.upstream_metadata.as_ref(), "claude_code"); + let mut bucket = expected_namespace_value + .as_ref() + .and_then(Value::as_object) + .cloned() + .unwrap_or_default(); + // Never let an older observation overwrite a newer active refresh. + if bucket + .get("updated_at") + .and_then(admin_provider_quota_pure::coerce_json_u64) + .is_some_and(|stored| stored > observed_at_unix_secs) + { + return Ok(false); + } + let Some(patch) = parsed.as_object() else { + return Ok(false); + }; + // Headers can describe only some windows; absent windows keep their stored value. + for (field, value) in patch { + bucket.insert(field.clone(), value.clone()); + } + let next_bucket = Value::Object(bucket); + if expected_namespace_value.as_ref() == Some(&next_bucket) { + return Ok(false); + } + + let updated_upstream_metadata = merge_metadata_object( + key.upstream_metadata.as_ref(), + "claude_code", + next_bucket.clone(), + ); + let updated_status_snapshot = sync_provider_key_quota_status_snapshot( + key.status_snapshot.as_ref(), + provider.provider_type.as_str(), + updated_upstream_metadata.as_ref(), + "response_headers", + ); + let updated = state + .update_provider_catalog_key_runtime_metadata( + &ProviderCatalogKeyRuntimeMetadataUpdate { + key_id: key_id.clone(), + namespace: "claude_code".to_string(), + expected_upstream_metadata_value: expected_namespace_value, + upstream_metadata_value: next_bucket, + status_snapshot_patch: quota_status_snapshot_patch( + updated_status_snapshot.as_ref(), + ), + updated_at_unix_secs: Some(observed_at_unix_secs), + }, + ) + .await?; + if updated { + return Ok(true); + } + if attempt + 1 < RUNTIME_METADATA_CAS_MAX_ATTEMPTS { + let backoff_us = 50_u64.saturating_mul((attempt + 1) as u64).min(1_000); + tokio::time::sleep(Duration::from_micros(backoff_us)).await; + } + } + Ok(false) +} + async fn sync_codex_websocket_quota_from_stream_payload( state: &AppState, payload: &GatewayStreamReportRequest, diff --git a/crates/aether-admin/src/provider/quota.rs b/crates/aether-admin/src/provider/quota.rs index f515fe156..43d0c65f6 100644 --- a/crates/aether-admin/src/provider/quota.rs +++ b/crates/aether-admin/src/provider/quota.rs @@ -466,6 +466,65 @@ pub fn parse_claude_code_oauth_usage_response( Some(serde_json::Value::Object(bucket)) } +/// Parses the `anthropic-ratelimit-unified-*` response headers into a partial `claude_code` +/// metadata bucket (same field names as [`parse_claude_code_oauth_usage_response`]). +/// +/// Utilization headers are 0-1 fractions and are stored as percent; reset headers are Unix +/// seconds (millisecond values are normalized). Returns `None` when no window is reported. +pub fn parse_claude_code_usage_headers( + headers: &BTreeMap, + updated_at_unix_secs: u64, +) -> Option { + let normalized = headers + .iter() + .map(|(key, value)| (key.trim().to_ascii_lowercase(), value.trim().to_string())) + .collect::>(); + let mut bucket = serde_json::Map::new(); + for (header_window, prefix) in [ + ("5h", "five_hour"), + ("7d", "seven_day"), + ("7d_oi", "seven_day_fable"), + ] { + let header = |suffix: &str| { + normalized + .get(&format!( + "anthropic-ratelimit-unified-{header_window}-{suffix}" + )) + .map(String::as_str) + }; + if let Some(utilization) = header("utilization") + .and_then(|value| value.parse::().ok()) + .filter(|value| value.is_finite()) + { + bucket.insert( + format!("{prefix}_used_percent"), + serde_json::json!((utilization * 100.0).clamp(0.0, 100.0)), + ); + } + if let Some(reset_at) = header("reset") + .and_then(|value| value.parse::().ok()) + .map(|value| { + if value > 100_000_000_000 { + value / 1_000 + } else { + value + } + }) + .filter(|value| *value > 0) + { + bucket.insert(format!("{prefix}_reset_at"), serde_json::json!(reset_at)); + } + } + if bucket.is_empty() { + return None; + } + bucket.insert( + "updated_at".to_string(), + serde_json::json!(updated_at_unix_secs), + ); + Some(serde_json::Value::Object(bucket)) +} + pub fn parse_gemini_cli_retrieve_user_quota_response( value: &serde_json::Value, updated_at_unix_secs: u64, @@ -7769,8 +7828,53 @@ mod tests { #[cfg(test)] mod claude_code_quota_tests { - use super::parse_claude_code_oauth_usage_response; + use super::{parse_claude_code_oauth_usage_response, parse_claude_code_usage_headers}; use serde_json::json; + use std::collections::BTreeMap; + + #[test] + fn parses_unified_ratelimit_headers_into_percent_and_reset() { + let headers = BTreeMap::from([ + ( + "Anthropic-Ratelimit-Unified-5h-Utilization".to_string(), + "0.42".to_string(), + ), + ( + "anthropic-ratelimit-unified-5h-reset".to_string(), + "1800003600".to_string(), + ), + ( + "anthropic-ratelimit-unified-7d-utilization".to_string(), + "1.0".to_string(), + ), + ( + "anthropic-ratelimit-unified-7d-reset".to_string(), + "1800400000000".to_string(), + ), + ( + "anthropic-ratelimit-unified-7d_oi-utilization".to_string(), + "0.1".to_string(), + ), + ]); + let parsed = + parse_claude_code_usage_headers(&headers, 1_800_000_000).expect("headers should parse"); + + assert_eq!(parsed["five_hour_used_percent"], json!(42.0)); + assert_eq!(parsed["five_hour_reset_at"], json!(1_800_003_600u64)); + assert_eq!(parsed["seven_day_used_percent"], json!(100.0)); + // Millisecond timestamps are normalized to seconds. + assert_eq!(parsed["seven_day_reset_at"], json!(1_800_400_000u64)); + assert_eq!(parsed["seven_day_fable_used_percent"], json!(10.0)); + assert!(parsed.get("seven_day_fable_reset_at").is_none()); + assert_eq!(parsed["updated_at"], json!(1_800_000_000u64)); + } + + #[test] + fn unified_ratelimit_headers_absent_yield_none() { + let headers = + BTreeMap::from([("content-type".to_string(), "application/json".to_string())]); + assert!(parse_claude_code_usage_headers(&headers, 1).is_none()); + } #[test] fn parses_windows_and_skips_null_ones() {