diff --git a/crates/aisix-core/src/models/model.rs b/crates/aisix-core/src/models/model.rs index 286f429de..b954a82a9 100644 --- a/crates/aisix-core/src/models/model.rs +++ b/crates/aisix-core/src/models/model.rs @@ -310,6 +310,19 @@ pub struct Model { #[schemars(length(min = 1, max = 255))] pub pricing_key: Option, + /// Opaque control-plane-issued canonical, non-nil UUID that authorizes a + /// concrete wildcard-upstream model name for pricing. It is relevant only + /// to a direct-shaped wildcard model (chat or embedding); without it the + /// data plane deliberately omits the resolved model from terminal + /// telemetry so the control plane can leave the call unpriced rather than + /// trust mutable configuration. + #[serde(default, skip_serializing_if = "Option::is_none")] + #[schemars( + regex(pattern = "^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$"), + length(min = 36, max = 64) + )] + pub pricing_authority_id: Option, + /// Direct-model-only background health-check configuration. #[serde(default, skip_serializing_if = "Option::is_none")] pub background_model_check: Option, @@ -421,6 +434,9 @@ impl Model { if self.effort_mapping.take().is_some() { stripped.push("effort_mapping"); } + if self.pricing_authority_id.take().is_some() { + stripped.push("pricing_authority_id"); + } if self.pricing_key.take().is_some() { stripped.push("pricing_key"); } @@ -533,8 +549,8 @@ pub fn model_one_of() -> Value { /// [`Model::strip_kind_inapplicable`]). Kind policy (project decision): /// generic call knobs (`timeout`/`stream_timeout`/`retries`) resolve /// member → group → deployment default wherever a group slot exists; -/// model-specific knobs (`auto_prompt_caching`, `cost`, `pricing_key`) are -/// direct-only. +/// model-specific knobs (`auto_prompt_caching`, `cost`, `pricing_key`, +/// `pricing_authority_id`) are direct-shaped-only. pub fn model_one_of_strict() -> Value { model_one_of_variant(true) } @@ -559,6 +575,7 @@ fn model_one_of_variant(strict: bool) -> Value { "auto_prompt_caching", "cost", "pricing_key", + "pricing_authority_id", "effort_mapping", ], ); @@ -580,6 +597,7 @@ fn model_one_of_variant(strict: bool) -> Value { "auto_prompt_caching", "cost", "pricing_key", + "pricing_authority_id", "effort_mapping", ], ); @@ -592,6 +610,7 @@ fn model_one_of_variant(strict: bool) -> Value { "auto_prompt_caching", "cost", "pricing_key", + "pricing_authority_id", "effort_mapping", ], ); @@ -694,6 +713,23 @@ mod tests { assert_eq!(m.rate_limit.as_ref().unwrap().rpm, Some(100)); } + #[test] + fn pricing_authority_id_round_trips_only_when_set() { + let mut model: Model = serde_json::from_str(sample_json()).unwrap(); + assert!(model.pricing_authority_id.is_none()); + assert!(serde_json::to_value(&model) + .unwrap() + .get("pricing_authority_id") + .is_none()); + + model.pricing_authority_id = Some("a3ebdc63-e921-4323-a75c-3b911f950046".to_string()); + let encoded = serde_json::to_value(&model).unwrap(); + assert_eq!( + encoded["pricing_authority_id"], + serde_json::json!("a3ebdc63-e921-4323-a75c-3b911f950046") + ); + } + #[test] fn deserialises_stream_timeout_and_helpers_fold_zero() { let m: Model = serde_json::from_str( @@ -762,11 +798,20 @@ mod tests { "retries": 2, "timeout": 1000, "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0}, - "auto_prompt_caching": {"enabled": true} + "auto_prompt_caching": {"enabled": true}, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046" })); let mut stripped = group.strip_kind_inapplicable(); stripped.sort_unstable(); - assert_eq!(stripped, ["auto_prompt_caching", "cost", "retries"]); + assert_eq!( + stripped, + [ + "auto_prompt_caching", + "cost", + "pricing_authority_id", + "retries" + ] + ); assert!(group.retries.is_none() && group.cost.is_none()); assert_eq!(group.timeout, Some(1000)); // Semantic parent: timeout/retries are the group slots and stay. @@ -780,9 +825,13 @@ mod tests { }, "retries": 2, "timeout": 1000, - "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0} + "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0}, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046" })); - assert_eq!(sem.strip_kind_inapplicable(), ["cost"]); + assert_eq!( + sem.strip_kind_inapplicable(), + ["cost", "pricing_authority_id"] + ); assert_eq!(sem.retries, Some(2)); assert_eq!(sem.timeout, Some(1000)); // Direct: nothing strips. @@ -792,10 +841,15 @@ mod tests { "model_name": "gpt-4o", "provider_key_id": "pk-1", "retries": 2, - "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0} + "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0}, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046" })); assert!(direct.strip_kind_inapplicable().is_empty()); assert_eq!(direct.retries, Some(2)); + assert_eq!( + direct.pricing_authority_id.as_deref(), + Some("a3ebdc63-e921-4323-a75c-3b911f950046") + ); // Ensemble parent: the whole generic set strips (its own // deadline knob is `ensemble.timeout_ms`). let mut ens = load(serde_json::json!({ @@ -804,7 +858,9 @@ mod tests { "timeout": 1000, "stream_timeout": 500, "retries": 1, - "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0} + "cost": {"input_per_1k": 0.0, "output_per_1k": 0.0}, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + "pricing_key": "catalog-gpt" })); // Asserted WITHOUT a pre-sort: the strip output is already // lexicographic (a pure-strip loader row keeps the fields @@ -812,7 +868,14 @@ mod tests { let ens_stripped = ens.strip_kind_inapplicable(); assert_eq!( ens_stripped, - ["cost", "retries", "stream_timeout", "timeout"] + [ + "cost", + "pricing_authority_id", + "pricing_key", + "retries", + "stream_timeout", + "timeout" + ] ); } diff --git a/crates/aisix-core/src/models/schema.rs b/crates/aisix-core/src/models/schema.rs index 3b394ab0c..fa0a52c0c 100644 --- a/crates/aisix-core/src/models/schema.rs +++ b/crates/aisix-core/src/models/schema.rs @@ -845,6 +845,19 @@ pub fn model_root_schema(strict: bool) -> Value { .as_object_mut() .expect("model root schema is a JSON object") .insert("oneOf".to_string(), one_of); + if strict { + // CP-issued pricing authorities are canonical UUIDs, but UUID nil is + // not an authority. Keep this write-only so an already-projected + // legacy row still loads and the DP can safely emit it unpriced. + schema + .pointer_mut("/properties/pricing_authority_id") + .and_then(Value::as_object_mut) + .expect("model schema declares pricing_authority_id") + .insert( + "not".to_string(), + json!({"const": "00000000-0000-0000-0000-000000000000"}), + ); + } // `OnEmbeddingFailure` is `#[serde(untagged)]` with an object variant // (`{ "target": … }`): serde buffers untagged content and silently // swallows unknown fields inside it, invisible to the write path's @@ -6065,10 +6078,12 @@ mod tests { "display_name": "g", "routing": {"strategy": "failover", "targets": [{"model": "m"}]}, "retries": 3, - "cost": {"input_per_1k": 0.5, "output_per_1k": 1.5} + "cost": {"input_per_1k": 0.5, "output_per_1k": 1.5}, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046" }); let msg = validate_model(&group).unwrap_err().message; assert!(msg.contains("`cost`"), "{msg}"); + assert!(msg.contains("`pricing_authority_id`"), "{msg}"); assert!(msg.contains("`retries`"), "{msg}"); assert!(msg.contains("model group"), "{msg}"); assert!( @@ -6114,6 +6129,42 @@ mod tests { assert!(!msg.contains("semantic router"), "{msg}"); } + #[test] + fn pricing_authority_id_is_canonical_non_nil_on_write_and_lenient_on_read() { + let base = json!({ + "display_name": "catalog/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": "pk" + }); + let mut valid = base.clone(); + valid["pricing_authority_id"] = json!("a3ebdc63-e921-4323-a75c-3b911f950046"); + validate_model(&valid).expect("canonical authority UUID writes"); + + for invalid in [ + "00000000-0000-0000-0000-000000000000", + "A3EBDC63-E921-4323-A75C-3B911F950046", + "not-a-uuid", + ] { + let mut model = base.clone(); + model["pricing_authority_id"] = json!(invalid); + assert!( + validate_model(&model).is_err(), + "strict write must reject {invalid:?}" + ); + } + + let mut legacy = base; + legacy["pricing_authority_id"] = json!("00000000-0000-0000-0000-000000000000"); + validate_model_lenient(&legacy).expect("legacy nil authority still loads unpriced"); + + let schema = model_root_schema(true); + assert_eq!( + schema["properties"]["pricing_authority_id"]["maxLength"], + json!(64) + ); + } + #[test] fn model_non_dead_knob_failures_keep_the_generic_message() { // A failure that is NOT a dead knob must not be relabelled: the diff --git a/crates/aisix-core/tests/resource_schema_characterization.rs b/crates/aisix-core/tests/resource_schema_characterization.rs index 20bc05bf2..b90e33c14 100644 --- a/crates/aisix-core/tests/resource_schema_characterization.rs +++ b/crates/aisix-core/tests/resource_schema_characterization.rs @@ -1061,6 +1061,10 @@ const EXTRA_RELAXATIONS: &[(&str, &[&str])] = &[ "/oneOf/3/not/anyOf", "/properties/effort_mapping/additionalProperties/minLength", "/properties/effort_mapping/properties//minLength", + // CP may have projected an older nil authority. The strict contract + // rejects it, while a DP must keep loading that row and emit its + // usage unpriced during a rolling upgrade. + "/properties/pricing_authority_id/not", ], ), ]; diff --git a/crates/aisix-obs/src/usage.rs b/crates/aisix-obs/src/usage.rs index c0f02f9c3..6d50f3e6f 100644 --- a/crates/aisix-obs/src/usage.rs +++ b/crates/aisix-obs/src/usage.rs @@ -181,6 +181,23 @@ pub struct UsageEvent { #[serde(default, skip_serializing_if = "String::is_empty")] pub requested_model: String, + /// Concrete upstream model selected through a wildcard upstream template. + /// + /// `model_id` remains the configured wildcard row so policy and + /// attribution remain stable, while this value is the provider model name + /// that was actually dispatched. The control plane uses it only when the + /// configured row's `model_name` contains `*`, to look up its catalog + /// price; it is otherwise absent so an alias over a fixed upstream can + /// never override its configured pricing identity through telemetry. + #[serde(default, skip_serializing_if = "String::is_empty")] + pub resolved_pricing_model: String, + + /// Canonical non-nil CP-issued UUID paired with `resolved_pricing_model`. + /// Both values are absent for older DP configuration or any request that + /// did not dispatch a concrete wildcard-template model. + #[serde(default, skip_serializing_if = "String::is_empty")] + pub pricing_authority_id: String, + pub prompt_tokens: u32, pub completion_tokens: u32, @@ -1695,6 +1712,8 @@ mod tests { model_id: "mod-uuid".into(), api_key_id: "ak-uuid".into(), requested_model: "smart-group".into(), + resolved_pricing_model: "gpt-4o-2024-08-06".into(), + pricing_authority_id: "a3ebdc63-e921-4323-a75c-3b911f950046".into(), prompt_tokens: 12, completion_tokens: 34, upstream_latency_ms: 56, @@ -1709,6 +1728,8 @@ mod tests { // AISIX-Cloud#790: the client-sent alias rides next to model_id // so the dashboard can show the group a routed request used. assert!(json.contains(r#""requested_model":"smart-group""#)); + assert!(json.contains(r#""resolved_pricing_model":"gpt-4o-2024-08-06""#)); + assert!(json.contains(r#""pricing_authority_id":"a3ebdc63-e921-4323-a75c-3b911f950046""#)); assert!(json.contains(r#""prompt_tokens":12"#)); assert!(json.contains(r#""completion_tokens":34"#)); assert!(json.contains(r#""guardrail_blocked":false"#)); @@ -1767,6 +1788,7 @@ mod tests { assert!(!json.contains("reasoning_tokens")); assert!(!json.contains("cache_creation_tokens")); assert!(!json.contains("cache_read_tokens")); + assert!(!json.contains("pricing_authority_id")); assert!(!json.contains("provider_request_id")); assert!(!json.contains("provider_model_version")); assert!(!json.contains("finish_reason")); diff --git a/crates/aisix-proxy/src/attribution.rs b/crates/aisix-proxy/src/attribution.rs index 681726052..abef56338 100644 --- a/crates/aisix-proxy/src/attribution.rs +++ b/crates/aisix-proxy/src/attribution.rs @@ -83,6 +83,17 @@ pub(crate) struct Resolved { /// it at read time, so the pair is byte-identical to the one the /// success path emits. pub provider_key_id: String, + /// The concrete upstream model produced by the wildcard-resolution + /// branch, paired with the configured wildcard row that produced it. + /// Empty for exact model resolution, including a request that literally + /// names a wildcard row. Telemetry uses this only for an event that + /// actually dispatched that same row; it is not an access-log identity. + pub wildcard_pricing_model_id: String, + /// Canonical non-nil UUID issued by CP that makes the concrete model below + /// billable. This stays beside the captured model id so a terminal emitter + /// can either send the complete authority tuple or omit it entirely. + pub wildcard_pricing_authority_id: String, + pub wildcard_pricing_model: String, /// Which cache layer answered this request, once one has — `Some` /// exactly when the response came out of the cache. /// @@ -840,6 +851,55 @@ pub(crate) fn note_target(model: &Model, provider_key_id: &str) { }); } +/// Record the complete pricing authority a caller-addressed wildcard row +/// resolved to. +/// +/// This is deliberately separate from [`note_target`]: an exact request for +/// the literal wildcard row also has an upstream model name, but it never +/// passed wildcard capture and must not be used as a pricing identity. The +/// authority is optional for rolling upgrades; without it, retain nothing so +/// the terminal event cannot assert a concrete wildcard price. +pub(crate) fn note_wildcard_pricing_identity( + model_id: &str, + pricing_authority_id: Option<&str>, + concrete_model: &str, +) { + let Some(pricing_authority_id) = pricing_authority_id.filter(|id| !id.is_empty()) else { + return; + }; + if model_id.is_empty() + || !valid_wildcard_pricing_model(concrete_model) + || !valid_pricing_authority_id(pricing_authority_id) + { + return; + } + with(|r| { + r.wildcard_pricing_model_id = model_id.to_string(); + r.wildcard_pricing_authority_id = pricing_authority_id.to_string(); + r.wildcard_pricing_model = concrete_model.to_string(); + }); +} + +// Keep this in step with AISIX Cloud's model-pricing name bound. A wildcard +// capture can be a valid upstream model name while still being too long to +// become a safe CP pricing lookup key; omit the entire optional tuple so its +// terminal parent event remains observable and explicitly unpriced. +const MAX_WILDCARD_PRICING_MODEL_CHARS: usize = 120; + +fn valid_wildcard_pricing_model(model: &str) -> bool { + !model.is_empty() + && !model.contains('*') + && !model.contains('\0') + && model.chars().count() <= MAX_WILDCARD_PRICING_MODEL_CHARS +} + +fn valid_pricing_authority_id(id: &str) -> bool { + match uuid::Uuid::parse_str(id) { + Ok(parsed) => !parsed.is_nil() && parsed.to_string() == id, + Err(_) => false, + } +} + /// Overwrite the target half with what a CACHE HIT may honestly claim. /// /// A hit contacts no upstream, so nothing was dispatched to and the line @@ -858,7 +918,12 @@ pub(crate) fn note_target(model: &Model, provider_key_id: &str) { /// group, whose candidate is no more the producer than any other. pub(crate) fn note_cache_hit_entry(entry: &Model, hit_layer: &'static str) { note_target(entry, entry.provider_key_id.as_deref().unwrap_or_default()); - with(|r| r.cache_hit_layer = Some(hit_layer)); + with(|r| { + r.cache_hit_layer = Some(hit_layer); + r.wildcard_pricing_model_id.clear(); + r.wildcard_pricing_authority_id.clear(); + r.wildcard_pricing_model.clear(); + }); } /// What the current request has resolved, or `None` outside a request. diff --git a/crates/aisix-proxy/src/count_tokens.rs b/crates/aisix-proxy/src/count_tokens.rs index f15f4d068..37b54e808 100644 --- a/crates/aisix-proxy/src/count_tokens.rs +++ b/crates/aisix-proxy/src/count_tokens.rs @@ -1242,6 +1242,79 @@ mod tests { ); } + /// A wildcard alias still resolves and dispatches on the counting route, + /// but token counting measures a later inference request rather than + /// consuming model tokens itself. Its zero-token row therefore must not + /// carry the wildcard row's concrete pricing authority. + #[tokio::test] + async fn wildcard_count_tokens_routes_upstream_without_pricing_identity() { + use aisix_obs::UsageSink; + + let upstream = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/messages/count_tokens")) + .respond_with( + ResponseTemplate::new(200).set_body_json(serde_json::json!({"input_tokens": 42})), + ) + .expect(1) + .mount(&upstream) + .await; + + let snap = new_snap(&upstream.uri()); + let wildcard: Model = serde_json::from_value(serde_json::json!({ + "display_name": "anthropic/*", + "provider": "anthropic", + "model_name": "*", + "provider_key_id": PK_ID, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + snap.models + .insert(ResourceEntry::new("wildcard", wildcard, 1)); + snap.apikeys.insert(apikey_entry(&["*"])); + + let hub = Arc::new(Hub::new()); + hub.register_specialized( + "anthropic", + Arc::new(aisix_provider_anthropic::AnthropicBridge::new()), + ); + let handle = SnapshotHandle::new(snap); + let (tx, mut rx) = tokio::sync::mpsc::channel(8); + let app = crate::build_router( + crate::ProxyState::new(handle, hub, &cfg()) + .without_cache() + .with_usage_sink(UsageSink::new(tx)), + ); + + let requested_model = "anthropic/claude-haiku-4-5-20251001"; + let resp = app + .oneshot(make_req(serde_json::json!({ + "model": requested_model, + "messages": [{"role": "user", "content": "hello"}], + }))) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let received = upstream.received_requests().await.unwrap(); + assert_eq!(received.len(), 1); + assert_eq!(received[0].url.path(), "/v1/messages/count_tokens"); + let sent: serde_json::Value = serde_json::from_slice(&received[0].body).unwrap(); + assert_eq!(sent["model"], "claude-haiku-4-5-20251001"); + + let event = tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()) + .await + .expect("count_tokens must emit a usage event") + .expect("channel open"); + assert_eq!(event.operation, "count_tokens"); + assert_eq!(event.model_id, "wildcard"); + assert_eq!(event.requested_model, requested_model); + assert_eq!(event.prompt_tokens, 0); + assert_eq!(event.completion_tokens, 0); + assert!(event.pricing_authority_id.is_empty()); + assert!(event.resolved_pricing_model.is_empty()); + } + /// Same gate as `/v1/messages`: a body the scan parser rejects is /// refused only when a guardrail would have read it. An output-hook-only /// row resolves into the chain but never sees the request, so the body diff --git a/crates/aisix-proxy/src/jobs.rs b/crates/aisix-proxy/src/jobs.rs index dd0740b0f..f1b0a6084 100644 --- a/crates/aisix-proxy/src/jobs.rs +++ b/crates/aisix-proxy/src/jobs.rs @@ -144,6 +144,10 @@ pub(crate) struct JobTarget { pub pk_entry: Arc>, pub secret: String, pub adapter: Adapter, + /// The exact, authorized routing selector that dispatched this request. + /// It is embedded in returned resource ids so later batch/file retrievals + /// resolve the same concrete wildcard model rather than the `*` row. + routing_model: String, /// The ProviderKey's rendered `default_headers`, resolved once when the /// target is resolved so every round-trip on this surface (upload, /// poll, download) sends the same set (AISIX-Cloud#1112). @@ -161,6 +165,9 @@ impl JobTarget { fn display_name(&self) -> &str { &self.model_entry.value.display_name } + fn routing_model(&self) -> &str { + &self.routing_model + } fn provider_label(&self) -> &str { self.model_entry.value.provider.as_deref().unwrap_or("") } @@ -199,6 +206,7 @@ pub(crate) fn resolve_target( wanted: Option<&str>, client_ctx: &ClientContext, ) -> Result { + let routing_model = wanted.map(str::to_string); let model_entry = match wanted { Some(name) => { let entry = crate::model_resolve::resolve_model(snapshot, name) @@ -265,7 +273,9 @@ pub(crate) fn resolve_target( ); let forwarded_client = aisix_gateway::ForwardedClientHeaders::resolve(&header_ctx); let extra_headers = aisix_gateway::resolve_default_headers(&header_ctx); + let routing_model = routing_model.unwrap_or_else(|| model.display_name.clone()); Ok(JobTarget { + routing_model, model_entry, pk_entry, secret, @@ -1035,7 +1045,7 @@ pub(crate) async fn create_file( &request_id, ) .await?; - let model = target.display_name().to_string(); + let model = target.routing_model().to_string(); Ok(( json_response(status, &resp_headers, bytes, Some(&model)), target, @@ -1308,7 +1318,7 @@ pub(crate) async fn create_batch( &mut applied, ) .await?; - let model = target.display_name().to_string(); + let model = target.routing_model().to_string(); Ok(( json_response(status, &resp_headers, bytes, Some(&model)), target, @@ -1402,7 +1412,7 @@ pub(crate) async fn get_batch( } } - let model = target.display_name().to_string(); + let model = target.routing_model().to_string(); Ok(( json_response(status, &resp_headers, bytes, Some(&model)), target, @@ -1610,7 +1620,7 @@ pub(crate) async fn create_ft_job( &mut applied, ) .await?; - let model = target.display_name().to_string(); + let model = target.routing_model().to_string(); Ok(( json_response(status, &resp_headers, bytes, Some(&model)), target, @@ -1841,7 +1851,7 @@ async fn forward_simple( } resp } else { - let model = spec.rewrite_ids.then(|| target.display_name().to_string()); + let model = spec.rewrite_ids.then(|| target.routing_model().to_string()); json_response(status, &resp_headers, bytes, model.as_deref()) }; Ok((resp, target)) @@ -1923,6 +1933,11 @@ fn maybe_attribute_batch( let jwt = auth.jwt.clone(); let model_id = target.model_entry.id.clone(); let display_name = target.display_name().to_string(); + // This runs while the retrieval request still owns its attribution + // scope. The download/aggregation task below is detached, so it must + // carry the verified dispatch-time tuple rather than trying to derive a + // price from the provider's output line or a later snapshot. + let wildcard_pricing = crate::usage_attr::capture_wildcard_pricing_identity(); // Resolved before the spawn, off the live snapshot, through the same // index `least_cost` ranks with: `pricing_key` first, inline `cost` // second. The completed batch is priced at what the model costs when @@ -1950,6 +1965,7 @@ fn maybe_attribute_batch( user_name.as_deref(), &model_id, &display_name, + wildcard_pricing, cost.as_ref(), &pk_id, &secret, @@ -1984,6 +2000,7 @@ async fn attribute_batch_usage( user_name: Option<&str>, model_id: &str, display_name: &str, + wildcard_pricing: Option, cost: Option<&aisix_core::models::model::ModelCost>, pk_id: &str, secret: &str, @@ -2028,7 +2045,9 @@ async fn attribute_batch_usage( } let body = resp.bytes().await.map_err(|e| e.to_string())?; - // Aggregate per provider-billed model (`response.body.model`). + // Keep the provider-reported model only as diagnostic version data. The + // captured dispatch tuple, not an upstream output field, selects the + // wildcard catalog price for every completed batch slice. #[derive(Default)] struct Agg { prompt: u64, @@ -2069,8 +2088,6 @@ async fn attribute_batch_usage( let snap = state.snapshot.load(); let pk = crate::usage_attr::ResolvedPk::resolve(&snap, pk_id); - // Same exporter set for every model slice of one batch. - let exporters = crate::usage_attr::live_exporters(state, &snap); let multi = per_model.len() > 1; for (idx, (provider_model, agg)) in per_model.iter().enumerate() { let request_id = batch_attribution_request_id(raw_batch_id, idx, multi); @@ -2092,29 +2109,35 @@ async fn attribute_batch_usage( .map(|c| c.calculate(agg.prompt, agg.completion)) .unwrap_or(0.0), inbound_protocol: "batch".to_string(), - // Set here rather than by the emit chokepoint: this path - // deliberately bypasses it (no live request, so no trace - // bundle), and both labels still come from one constant. - operation: crate::operation::BATCH_COMPLETION.operation.to_string(), ..Default::default() }; + crate::usage_attr::apply_captured_wildcard_pricing_identity( + &mut event, + crate::operation::BATCH_COMPLETION, + /* dispatched */ true, + wildcard_pricing.as_ref(), + ); crate::usage_attr::apply_pk_telemetry(&mut event, &pk); // Attribution names the identity that observed completion — the // same caller the event's api_key_id already reflects. crate::usage_attr::apply_caller_identity(&mut event, jwt, user_id, user_name); let usage_model = crate::usage_attr::usage_event_model_label(&snap, &event.requested_model).into_owned(); - state.usage_sink.try_emit( - crate::operation::BATCH_COMPLETION.handler, - event.clone(), + // Completion attribution is a real terminal inference event even + // though it has no live request trace. Send it through the one + // chokepoint so CP telemetry and exporter fan-out receive the same + // captured pricing tuple. + crate::usage_attr::emit_usage( + state, + &snap, + crate::operation::BATCH_COMPLETION, + event, crate::usage_attr::usage_event_labels(&usage_model, &pk), + None, + None, + /* terminal */ true, + /* dispatched */ true, ); - // A background poll attributes usage after the fact — there is no - // live request and therefore no trace bundle; the exporter falls - // back to the legacy flat span (AISIX-Cloud#1279). - state - .otlp_fan_out - .fan_out(&event, None, None, exporters.iter().map(|e| &e.value)); tracing::info!( batch_id = %raw_batch_id, provider_model = %provider_model, @@ -2419,6 +2442,90 @@ mod tests { assert_eq!(&bytes[..], b"{\"custom_id\":\"r1\"}\n"); } + #[tokio::test] + async fn wildcard_selected_file_and_fine_tuning_management_events_are_unpriced() { + let upstream = MockServer::start().await; + Mock::given(wm_method("GET")) + .and(path("/v1/files/file-abc")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "id": "file-abc", + "object": "file" + }))) + .expect(1) + .mount(&upstream) + .await; + Mock::given(wm_method("GET")) + .and(path("/v1/fine_tuning/jobs/ftjob-9")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "id": "ftjob-9", + "object": "fine_tuning.job" + }))) + .expect(1) + .mount(&upstream) + .await; + + let snap = AisixSnapshot::new(); + snap.provider_keys.insert(openai_pk(PK_A, &upstream.uri())); + let wildcard: Model = serde_json::from_value(serde_json::json!({ + "display_name": "jobs/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": PK_A, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + snap.models + .insert(ResourceEntry::new("wildcard", wildcard, 1)); + snap.apikeys.insert(apikey_entry(&["*"])); + let (app, mut rx) = build_app_with_sink(snap); + + for (management_kind, path, operation) in [ + ( + "file", + format!( + "/v1/files/{}", + encode_routed_id("file-abc", "jobs/gpt-4o-2024-08-06") + ), + "files", + ), + ( + "fine-tuning", + format!( + "/v1/fine_tuning/jobs/{}", + encode_routed_id("ftjob-9", "jobs/gpt-4o-2024-08-06") + ), + "fine_tuning", + ), + ] { + let req = Request::builder() + .method("GET") + .uri(path) + .header("authorization", "Bearer sk-caller") + .body(axum::body::Body::empty()) + .unwrap(); + let resp = app.clone().oneshot(req).await.unwrap(); + assert_eq!(resp.status(), StatusCode::OK, "{management_kind}"); + + let event = tokio::time::timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("management usage event must arrive") + .expect("usage sink must remain open"); + assert_eq!(event.model_id, "wildcard", "{management_kind}"); + assert_eq!(event.requested_model, "jobs/*", "{management_kind}"); + assert_eq!(event.operation, operation, "{management_kind}"); + assert_eq!(event.prompt_tokens, 0, "{management_kind}"); + assert_eq!(event.completion_tokens, 0, "{management_kind}"); + assert!( + event.pricing_authority_id.is_empty(), + "{management_kind} management event must not select wildcard pricing" + ); + assert!( + event.resolved_pricing_model.is_empty(), + "{management_kind} management event must not select wildcard pricing" + ); + } + } + // ---- batches ---- #[tokio::test] @@ -2445,11 +2552,18 @@ mod tests { snap.provider_keys .insert(openai_pk(PK_B, &upstream_b.uri())); snap.models.insert(model("m-a", "jobs-a", PK_A)); - snap.models.insert(model("m-b", "jobs-b", PK_B)); + let wildcard: Model = serde_json::from_value(serde_json::json!({ + "display_name": "jobs/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": PK_B, + })) + .unwrap(); + snap.models.insert(ResourceEntry::new("m-b", wildcard, 1)); snap.apikeys.insert(apikey_entry(&["*"])); let app = build_app(snap); - let encoded = encode_routed_id("file-realB", "jobs-b"); + let encoded = encode_routed_id("file-realB", "jobs/gpt-4o-2024-08-06"); let body = serde_json::json!({ "input_file_id": encoded, "endpoint": "/v1/chat/completions", @@ -2475,11 +2589,17 @@ mod tests { let v: Value = serde_json::from_slice(&bytes).unwrap(); assert_eq!( decode_routed_id(v["id"].as_str().unwrap()), - Some(("batch_777".to_string(), "jobs-b".to_string())) + Some(( + "batch_777".to_string(), + "jobs/gpt-4o-2024-08-06".to_string() + )) ); assert_eq!( decode_routed_id(v["input_file_id"].as_str().unwrap()), - Some(("file-realB".to_string(), "jobs-b".to_string())) + Some(( + "file-realB".to_string(), + "jobs/gpt-4o-2024-08-06".to_string() + )) ); assert!( upstream_a.received_requests().await.unwrap().is_empty(), @@ -2505,7 +2625,7 @@ mod tests { "id": "batch_req_1", "custom_id": "r1", "response": {"status_code": 200, "body": { - "model": "gpt-4o-2024-08-06", + "model": "provider-reported-batch-model", "usage": {"prompt_tokens": 10, "completion_tokens": 5, "prompt_tokens_details": {"cached_tokens": 2}} }} @@ -2514,7 +2634,7 @@ mod tests { "id": "batch_req_2", "custom_id": "r2", "response": {"status_code": 200, "body": { - "model": "gpt-4o-2024-08-06", + "model": "provider-reported-batch-model", "usage": {"prompt_tokens": 7, "completion_tokens": 3} }} }); @@ -2529,11 +2649,22 @@ mod tests { let snap = AisixSnapshot::new(); snap.provider_keys.insert(openai_pk(PK_A, &upstream.uri())); - snap.models.insert(model("m-a", "jobs-a", PK_A)); + let wildcard: Model = serde_json::from_value(serde_json::json!({ + "display_name": "jobs/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": PK_A, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + snap.models.insert(ResourceEntry::new("m-a", wildcard, 1)); snap.apikeys.insert(apikey_entry(&["*"])); let (app, mut rx) = build_app_with_sink(snap); - let encoded = encode_routed_id("batch_1", "jobs-a"); + // A concrete caller hint enters the wildcard-capture branch. The + // zero-token management event must still stay unpriced even though + // this GET genuinely contacts the upstream batch API. + let encoded = encode_routed_id("batch_1", "jobs/gpt-4o-2024-08-06"); let mk_req = || { Request::builder() .method("GET") @@ -2550,12 +2681,12 @@ mod tests { let v: Value = serde_json::from_slice(&bytes).unwrap(); assert_eq!( decode_routed_id(v["output_file_id"].as_str().unwrap()), - Some(("file-out".to_string(), "jobs-a".to_string())) + Some(("file-out".to_string(), "jobs/gpt-4o-2024-08-06".to_string())) ); // Two events expected: the zero-token management event plus ONE // aggregated batch event from the detached attribution task. - let mut mgmt = 0u32; + let mut mgmt: Option = None; let mut agg: Option = None; for _ in 0..2 { let ev = tokio::time::timeout(Duration::from_secs(3), rx.recv()) @@ -2565,10 +2696,21 @@ mod tests { if ev.inbound_protocol == "batch" { agg = Some(ev); } else { - mgmt += 1; + assert!( + mgmt.replace(ev).is_none(), + "only one management event expected" + ); } } - assert_eq!(mgmt, 1); + let mgmt = mgmt.expect("management event must be emitted"); + assert!( + mgmt.pricing_authority_id.is_empty(), + "a batch-management request must not select wildcard pricing" + ); + assert!( + mgmt.resolved_pricing_model.is_empty(), + "a batch-management request must not select wildcard pricing" + ); let agg = agg.expect("aggregated batch event must be emitted"); // cp-api accepts any visible-ASCII request_id since // AISIX-Cloud#1288, so this is no longer a wire constraint — but a @@ -2584,8 +2726,16 @@ mod tests { assert_eq!(agg.prompt_tokens, 17); assert_eq!(agg.completion_tokens, 8); assert_eq!(agg.cached_prompt_tokens, 2); - assert_eq!(agg.provider_model_version, "gpt-4o-2024-08-06"); - assert_eq!(agg.requested_model, "jobs-a"); + assert_eq!(agg.provider_model_version, "provider-reported-batch-model"); + assert_eq!(agg.requested_model, "jobs/*"); + assert_eq!( + agg.pricing_authority_id, "a3ebdc63-e921-4323-a75c-3b911f950046", + "the aggregate must retain the captured wildcard authority" + ); + assert_eq!( + agg.resolved_pricing_model, "gpt-4o-2024-08-06", + "the aggregate must use the dispatch model, not provider output" + ); // Second retrieve: management event only — the attribution is // process-deduped. diff --git a/crates/aisix-proxy/src/messages.rs b/crates/aisix-proxy/src/messages.rs index 8f3cad941..7e96bb561 100644 --- a/crates/aisix-proxy/src/messages.rs +++ b/crates/aisix-proxy/src/messages.rs @@ -5309,7 +5309,7 @@ data: [DONE]\n\n"; } #[tokio::test] - async fn non_anthropic_streaming_records_anthropic_usage_event_with_ttft() { + async fn wildcard_streaming_records_anthropic_usage_event_with_pricing_authority_and_ttft() { use aisix_obs::UsageSink; use aisix_provider_openai::OpenAiBridge; @@ -5331,7 +5331,16 @@ data: [DONE]\n\n"; let (tx, mut rx) = tokio::sync::mpsc::channel(4); let snap = new_snap_openai(&upstream.uri()); - snap.models.insert(openai_model("my-claude-alias")); + let wildcard: Model = serde_json::from_value(serde_json::json!({ + "display_name": "my-claude/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": OPENAI_PK_ID, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .expect("wildcard model parses"); + snap.models + .insert(ResourceEntry::new("m-wildcard", wildcard, 1)); snap.apikeys.insert(apikey_entry(&["*"])); let hub = Arc::new(Hub::new()); @@ -5344,7 +5353,7 @@ data: [DONE]\n\n"; let app = crate::build_router(state); let body = serde_json::json!({ - "model": "my-claude-alias", + "model": "my-claude/gpt-4o-2024-08-06", "messages": [{"role": "user", "content": "hi"}], "max_tokens": 100, "stream": true, @@ -5361,6 +5370,12 @@ data: [DONE]\n\n"; .expect("usage event was never emitted") .expect("usage event sender dropped"); assert_eq!(event.inbound_protocol, "anthropic"); + assert_eq!(event.model_id, "m-wildcard"); + assert_eq!( + event.pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(event.resolved_pricing_model, "gpt-4o-2024-08-06"); assert_eq!(event.prompt_tokens, 13); assert_eq!(event.completion_tokens, 4); assert_eq!(event.provider_request_id, "cmpl-359"); diff --git a/crates/aisix-proxy/src/model_resolve.rs b/crates/aisix-proxy/src/model_resolve.rs index c91d3a791..ae11ec14f 100644 --- a/crates/aisix-proxy/src/model_resolve.rs +++ b/crates/aisix-proxy/src/model_resolve.rs @@ -37,9 +37,22 @@ pub(crate) fn resolve_model( } let (entry, upstream) = best_wildcard_row(snapshot, requested)?; crate::attribution::note_requested_model(requested); - // A wildcard row is a direct model, and attribution stays on the ROW + // A wildcard row is direct-shaped, and attribution stays on the ROW // (see the module docs), so the synthetic clone below inherits its id. note_dispatchable_entry(&entry); + // Keep the concrete value separate from ordinary target attribution: + // an exact request for the literal wildcard row does not run this branch + // and must never make its static template eligible for pricing. This is + // deliberately decided against the dispatch snapshot: a later terminal + // emitter may see a refreshed configuration where the row changed or was + // deleted, but it must price the concrete model this request dispatched. + if wildcard_pricing_eligible(&entry.value, &upstream) { + crate::attribution::note_wildcard_pricing_identity( + &entry.id, + entry.value.pricing_authority_id.as_deref(), + &upstream, + ); + } let mut model = entry.value.clone(); model.model_name = Some(upstream); Some(Arc::new(ResourceEntry::new( @@ -76,8 +89,9 @@ fn best_wildcard_row( if !model.display_name.contains('*') { continue; } - // Only direct Models can serve a wildcard alias — routers / ensembles / - // semantic routers have no upstream `model_name` to dispatch. + // Only direct-shaped Models can serve a wildcard alias — routers / + // ensembles / semantic routers have no upstream `model_name` to + // dispatch. if model.is_routing() || model.is_ensemble() || model.is_semantic() { continue; } @@ -100,7 +114,7 @@ fn best_wildcard_row( /// body (a client-supplied video id) and must not be echoed back as though /// the gateway had attested it. pub(crate) fn row_serves_name(model: &Model, requested: &str) -> bool { - // The same kind gate `best_wildcard_row` applies: only a direct row can + // The same kind gate `best_wildcard_row` applies: only a direct-shaped row can // serve a caller-minted alias. Unreachable today on the one surface that // calls this — `dispatch::require_provider` rejects those kinds first — // but this sits beside the function it mirrors, and it judges an entry @@ -159,9 +173,26 @@ fn resolve_upstream_model_name(model: &Model, capture: &str) -> String { } } +/// Whether a wildcard dispatch produced a concrete provider-model identity +/// that is safe to send to CP for pricing. +/// +/// A wildcard display alias with a fixed upstream is already priced by its +/// configured model name. Only an upstream template with exactly one `*` +/// yields a caller-specific provider-model identity, and a literal `*` is +/// never a concrete catalog key. +fn wildcard_pricing_eligible(model: &Model, upstream: &str) -> bool { + model + .model_name + .as_deref() + .is_some_and(|template| template.bytes().filter(|&byte| byte == b'*').count() == 1) + && !upstream.is_empty() + && !upstream.contains('*') +} + #[cfg(test)] mod tests { use super::*; + use aisix_core::models::EmbeddingConfig; use aisix_core::snapshot::ResourceTable; fn direct_model(display_name: &str, model_name: Option<&str>) -> Model { @@ -174,6 +205,12 @@ mod tests { .unwrap() } + fn priced_direct_model(display_name: &str, model_name: Option<&str>) -> Model { + let mut model = direct_model(display_name, model_name); + model.pricing_authority_id = Some("a3ebdc63-e921-4323-a75c-3b911f950046".to_string()); + model + } + /// `row_serves_name` gates a name that did NOT arrive on a live request /// body — it decides whether the gateway will echo a caller-supplied /// string back as its own `model`. A row must accept every name it @@ -239,6 +276,104 @@ mod tests { assert_eq!(resolved.value.display_name, "openai/*"); } + /// Pricing eligibility belongs to the snapshot that dispatched the + /// request. In particular, a wildcard display alias over a fixed model, + /// a fixed direct or embedding model with a stale authority, a stale + /// multiple-star template, a literal wildcard row, a capture that is still + /// itself a wildcard, or an absent, nil, or noncanonical authority must + /// not manufacture a concrete provider-model price identity. + #[tokio::test] + async fn only_concrete_template_capture_with_authority_sets_wildcard_pricing_identity() { + use std::sync::Arc; + + let mut embedding = priced_direct_model("embedding", Some("text-embedding-3-small")); + embedding.embedding = Some(EmbeddingConfig { + dimensions: 4, + normalize: true, + }); + let multiple_stars = priced_direct_model("multiple/*", Some("gpt-*-*")); + assert!( + !wildcard_pricing_eligible(&multiple_stars, "gpt-4o"), + "an authority left on an unsupported multiple-star template must not mint pricing" + ); + let snap = snapshot_with(vec![ + ("wildcard", priced_direct_model("openrouter/*", Some("*"))), + ("fixed", priced_direct_model("fixed/*", Some("gpt-4o"))), + ("exact", priced_direct_model("exact", Some("gpt-4o-mini"))), + ("embedding", embedding), + ("multiple", multiple_stars), + ("unpriced", direct_model("unpriced/*", Some("*"))), + ]); + let mut nil_authority = direct_model("nil/*", Some("*")); + nil_authority.pricing_authority_id = Some("00000000-0000-0000-0000-000000000000".into()); + snap.models + .insert(ResourceEntry::new("nil", nil_authority, 1)); + let mut noncanonical_authority = direct_model("noncanonical/*", Some("*")); + noncanonical_authority.pricing_authority_id = + Some("A3EBDC63-E921-4323-A75C-3B911F950046".into()); + snap.models.insert(ResourceEntry::new( + "noncanonical", + noncanonical_authority, + 1, + )); + + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + resolve_model(&snap, "openrouter/gpt-4o-2024-08-06") + .expect("concrete wildcard request resolves"); + let resolved = crate::attribution::current().expect("in request scope"); + assert_eq!(resolved.wildcard_pricing_model_id, "wildcard"); + assert_eq!( + resolved.wildcard_pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(resolved.wildcard_pricing_model, "gpt-4o-2024-08-06"); + }, + ) + .await; + + for requested in [ + "fixed/anything", + "exact", + "embedding", + "multiple/gpt-4o", + "openrouter/*", + "openrouter/gpt-*", + "unpriced/gpt-4o", + "nil/gpt-4o", + "noncanonical/gpt-4o", + ] { + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + resolve_model(&snap, requested).expect("configured request resolves"); + let resolved = crate::attribution::current().expect("in request scope"); + assert!(resolved.wildcard_pricing_model_id.is_empty(), "{requested}"); + assert!( + resolved.wildcard_pricing_authority_id.is_empty(), + "{requested}" + ); + assert!(resolved.wildcard_pricing_model.is_empty(), "{requested}"); + }, + ) + .await; + } + + let overlong = format!("openrouter/{}", "界".repeat(121)); + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + resolve_model(&snap, &overlong).expect("overlong wildcard request still resolves"); + let resolved = crate::attribution::current().expect("in request scope"); + assert!(resolved.wildcard_pricing_model_id.is_empty()); + assert!(resolved.wildcard_pricing_authority_id.is_empty()); + assert!(resolved.wildcard_pricing_model.is_empty()); + }, + ) + .await; + } + #[test] fn wildcard_row_name_bounds_caller_minted_aliases() { let snap = snapshot_with(vec![ diff --git a/crates/aisix-proxy/src/realtime.rs b/crates/aisix-proxy/src/realtime.rs index 248866efe..fb5a53fe0 100644 --- a/crates/aisix-proxy/src/realtime.rs +++ b/crates/aisix-proxy/src/realtime.rs @@ -270,6 +270,11 @@ pub(crate) async fn realtime( Ok((ws, prep)) => { let state2 = state.clone(); let client2 = client.clone(); + // `on_upgrade` runs on a new Tokio task, which does not inherit + // the request task-local. Copy the verified wildcard price + // selection now; the HTTP cell is already owned by the completed + // upgrade response and cannot own the session terminal event. + let wildcard_pricing = crate::usage_attr::capture_wildcard_pricing_identity(); // `on_upgrade` runs the session on a detached task, so the // request span has to be attached to the future rather than // inherited — without it the session's guardrail checks log @@ -278,7 +283,16 @@ pub(crate) async fn realtime( ws.protocols(["realtime"]).on_upgrade(move |socket| { use tracing::Instrument as _; async move { - run_session(state2, prep, socket, client2, request_id, started).await; + run_session( + state2, + prep, + socket, + client2, + request_id, + started, + wildcard_pricing, + ) + .await; } .instrument(span) }) @@ -725,6 +739,28 @@ async fn run_session( client: ClientContext, request_id: String, started: Instant, + wildcard_pricing: Option, +) { + run_session_inner( + state, + prep, + client_ws, + client, + request_id, + started, + wildcard_pricing, + ) + .await; +} + +async fn run_session_inner( + state: ProxyState, + prep: Prepared, + client_ws: WebSocket, + client: ClientContext, + request_id: String, + started: Instant, + wildcard_pricing: Option, ) { let Prepared { auth, @@ -1093,6 +1129,12 @@ async fn run_session( guardrail_bypassed_reason: crate::usage_attr::bypass_reason(&audit), ..Default::default() }; + crate::usage_attr::apply_captured_wildcard_pricing_identity( + &mut event, + crate::operation::REALTIME, + /* dispatched */ true, + wildcard_pricing.as_ref(), + ); crate::usage_attr::apply_pk_telemetry(&mut event, &pk); crate::usage_attr::apply_caller_identity( &mut event, @@ -1582,6 +1624,60 @@ mod tests { assert_eq!(ev.api_key_id, "k-1"); } + /// `on_upgrade` moves the session to a task that does not inherit the + /// request task-local. A wildcard model resolved before the upgrade must + /// still reach its terminal session UsageEvent as the concrete provider + /// pricing identity, rather than being lost when the HTTP handler ends. + #[tokio::test] + async fn terminal_realtime_usage_keeps_wildcard_pricing_identity() { + let (up_addr, _handshake, _frames) = spawn_upstream().await; + let snap = snapshot(&format!("http://{up_addr}/v1"), "openai", "openai"); + let wildcard: Model = serde_json::from_value(serde_json::json!({ + "display_name": "rt-*", + "provider": "openai", + "model_name": "gpt-realtime-*", + "provider_key_id": PK_ID, + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + snap.models.insert(ResourceEntry::new("m-rt", wildcard, 2)); + let (addr, _state, mut rx) = serve(snap).await; + + let mut req = format!("ws://{addr}/v1/realtime?model=rt-2026-01-01") + .into_client_request() + .unwrap(); + req.headers_mut() + .insert("authorization", "Bearer sk-caller".parse().unwrap()); + let (ws, _) = tokio_tungstenite::connect_async(req) + .await + .expect("handshake"); + let (mut tx, mut client_rx) = ws.split(); + tx.send(TgMessage::Text( + serde_json::json!({"type": "session.update", "session": {"instructions": "hi"}}) + .to_string(), + )) + .await + .unwrap(); + + while let Some(Ok(message)) = client_rx.next().await { + if matches!(message, TgMessage::Close(_)) { + break; + } + } + + let event = tokio::time::timeout(Duration::from_secs(3), rx.recv()) + .await + .expect("terminal usage event expected") + .expect("usage sink closed"); + assert_eq!(event.model_id, "m-rt"); + assert_eq!(event.requested_model, "rt-2026-01-01"); + assert_eq!( + event.pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(event.resolved_pricing_model, "gpt-realtime-2026-01-01"); + } + /// `/v1/realtime` builds its upstream handshake by hand rather than /// through the shared bridge pipeline, which is exactly where a /// per-request mechanism goes silently missing: the operator diff --git a/crates/aisix-proxy/src/request_metrics.rs b/crates/aisix-proxy/src/request_metrics.rs index e0a72d82d..099253445 100644 --- a/crates/aisix-proxy/src/request_metrics.rs +++ b/crates/aisix-proxy/src/request_metrics.rs @@ -171,7 +171,9 @@ impl Default for Upstream<'_> { /// Whether the caller addressed an ensemble model. See [`LastTarget::new`]. fn is_ensemble(snap: &AisixSnapshot, requested_model: &str) -> bool { !requested_model.is_empty() - && crate::model_resolve::resolve_model(snap, requested_model) + && snap + .models + .get_by_name(requested_model) .is_some_and(|entry| entry.value.is_ensemble()) } @@ -702,6 +704,9 @@ mod tests { provider: "OpenAI".to_string(), upstream_model: "gpt-4o-mini".to_string(), provider_key_id: "pk-1".to_string(), + wildcard_pricing_model_id: String::new(), + wildcard_pricing_authority_id: String::new(), + wildcard_pricing_model: String::new(), cache_hit_layer: None, } } @@ -753,6 +758,62 @@ mod tests { assert_eq!(upstream.pk.name(), UNKNOWN); } + /// E2E metric recording loads the latest snapshot after dispatch. Its + /// ensemble check is classification only: resolving a wildcard again + /// there would replace the concrete pricing identity captured from the + /// snapshot that actually dispatched the request. + #[tokio::test] + async fn ensemble_metric_check_does_not_replace_captured_wildcard_identity() { + use std::sync::Arc; + + let dispatched = snapshot_with( + "wildcard", + serde_json::json!({ + "display_name": "openrouter/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": "pk-1", + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + }), + ); + let refreshed = snapshot_with( + "wildcard", + serde_json::json!({ + "display_name": "openrouter/*", + "provider": "openai", + "model_name": "replacement-*", + "provider_key_id": "pk-1", + }), + ); + + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + crate::model_resolve::resolve_model(&dispatched, "openrouter/gpt-4o") + .expect("wildcard model resolves at dispatch"); + let captured = + crate::attribution::current().expect("request attribution is installed"); + assert_eq!(captured.wildcard_pricing_model_id, "wildcard"); + assert_eq!( + captured.wildcard_pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(captured.wildcard_pricing_model, "gpt-4o"); + + assert!(!is_ensemble(&refreshed, "openrouter/gpt-4o")); + let after_metrics = + crate::attribution::current().expect("request attribution is installed"); + assert_eq!(after_metrics.wildcard_pricing_model_id, "wildcard"); + assert_eq!( + after_metrics.wildcard_pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(after_metrics.wildcard_pricing_model, "gpt-4o"); + }, + ) + .await; + } + /// A request that never selected a target keeps the placeholder, and the /// ProviderKey id must be `unknown` rather than the empty string an /// unresolved `ResolvedPk` reports verbatim — an empty label value would diff --git a/crates/aisix-proxy/src/usage_attr.rs b/crates/aisix-proxy/src/usage_attr.rs index adafbf983..e6eeada8c 100644 --- a/crates/aisix-proxy/src/usage_attr.rs +++ b/crates/aisix-proxy/src/usage_attr.rs @@ -439,6 +439,104 @@ pub(crate) fn metric_model_label_pair<'a>( } } +/// Whether a usage surface is a model-inference request whose dispatch +/// identity can name a concrete wildcard price. +/// +/// This is intentionally an allowlist rather than a list of the management +/// surfaces we currently know about. A new zero-token or control-plane route +/// must remain unpriced until it explicitly establishes the same billing +/// contract as a model-inference surface. `BATCH_COMPLETION` is the one +/// detached terminal surface: it may use an identity captured from the +/// completed batch's dispatch request, while the `BATCHES` management route +/// remains unpriced. +fn is_billable_inference_surface(surface: Surface) -> bool { + [ + crate::operation::CHAT, + crate::operation::COMPLETIONS, + crate::operation::MESSAGES, + crate::operation::RESPONSES, + crate::operation::EMBEDDINGS, + crate::operation::RERANK, + crate::operation::REALTIME, + crate::operation::IMAGE_GENERATION, + crate::operation::IMAGE_EDIT, + crate::operation::TRANSCRIPTION, + crate::operation::TRANSLATION, + crate::operation::SPEECH, + crate::operation::VIDEO_GENERATION, + crate::operation::BATCH_COMPLETION, + ] + .contains(&surface) +} + +/// A complete wildcard pricing selection captured at dispatch time. +/// +/// Detached terminal work cannot depend on a Tokio task-local surviving past +/// its originating request. Keep only the price-selection tuple: unlike the +/// rest of request attribution, it is safe and necessary to carry to the +/// terminal event, and its `model_id` prevents it being applied to another +/// attempt or model row. +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) struct WildcardPricingIdentity { + model_id: String, + pricing_authority_id: String, + resolved_pricing_model: String, +} + +/// Snapshot the current request's dispatch-time wildcard price selection for +/// a detached terminal emitter. Cache hits, fixed rows, incomplete legacy +/// projections, and code outside a request scope intentionally yield `None`. +pub(crate) fn capture_wildcard_pricing_identity() -> Option { + let resolved = crate::attribution::current()?; + if resolved.cache_hit_layer.is_some() + || resolved.wildcard_pricing_model_id.is_empty() + || resolved.wildcard_pricing_authority_id.is_empty() + || resolved.wildcard_pricing_model.is_empty() + { + return None; + } + Some(WildcardPricingIdentity { + model_id: resolved.wildcard_pricing_model_id, + pricing_authority_id: resolved.wildcard_pricing_authority_id, + resolved_pricing_model: resolved.wildcard_pricing_model, + }) +} + +/// Apply a wildcard pricing identity captured at dispatch time. The caller +/// still has to prove that this terminal event describes an upstream dispatch +/// on a priced inference surface; a captured tuple can never price a +/// management, cached, failed-before-dispatch, or different-model event. +pub(crate) fn apply_captured_wildcard_pricing_identity( + event: &mut UsageEvent, + surface: Surface, + dispatched: bool, + identity: Option<&WildcardPricingIdentity>, +) { + if !dispatched || !is_billable_inference_surface(surface) { + return; + } + let Some(identity) = identity else { + return; + }; + if identity.model_id != event.model_id { + return; + } + event.pricing_authority_id = identity.pricing_authority_id.clone(); + event.resolved_pricing_model = identity.resolved_pricing_model.clone(); +} + +/// Fill the optional DP-to-CP wildcard-pricing authority at the one usage +/// emission chokepoint. Model resolution establishes eligibility from the +/// dispatch snapshot and records a complete authority tuple, so this path must +/// not consult an emission-time snapshot that may have changed while a stream +/// or realtime session was still running. A cache hit or pre-dispatch failure +/// has no upstream call to price. Detached gateway work has no caller +/// attribution and therefore cannot price itself as the parent request. +fn apply_wildcard_pricing_model(event: &mut UsageEvent, surface: Surface, dispatched: bool) { + let identity = capture_wildcard_pricing_identity(); + apply_captured_wildcard_pricing_identity(event, surface, dispatched, identity.as_ref()); +} + /// Stamp the five per-PK attribution fields onto an in-progress UsageEvent, /// sanitising the operator-controlled tag strings (control-char strip + length /// cap) before they hit the wire. One source of truth for the mapping so the @@ -958,6 +1056,7 @@ pub(crate) fn emit_usage( if terminal && event.guardrail_blocked { state.metrics.record_guardrail_blocked_request(); } + apply_wildcard_pricing_model(&mut event, surface, dispatched); // Gateway-initiated semantic embeddings are child work of this request, // never attempts for the model the caller addressed. Attach their bounded // ledger only to the one terminal parent event: a retry can emit @@ -1118,6 +1217,528 @@ mod tests { assert_eq!((m.as_ref(), u.as_ref()), ("no-such/model", "raw-upstream")); } + #[test] + fn billable_inference_surface_allowlist_is_exhaustive() { + for surface in [ + crate::operation::CHAT, + crate::operation::COMPLETIONS, + crate::operation::MESSAGES, + crate::operation::RESPONSES, + crate::operation::EMBEDDINGS, + crate::operation::RERANK, + crate::operation::REALTIME, + crate::operation::IMAGE_GENERATION, + crate::operation::IMAGE_EDIT, + crate::operation::TRANSCRIPTION, + crate::operation::TRANSLATION, + crate::operation::SPEECH, + crate::operation::VIDEO_GENERATION, + crate::operation::BATCH_COMPLETION, + ] { + assert!( + is_billable_inference_surface(surface), + "{} must retain wildcard pricing attribution", + surface.operation, + ); + } + + for surface in [ + crate::operation::COUNT_TOKENS, + crate::operation::FILES, + crate::operation::BATCHES, + crate::operation::FINE_TUNING, + crate::operation::MCP, + crate::operation::A2A, + crate::operation::PASSTHROUGH, + ] { + assert!( + !is_billable_inference_surface(surface), + "{} must not select a wildcard pricing identity", + surface.operation, + ); + } + } + + #[test] + fn captured_wildcard_pricing_identity_only_prices_the_matching_dispatched_terminal() { + let identity = WildcardPricingIdentity { + model_id: "wildcard".to_string(), + pricing_authority_id: "a3ebdc63-e921-4323-a75c-3b911f950046".to_string(), + resolved_pricing_model: "gpt-4o-2024-08-06".to_string(), + }; + + let mut completed_batch = UsageEvent { + // NO-GUARDRAIL-CHAIN: unit-only value used to exercise the + // captured-identity gate, not an emitted gateway event. + model_id: "wildcard".to_string(), + ..Default::default() + }; + apply_captured_wildcard_pricing_identity( + &mut completed_batch, + crate::operation::BATCH_COMPLETION, + true, + Some(&identity), + ); + assert_eq!( + completed_batch.pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(completed_batch.resolved_pricing_model, "gpt-4o-2024-08-06"); + + for (surface, dispatched, model_id, case) in [ + (crate::operation::BATCHES, true, "wildcard", "management"), + ( + crate::operation::BATCH_COMPLETION, + false, + "wildcard", + "undispatched", + ), + ( + crate::operation::BATCH_COMPLETION, + true, + "other", + "different model", + ), + ] { + let mut event = UsageEvent { + // NO-GUARDRAIL-CHAIN: unit-only value used to exercise the + // captured-identity gate, not an emitted gateway event. + model_id: model_id.to_string(), + ..Default::default() + }; + apply_captured_wildcard_pricing_identity( + &mut event, + surface, + dispatched, + Some(&identity), + ); + assert!( + event.pricing_authority_id.is_empty(), + "{case} event must not select wildcard pricing" + ); + assert!( + event.resolved_pricing_model.is_empty(), + "{case} event must not select wildcard pricing" + ); + } + } + + #[tokio::test] + async fn request_attribution_stamps_only_concrete_wildcard_pricing_authorities() { + use aisix_core::resource::ResourceEntry; + use aisix_core::snapshot::ResourceTable; + + let table = ResourceTable::default(); + let configured: aisix_core::Model = serde_json::from_value(serde_json::json!({ + "display_name": "openrouter/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": "pk-1", + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + let cache_entry = configured.clone(); + table.insert(ResourceEntry::new("wildcard", configured, 1)); + let snap = AisixSnapshot { + models: table, + ..Default::default() + }; + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + let served = + crate::model_resolve::resolve_model(&snap, "openrouter/gpt-4o-2024-08-06") + .expect("wildcard model resolves"); + crate::attribution::note_target(&served.value, "pk-1"); + let attribution = + crate::attribution::current().expect("in request attribution scope"); + assert_eq!(attribution.wildcard_pricing_model_id, "wildcard"); + assert_eq!( + attribution.wildcard_pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(attribution.wildcard_pricing_model, "gpt-4o-2024-08-06"); + let mut event = UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused unit test constructs + // a synthetic pricing event, not a gateway request. + model_id: "wildcard".to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }; + apply_wildcard_pricing_model(&mut event, crate::operation::CHAT, true); + assert_eq!( + event.pricing_authority_id, + "a3ebdc63-e921-4323-a75c-3b911f950046" + ); + assert_eq!(event.resolved_pricing_model, "gpt-4o-2024-08-06"); + + for (surface, non_inference_kind) in [ + (crate::operation::COUNT_TOKENS, "count-tokens"), + (crate::operation::FILES, "file"), + (crate::operation::BATCHES, "batch"), + (crate::operation::FINE_TUNING, "fine-tuning"), + ] { + let mut management_event = UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused unit test constructs + // a synthetic pricing event, not a gateway request. + model_id: "wildcard".to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }; + apply_wildcard_pricing_model(&mut management_event, surface, true); + assert!( + management_event.pricing_authority_id.is_empty(), + "{non_inference_kind} non-inference event must not select wildcard pricing" + ); + assert!( + management_event.resolved_pricing_model.is_empty(), + "{non_inference_kind} non-inference event must not select wildcard pricing" + ); + } + + crate::attribution::note_cache_hit_entry(&cache_entry, "exact"); + let cached_attribution = + crate::attribution::current().expect("in request attribution scope"); + assert_eq!(cached_attribution.cache_hit_layer, Some("exact")); + assert_eq!(cached_attribution.upstream_model, "*"); + assert!( + capture_wildcard_pricing_identity().is_none(), + "a cache hit must not hand a detached emitter a wildcard price" + ); + // `note_cache_hit_entry` clears the identity above. Restore + // one here to pin the separate emission gate too: a cache + // hit is never billable as an upstream wildcard dispatch. + crate::attribution::note_wildcard_pricing_identity( + "wildcard", + Some("a3ebdc63-e921-4323-a75c-3b911f950046"), + "gpt-4o-2024-08-06", + ); + let mut cached_event = UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused unit test constructs + // a synthetic pricing event, not a gateway request. + model_id: "wildcard".to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }; + apply_wildcard_pricing_model(&mut cached_event, crate::operation::CHAT, true); + assert!(cached_event.pricing_authority_id.is_empty()); + assert!(cached_event.resolved_pricing_model.is_empty()); + }, + ) + .await; + + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + let literal = crate::model_resolve::resolve_model(&snap, "openrouter/*") + .expect("literal configured wildcard row resolves exactly"); + crate::attribution::note_target(&literal.value, "pk-1"); + let attribution = + crate::attribution::current().expect("in request attribution scope"); + assert!(attribution.wildcard_pricing_model_id.is_empty()); + assert!(attribution.wildcard_pricing_authority_id.is_empty()); + assert!(attribution.wildcard_pricing_model.is_empty()); + + let mut event = UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused unit test constructs + // a synthetic pricing event, not a gateway request. This + // exact alias is dispatchable but its static `*` template + // is never a concrete catalog price. + model_id: "wildcard".to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }; + apply_wildcard_pricing_model(&mut event, crate::operation::CHAT, true); + assert!(event.pricing_authority_id.is_empty()); + assert!(event.resolved_pricing_model.is_empty()); + }, + ) + .await; + + // A configuration projected by an older control plane has the + // wildcard template but no authority UUID. It may dispatch, but its + // terminal event must not claim a concrete catalog price. + let legacy_table = ResourceTable::default(); + let mut legacy_configured = cache_entry; + legacy_configured.pricing_authority_id = None; + legacy_table.insert(ResourceEntry::new("legacy", legacy_configured, 1)); + let legacy_snap = AisixSnapshot { + models: legacy_table, + ..Default::default() + }; + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + crate::model_resolve::resolve_model(&legacy_snap, "openrouter/gpt-4o-2024-08-06") + .expect("legacy wildcard model resolves"); + let attribution = + crate::attribution::current().expect("in request attribution scope"); + assert!(attribution.wildcard_pricing_model_id.is_empty()); + assert!(attribution.wildcard_pricing_authority_id.is_empty()); + assert!(attribution.wildcard_pricing_model.is_empty()); + + let mut event = UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused test constructs a + // synthetic terminal usage event after dispatch. + model_id: "legacy".to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }; + apply_wildcard_pricing_model(&mut event, crate::operation::CHAT, true); + let wire = serde_json::to_value(event).expect("usage event serialises"); + assert!(wire.get("pricing_authority_id").is_none()); + assert!(wire.get("resolved_pricing_model").is_none()); + }, + ) + .await; + + // A pointer left on a fixed direct or embedding row after a wildcard + // alias is changed must not broaden the wildcard-only contract. + let fixed_table = ResourceTable::default(); + let direct: aisix_core::Model = serde_json::from_value(serde_json::json!({ + "display_name": "fixed-chat", + "provider": "openai", + "model_name": "gpt-4o-mini", + "provider_key_id": "pk-1", + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + let embedding: aisix_core::Model = serde_json::from_value(serde_json::json!({ + "display_name": "fixed-embedding", + "provider": "openai", + "model_name": "text-embedding-3-small", + "provider_key_id": "pk-1", + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + "embedding": {"dimensions": 4}, + })) + .unwrap(); + fixed_table.insert(ResourceEntry::new("fixed-chat", direct, 1)); + fixed_table.insert(ResourceEntry::new("fixed-embedding", embedding, 1)); + let fixed_snap = AisixSnapshot { + models: fixed_table, + ..Default::default() + }; + for (model_id, surface) in [ + ("fixed-chat", crate::operation::CHAT), + ("fixed-embedding", crate::operation::EMBEDDINGS), + ] { + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + let served = crate::model_resolve::resolve_model(&fixed_snap, model_id) + .expect("fixed direct-shaped model resolves"); + crate::attribution::note_target(&served.value, "pk-1"); + let mut event = UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused test constructs a + // synthetic terminal usage event after dispatch. + model_id: model_id.to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }; + apply_wildcard_pricing_model(&mut event, surface, true); + assert_eq!(event.model_id, model_id); + assert!(event.pricing_authority_id.is_empty()); + assert!(event.resolved_pricing_model.is_empty()); + }, + ) + .await; + } + } + + /// A stream can outlive the configuration generation that dispatched it. + /// The terminal event must retain the concrete provider model captured at + /// dispatch, even if the same row becomes a fixed alias or disappears + /// before the terminal emit runs. + #[tokio::test] + async fn terminal_usage_keeps_dispatch_time_wildcard_identity_after_refresh_or_delete() { + use aisix_core::snapshot::{ResourceTable, SnapshotHandle}; + use aisix_core::ProxyConfig; + use aisix_gateway::Hub; + use aisix_obs::UsageSink; + + fn wildcard_snapshot(template: &str) -> AisixSnapshot { + let table = ResourceTable::default(); + let model: aisix_core::Model = serde_json::from_value(serde_json::json!({ + "display_name": "openrouter/*", + "provider": "openai", + "model_name": template, + "provider_key_id": "pk-1", + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + table.insert(ResourceEntry::new("wildcard", model, 1)); + AisixSnapshot { + models: table, + ..Default::default() + } + } + + let dispatched = wildcard_snapshot("*"); + let refreshed_fixed = wildcard_snapshot("gpt-4o"); + let deleted = AisixSnapshot::new(); + let cfg = ProxyConfig { + addr: "127.0.0.1:0".into(), + request_body_limit_bytes: 0, + real_ip: Default::default(), + request_id: Default::default(), + url_rewrites: Vec::new(), + tls: None, + listeners: Vec::new(), + thread_per_core: None, + workers: None, + }; + + for (case, terminal_snapshot) in [ + ("the row changed to a fixed alias", refreshed_fixed), + ("the row was deleted", deleted), + ] { + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + let state = ProxyState::new( + SnapshotHandle::new(terminal_snapshot.clone()), + Arc::new(Hub::new()), + &cfg, + ) + .with_usage_sink(UsageSink::new(tx)); + + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + crate::model_resolve::resolve_model( + &dispatched, + "openrouter/gpt-4o-2024-08-06", + ) + .expect("wildcard model resolves at dispatch"); + let pk = ResolvedPk::unresolved(); + emit_usage( + &state, + &terminal_snapshot, + crate::operation::CHAT, + UsageEvent { + // NO-GUARDRAIL-CHAIN: this focused test emits a + // synthetic terminal usage event after dispatch. + model_id: "wildcard".to_string(), + guardrail_bypassed_reason: String::new(), + applied_guardrails: Vec::new(), + ..Default::default() + }, + usage_event_labels("openrouter/*", &pk), + None, + None, + /* terminal */ true, + /* dispatched */ true, + ); + }, + ) + .await; + + let event = rx.try_recv().expect("terminal usage event emitted"); + assert_eq!( + event.resolved_pricing_model, "gpt-4o-2024-08-06", + "{case} must not rewrite the dispatched wildcard identity" + ); + assert_eq!( + event.pricing_authority_id, "a3ebdc63-e921-4323-a75c-3b911f950046", + "{case} must not rewrite the dispatched wildcard authority" + ); + } + } + + /// A failed attempt retains its target row so it can be observed, but a + /// target rate-limit or bridge-preparation refusal never reached a + /// provider and therefore must not claim a concrete wildcard price. + #[tokio::test] + async fn undispatched_wildcard_attempt_has_no_pricing_identity() { + use aisix_core::snapshot::{ResourceTable, SnapshotHandle}; + use aisix_core::ProxyConfig; + use aisix_gateway::Hub; + use aisix_obs::UsageSink; + + let table = ResourceTable::default(); + let configured: aisix_core::Model = serde_json::from_value(serde_json::json!({ + "display_name": "openrouter/*", + "provider": "openai", + "model_name": "*", + "provider_key_id": "pk-1", + "pricing_authority_id": "a3ebdc63-e921-4323-a75c-3b911f950046", + })) + .unwrap(); + table.insert(ResourceEntry::new("wildcard", configured, 1)); + let snap = AisixSnapshot { + models: table, + ..Default::default() + }; + let cfg = ProxyConfig { + addr: "127.0.0.1:0".into(), + request_body_limit_bytes: 0, + real_ip: Default::default(), + request_id: Default::default(), + url_rewrites: Vec::new(), + tls: None, + listeners: Vec::new(), + thread_per_core: None, + workers: None, + }; + let (tx, mut rx) = tokio::sync::mpsc::channel(1); + let state = ProxyState::new( + SnapshotHandle::new(snap.clone()), + Arc::new(Hub::new()), + &cfg, + ) + .with_usage_sink(UsageSink::new(tx)); + let client = ClientContext::default(); + let attempts = [crate::attempt::AttemptRecord { + index: 0, + kind: "initial", + target_model: String::new(), + target_model_id: "wildcard".to_string(), + provider_key_id: "pk-1".to_string(), + status: 429, + success: false, + error_class: "rate_limited".to_string(), + error_message: "target rate limit refused before dispatch".to_string(), + latency_ms: 0, + dispatched: false, + }]; + + crate::attribution::scope( + Arc::new(crate::attribution::RequestAttribution::default()), + async { + let served = + crate::model_resolve::resolve_model(&snap, "openrouter/gpt-4o-2024-08-06") + .expect("wildcard model resolves"); + crate::attribution::note_target(&served.value, "pk-1"); + emit_failed_attempts( + &state, + &snap, + crate::operation::CHAT, + "request-1746", + "openrouter/gpt-4o-2024-08-06", + "api-key", + &client, + &[], + &attempts, + true, + false, + Vec::new(), + crate::redact::RedactionCounts::new(), + &None, + ); + }, + ) + .await; + + let event = rx.try_recv().expect("failed attempt emits usage"); + assert_eq!(event.model_id, "wildcard"); + assert!(event.pricing_authority_id.is_empty()); + assert!(event.resolved_pricing_model.is_empty()); + } + /// AISIX-Cloud#1289: the id is upstream-controlled and reaches a log line /// and cp-api's `dpmgr_usage_events`. A newline in it would break the /// one-record-per-line shape every log consumer relies on, and an diff --git a/schemas/README.md b/schemas/README.md index 5b3938b9f..ef66d9cd9 100644 --- a/schemas/README.md +++ b/schemas/README.md @@ -116,7 +116,7 @@ so a consumer that models the lenient set as "the strict set with | `guardrail` | the `semantic` kind requires neither an embedding model (under either spelling) nor a threshold beside each example list | | `mcp_policy` | `allow` is not required; the `McpToolRef` relaxations above apply here too, as does the absent `deny`-beside-`deny_ids` guard (its team-scope guard is on both sets) | | `mcp_server` | the label pattern (`name`, and its former spelling `display_name`) still forbids `__` and a trailing `_`, but not a `*` | -| `model` | the per-kind `not`/`anyOf` lists that forbid a knob a kind never resolves are shorter — a stored row keeps loading and `Model::strip_kind_inapplicable` drops the dead knob; and an `effort_mapping` target value may be empty, which the write path refuses | +| `model` | the per-kind `not`/`anyOf` lists that forbid a knob a kind never resolves are shorter — a stored row keeps loading and `Model::strip_kind_inapplicable` drops the dead knob; an `effort_mapping` target value may be empty; and a legacy nil `pricing_authority_id` still loads so its usage can remain explicitly unpriced during a rolling upgrade | Three of those are worth spelling out. A half-written `McpToolRef` entry has to keep DESERIALIZING, not merely validating: the loader skips a row diff --git a/schemas/resources-lenient/model.schema.json b/schemas/resources-lenient/model.schema.json index 511df87a9..b5516b6a9 100644 --- a/schemas/resources-lenient/model.schema.json +++ b/schemas/resources-lenient/model.schema.json @@ -1365,6 +1365,13 @@ "minLength": 1, "type": "string" }, + "pricing_authority_id": { + "description": "Opaque control-plane-issued canonical, non-nil UUID that authorizes a concrete wildcard-upstream model name for pricing. It is relevant only to a direct-shaped wildcard model (chat or embedding); without it the data plane deliberately omits the resolved model from terminal telemetry so the control plane can leave the call unpriced rather than trust mutable configuration.", + "maxLength": 64, + "minLength": 36, + "pattern": "^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$", + "type": "string" + }, "pricing_key": { "description": "Name of a shared `pricing` document to take the per-token cost from, instead of setting `cost` on this model. The environment's own pricing documents are searched first and the shared catalog second; when the key matches neither, `cost` applies. Editing the pricing document repricies every model naming it, with no change to the models themselves.", "maxLength": 255, diff --git a/schemas/resources/model.schema.json b/schemas/resources/model.schema.json index 4beb8980b..094e646b4 100644 --- a/schemas/resources/model.schema.json +++ b/schemas/resources/model.schema.json @@ -1189,6 +1189,11 @@ "pricing_key" ] }, + { + "required": [ + "pricing_authority_id" + ] + }, { "required": [ "effort_mapping" @@ -1305,6 +1310,11 @@ "pricing_key" ] }, + { + "required": [ + "pricing_authority_id" + ] + }, { "required": [ "effort_mapping" @@ -1374,6 +1384,11 @@ "pricing_key" ] }, + { + "required": [ + "pricing_authority_id" + ] + }, { "required": [ "effort_mapping" @@ -1471,6 +1486,16 @@ "minLength": 1, "type": "string" }, + "pricing_authority_id": { + "description": "Opaque control-plane-issued canonical, non-nil UUID that authorizes a concrete wildcard-upstream model name for pricing. It is relevant only to a direct-shaped wildcard model (chat or embedding); without it the data plane deliberately omits the resolved model from terminal telemetry so the control plane can leave the call unpriced rather than trust mutable configuration.", + "maxLength": 64, + "minLength": 36, + "not": { + "const": "00000000-0000-0000-0000-000000000000" + }, + "pattern": "^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$", + "type": "string" + }, "pricing_key": { "description": "Name of a shared `pricing` document to take the per-token cost from, instead of setting `cost` on this model. The environment's own pricing documents are searched first and the shared catalog second; when the key matches neither, `cost` applies. Editing the pricing document repricies every model naming it, with no change to the models themselves.", "maxLength": 255, diff --git a/tests/e2e/src/cases/wildcard-pricing-telemetry-e2e.test.ts b/tests/e2e/src/cases/wildcard-pricing-telemetry-e2e.test.ts new file mode 100644 index 000000000..61e3bbc6f --- /dev/null +++ b/tests/e2e/src/cases/wildcard-pricing-telemetry-e2e.test.ts @@ -0,0 +1,573 @@ +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"; + +// Gateway/exporter boundary regression for AISIX-Cloud#1746: the configured +// wildcard row remains the model identity, while the concrete upstream model +// travels on the usage-export wire as `resolved_pricing_model`. This runs a +// real gateway process but `startMockSls` is only the exporter protocol +// receiver; it is not evidence for the CP/DPM/PostgreSQL ingestion contract, +// which belongs to AISIX-Cloud's mTLS telemetry E2E. + +const CALLER_PLAINTEXT = "sk-wildcard-pricing-telemetry-caller"; +const CALLER_KEY_HASH = createHash("sha256").update(CALLER_PLAINTEXT).digest("hex"); +const CREDENTIAL_REF = "wildcardpricing"; +const LOGSTORE = "wildcard-pricing-telemetry"; +const WILDCARD_ALIAS = "openrouter/*"; +const KNOWN_MODEL = "openai/gpt-4o-mini"; +const UNKNOWN_MODEL = "unknown/provider-model"; +const PRICING_AUTHORITY_ID = "a3ebdc63-e921-4323-a75c-3b911f950046"; +const EMBEDDING_WILDCARD_ALIAS = "embedding/*"; +const EMBEDDING_REQUEST_MODEL = "embedding/embedding-3-small"; +const EMBEDDING_UPSTREAM_MODEL = "text-embedding-3-small"; +const EMBEDDING_INPUT = "price this embedding"; +const EMBEDDING_VECTOR = [0.1, 0.2, 0.3]; +const COUNT_TOKENS_WILDCARD_ALIAS = "anthropic/*"; +const COUNT_TOKENS_REQUEST_MODEL = "anthropic/claude-haiku-4-5-20251001"; +const COUNT_TOKENS_UPSTREAM_MODEL = "claude-haiku-4-5-20251001"; +const STREAMING_WILDCARD_ALIAS = "stream/*"; +const STREAMING_REQUEST_MODEL = "stream/gpt-4o-2024-08-06"; +const STREAMING_UPSTREAM_MODEL = "gpt-4o-2024-08-06"; +const BATCH_WILDCARD_ALIAS = "batch/*"; +const BATCH_REQUEST_MODEL = "batch/gpt-4o-2024-08-06"; + +function upstreamResponse() { + return { + id: "chatcmpl-wildcard-pricing", + object: "chat.completion", + created: 0, + // Pricing must come from the dispatch attribution, never a provider + // response field that happens to look like a model identity. + model: "provider-response-model", + choices: [{ index: 0, message: { role: "assistant", content: "ok" }, finish_reason: "stop" }], + usage: { prompt_tokens: 10, completion_tokens: 5, total_tokens: 15 }, + }; +} + +function embeddingUpstreamResponse() { + return { + object: "list", + // Pricing must come from wildcard dispatch attribution rather than the + // provider response's optional model field. + model: "provider-response-embedding-model", + data: [{ object: "embedding", index: 0, embedding: EMBEDDING_VECTOR }], + usage: { prompt_tokens: 7, total_tokens: 7 }, + }; +} + +function streamEvents(): string[] { + return [ + JSON.stringify({ + id: "chatcmpl-wildcard-stream", + object: "chat.completion.chunk", + created: 0, + model: "provider-response-stream-model", + choices: [{ index: 0, delta: { role: "assistant", content: "ok" }, finish_reason: null }], + }), + JSON.stringify({ + id: "chatcmpl-wildcard-stream", + object: "chat.completion.chunk", + created: 0, + model: "provider-response-stream-model", + choices: [], + usage: { prompt_tokens: 11, completion_tokens: 6, total_tokens: 17 }, + }), + "[DONE]", + ]; +} + +function completedBatchResponse() { + return { + id: "batch-wildcard-completed", + object: "batch", + status: "completed", + input_file_id: "file-wildcard-input", + output_file_id: "file-wildcard-output", + }; +} + +function completedBatchOutput(): string { + return `${JSON.stringify({ + id: "batch-request-1", + custom_id: "request-1", + response: { + status_code: 200, + body: { + // The provider's response field is deliberately different: it is + // diagnostic only and must not override the dispatch price model. + model: "provider-reported-batch-model", + usage: { + prompt_tokens: 13, + completion_tokens: 7, + prompt_tokens_details: { cached_tokens: 3 }, + }, + }, + }, + })}\n`; +} + +function routedId(raw: string, model: string): string { + return `aisix-${Buffer.from(`${raw};model,${model}`).toString("base64url")}`; +} + +describe("wildcard pricing telemetry e2e", () => { + let app: SpawnedApp | undefined; + let sls: MockSls | undefined; + let upstream: OpenAiUpstream | undefined; + let embeddingUpstream: OpenAiUpstream | undefined; + let countTokensUpstream: OpenAiUpstream | undefined; + let streamingUpstream: OpenAiUpstream | undefined; + let batchUpstream: OpenAiUpstream | undefined; + let wildcardID = ""; + let embeddingWildcardID = ""; + let countTokensWildcardID = ""; + let streamingWildcardID = ""; + let batchWildcardID = ""; + let etcdReachable = false; + + beforeAll(async () => { + const etcd = new EtcdClient(); + etcdReachable = await etcd.ping(); + if (!etcdReachable) return; + + sls = await startMockSls(); + upstream = await startOpenAiUpstream({ nonStreamBody: upstreamResponse() }); + embeddingUpstream = await startOpenAiUpstream({ nonStreamBody: embeddingUpstreamResponse() }); + countTokensUpstream = await startOpenAiUpstream({ nonStreamBody: { input_tokens: 42 } }); + streamingUpstream = await startOpenAiUpstream({ streamEvents: streamEvents() }); + batchUpstream = await startOpenAiUpstream({ + scriptedResponses: [ + { nonStreamBody: completedBatchResponse() }, + { rawBody: completedBatchOutput(), rawContentType: "application/jsonl" }, + ], + }); + 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: "wildcard-pricing-sls", + enabled: true, + kind: "aliyun_sls", + endpoint: sls.url, + project: "aisix-e2e-obs", + logstore: LOGSTORE, + credential_ref: CREDENTIAL_REF, + content_mode: "metadata_only", + }); + const providerKey = await seed.createProviderKey({ + display_name: "wildcard-pricing-pk", + provider: "openrouter", + adapter: "openai", + secret: "sk-mock", + api_base: `${upstream.baseUrl}/v1`, + }); + const wildcard = await seed.createModel({ + display_name: WILDCARD_ALIAS, + provider: "openrouter", + model_name: "*", + provider_key_id: providerKey.id, + pricing_authority_id: PRICING_AUTHORITY_ID, + }); + wildcardID = wildcard.id; + const embeddingProviderKey = await seed.createProviderKey({ + display_name: "wildcard-pricing-embedding-pk", + provider: "openai", + adapter: "openai", + secret: "sk-mock", + api_base: `${embeddingUpstream.baseUrl}/v1`, + }); + const embeddingWildcard = await seed.createModel({ + display_name: EMBEDDING_WILDCARD_ALIAS, + provider: "openai", + model_name: "text-*", + provider_key_id: embeddingProviderKey.id, + pricing_authority_id: PRICING_AUTHORITY_ID, + embedding: { dimensions: EMBEDDING_VECTOR.length }, + }); + embeddingWildcardID = embeddingWildcard.id; + const countTokensProviderKey = await seed.createProviderKey({ + display_name: "wildcard-pricing-count-tokens-pk", + provider: "anthropic", + adapter: "anthropic", + secret: "sk-mock", + // The Anthropic bridge appends `/v1/messages/count_tokens` to its + // bare provider base URL. + api_base: countTokensUpstream.baseUrl, + }); + const countTokensWildcard = await seed.createModel({ + display_name: COUNT_TOKENS_WILDCARD_ALIAS, + provider: "anthropic", + model_name: "*", + provider_key_id: countTokensProviderKey.id, + pricing_authority_id: PRICING_AUTHORITY_ID, + }); + countTokensWildcardID = countTokensWildcard.id; + const streamingProviderKey = await seed.createProviderKey({ + display_name: "wildcard-pricing-stream-pk", + provider: "openai", + adapter: "openai", + secret: "sk-mock", + api_base: `${streamingUpstream.baseUrl}/v1`, + }); + const streamingWildcard = await seed.createModel({ + display_name: STREAMING_WILDCARD_ALIAS, + provider: "openai", + model_name: "*", + provider_key_id: streamingProviderKey.id, + pricing_authority_id: PRICING_AUTHORITY_ID, + }); + streamingWildcardID = streamingWildcard.id; + const batchProviderKey = await seed.createProviderKey({ + display_name: "wildcard-pricing-batch-pk", + provider: "openai", + adapter: "openai", + secret: "sk-mock", + api_base: `${batchUpstream.baseUrl}/v1`, + }); + const batchWildcard = await seed.createModel({ + display_name: BATCH_WILDCARD_ALIAS, + provider: "openai", + model_name: "*", + provider_key_id: batchProviderKey.id, + pricing_authority_id: PRICING_AUTHORITY_ID, + }); + batchWildcardID = batchWildcard.id; + + // Seeded last: a successful models-list gate proves that all preceding + // resources, including the exporter, are in the same gateway snapshot. + await seed.createApiKey({ key_hash: CALLER_KEY_HASH, allowed_models: ["*"] }); + await waitConfigPropagation(async () => { + const res = await fetch(`${app!.proxyUrl}/v1/models`, { + headers: { authorization: `Bearer ${CALLER_PLAINTEXT}` }, + }); + await res.arrayBuffer(); + return res.status === 200; + }); + }, 60_000); + + afterAll(async () => { + await app?.exit(); + await upstream?.close(); + await embeddingUpstream?.close(); + await countTokensUpstream?.close(); + await streamingUpstream?.close(); + await batchUpstream?.close(); + await sls?.close(); + }); + + async function requestModel(model: string): Promise { + if (!app) throw new Error("app not ready"); + const res = await fetch(`${app.proxyUrl}/v1/chat/completions`, { + method: "POST", + headers: { + authorization: `Bearer ${CALLER_PLAINTEXT}`, + "content-type": "application/json", + }, + body: JSON.stringify({ model, messages: [{ role: "user", content: "hi" }] }), + }); + const body = await res.text(); + expect(res.status, body).toBe(200); + } + + test("wildcard dispatch exports the concrete known and unknown pricing identities", async (ctx) => { + if (!etcdReachable || !app || !sls || !upstream || !wildcardID) { + ctx.skip(); + return; + } + + const knownRequest = `openrouter/${KNOWN_MODEL}`; + await requestModel(knownRequest); + expect(upstream.receivedRequests).toHaveLength(1); + expect(JSON.parse(upstream.receivedRequests[0]!.body)).toMatchObject({ model: KNOWN_MODEL }); + + const known = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("requested_model") === knownRequest, + `usage event for ${knownRequest}`, + ); + expect(known.get("model_id")).toBe(wildcardID); + expect(known.get("pricing_authority_id")).toBe(PRICING_AUTHORITY_ID); + expect(known.get("resolved_pricing_model")).toBe(KNOWN_MODEL); + expect(known.get("prompt_tokens")).toBe("10"); + expect(known.get("completion_tokens")).toBe("5"); + + const unknownRequest = `openrouter/${UNKNOWN_MODEL}`; + await requestModel(unknownRequest); + expect(upstream.receivedRequests).toHaveLength(2); + expect(JSON.parse(upstream.receivedRequests[1]!.body)).toMatchObject({ model: UNKNOWN_MODEL }); + + const unknown = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("requested_model") === unknownRequest, + `usage event for ${unknownRequest}`, + ); + expect(unknown.get("model_id")).toBe(wildcardID); + expect(unknown.get("pricing_authority_id")).toBe(PRICING_AUTHORITY_ID); + expect(unknown.get("resolved_pricing_model")).toBe(UNKNOWN_MODEL); + }); + + test("stream terminal telemetry keeps the wildcard dispatch price after the response scope ends", async (ctx) => { + if (!etcdReachable || !app || !sls || !streamingUpstream || !streamingWildcardID) { + ctx.skip(); + return; + } + + const before = streamingUpstream.receivedRequests.length; + const res = await fetch(`${app.proxyUrl}/v1/chat/completions`, { + method: "POST", + headers: { + authorization: `Bearer ${CALLER_PLAINTEXT}`, + "content-type": "application/json", + }, + body: JSON.stringify({ + model: STREAMING_REQUEST_MODEL, + stream: true, + messages: [{ role: "user", content: "stream this prompt" }], + }), + }); + const body = await res.text(); + expect(res.status, body).toBe(200); + expect(body).toContain("[DONE]"); + + const calls = streamingUpstream.receivedRequests.slice(before); + expect(calls).toHaveLength(1); + expect(calls[0]!.method).toBe("POST"); + expect(calls[0]!.path).toBe("/v1/chat/completions"); + expect(JSON.parse(calls[0]!.body)).toMatchObject({ + model: STREAMING_UPSTREAM_MODEL, + stream: true, + }); + + const event = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("requested_model") === STREAMING_REQUEST_MODEL, + `stream terminal usage event for ${STREAMING_REQUEST_MODEL}`, + ); + expect(event.get("model_id")).toBe(streamingWildcardID); + expect(event.get("pricing_authority_id")).toBe(PRICING_AUTHORITY_ID); + expect(event.get("resolved_pricing_model")).toBe(STREAMING_UPSTREAM_MODEL); + expect(event.get("prompt_tokens")).toBe("11"); + expect(event.get("completion_tokens")).toBe("6"); + }); + + test("direct embedding wildcard dispatch exports its concrete pricing identity", async (ctx) => { + if (!etcdReachable || !app || !sls || !embeddingUpstream || !embeddingWildcardID) { + ctx.skip(); + return; + } + + const baseline = embeddingUpstream.receivedRequests.length; + const res = await fetch(`${app.proxyUrl}/v1/embeddings`, { + method: "POST", + headers: { + authorization: `Bearer ${CALLER_PLAINTEXT}`, + "content-type": "application/json", + }, + body: JSON.stringify({ model: EMBEDDING_REQUEST_MODEL, input: EMBEDDING_INPUT }), + }); + const body = await res.text(); + expect(res.status, body).toBe(200); + expect(JSON.parse(body)).toMatchObject({ + object: "list", + model: EMBEDDING_REQUEST_MODEL, + data: [{ object: "embedding", index: 0, embedding: EMBEDDING_VECTOR }], + usage: { prompt_tokens: 7, total_tokens: 7 }, + }); + + const calls = embeddingUpstream.receivedRequests.slice(baseline); + expect(calls).toHaveLength(1); + expect(calls[0]!.method).toBe("POST"); + expect(calls[0]!.path).toBe("/v1/embeddings"); + expect(JSON.parse(calls[0]!.body)).toMatchObject({ + model: EMBEDDING_UPSTREAM_MODEL, + input: EMBEDDING_INPUT, + }); + + const event = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("requested_model") === EMBEDDING_REQUEST_MODEL, + `usage event for ${EMBEDDING_REQUEST_MODEL}`, + ); + expect(event.get("model_id")).toBe(embeddingWildcardID); + expect(event.get("pricing_authority_id")).toBe(PRICING_AUTHORITY_ID); + expect(event.get("resolved_pricing_model")).toBe(EMBEDDING_UPSTREAM_MODEL); + expect(event.get("prompt_tokens")).toBe("7"); + }); + + test("wildcard count_tokens remains unpriced after a real Anthropic dispatch", async (ctx) => { + if (!etcdReachable || !app || !sls || !countTokensUpstream || !countTokensWildcardID) { + ctx.skip(); + return; + } + + const baseline = countTokensUpstream.receivedRequests.length; + const res = await fetch(`${app.proxyUrl}/v1/messages/count_tokens`, { + method: "POST", + headers: { + "content-type": "application/json", + "x-api-key": CALLER_PLAINTEXT, + "anthropic-version": "2023-06-01", + }, + body: JSON.stringify({ + model: COUNT_TOKENS_REQUEST_MODEL, + messages: [{ role: "user", content: "count this prompt" }], + }), + }); + const body = await res.text(); + expect(res.status, body).toBe(200); + expect(JSON.parse(body)).toMatchObject({ input_tokens: 42 }); + + const calls = countTokensUpstream.receivedRequests.slice(baseline); + expect(calls).toHaveLength(1); + expect(calls[0]!.method).toBe("POST"); + expect(calls[0]!.path).toBe("/v1/messages/count_tokens"); + expect(JSON.parse(calls[0]!.body)).toMatchObject({ + model: COUNT_TOKENS_UPSTREAM_MODEL, + messages: [{ role: "user", content: "count this prompt" }], + }); + + const event = await waitForSlsLog( + sls, + LOGSTORE, + (log) => + log.get("operation") === "count_tokens" && + log.get("requested_model") === COUNT_TOKENS_REQUEST_MODEL, + `usage event for ${COUNT_TOKENS_REQUEST_MODEL}`, + ); + expect(event.get("model_id")).toBe(countTokensWildcardID); + expect(event.get("prompt_tokens")).toBe("0"); + expect(event.get("completion_tokens")).toBe("0"); + expect(event.get("cached_prompt_tokens")).toBeUndefined(); + expect(event.get("reasoning_tokens")).toBeUndefined(); + expect(event.get("total_tokens")).toBeUndefined(); + expect(event.get("pricing_authority_id")).toBeUndefined(); + expect(event.get("resolved_pricing_model")).toBeUndefined(); + }); + + test("wildcard-routed job management events remain unpriced", async (ctx) => { + if (!etcdReachable || !app || !sls || !upstream || !wildcardID) { + ctx.skip(); + return; + } + + const model = `openrouter/${KNOWN_MODEL}`; + const calls = [ + [ + "file", + "files", + `/v1/files/${routedId("file-wildcard", model)}`, + "/v1/files/file-wildcard", + ], + [ + "batch", + "batches", + `/v1/batches/${routedId("batch-wildcard", model)}`, + "/v1/batches/batch-wildcard", + ], + [ + "fine-tuning", + "fine_tuning", + `/v1/fine_tuning/jobs/${routedId("ftjob-wildcard", model)}`, + "/v1/fine_tuning/jobs/ftjob-wildcard", + ], + ] as const; + + for (const [kind, operation, gatewayPath, upstreamPath] of calls) { + const before = upstream.receivedRequests.length; + const res = await fetch(`${app.proxyUrl}${gatewayPath}`, { + headers: { authorization: `Bearer ${CALLER_PLAINTEXT}` }, + }); + const body = await res.text(); + expect(res.status, `${kind}: ${body}`).toBe(200); + + const forwarded = upstream.receivedRequests.slice(before); + expect(forwarded).toHaveLength(1); + expect(forwarded[0]!.method).toBe("GET"); + expect(forwarded[0]!.path).toBe(upstreamPath); + + const event = await waitForSlsLog( + sls, + LOGSTORE, + (log) => log.get("operation") === operation && log.get("model_id") === wildcardID, + `${kind} wildcard management usage event`, + ); + expect(event.get("requested_model")).toBe(WILDCARD_ALIAS); + expect(event.get("prompt_tokens")).toBe("0"); + expect(event.get("completion_tokens")).toBe("0"); + expect(event.get("pricing_authority_id")).toBeUndefined(); + expect(event.get("resolved_pricing_model")).toBeUndefined(); + } + }); + + test("completed batch telemetry retains dispatch pricing while its management event stays unpriced", async (ctx) => { + if (!etcdReachable || !app || !sls || !batchUpstream || !batchWildcardID) { + ctx.skip(); + return; + } + + const before = batchUpstream.receivedRequests.length; + const route = routedId("batch-wildcard-completed", BATCH_REQUEST_MODEL); + const res = await fetch(`${app.proxyUrl}/v1/batches/${route}`, { + headers: { authorization: `Bearer ${CALLER_PLAINTEXT}` }, + }); + const body = await res.text(); + expect(res.status, body).toBe(200); + expect(JSON.parse(body)).toMatchObject({ + output_file_id: routedId("file-wildcard-output", BATCH_REQUEST_MODEL), + }); + + const management = await waitForSlsLog( + sls, + LOGSTORE, + (log) => + log.get("operation") === "batches" && + log.get("model_id") === batchWildcardID && + log.get("requested_model") === BATCH_WILDCARD_ALIAS, + "wildcard batch management usage event", + ); + expect(management.get("prompt_tokens")).toBe("0"); + expect(management.get("completion_tokens")).toBe("0"); + expect(management.get("pricing_authority_id")).toBeUndefined(); + expect(management.get("resolved_pricing_model")).toBeUndefined(); + + const aggregate = await waitForSlsLog( + sls, + LOGSTORE, + (log) => + log.get("inbound_protocol") === "batch" && + log.get("model_id") === batchWildcardID && + log.get("requested_model") === BATCH_WILDCARD_ALIAS, + "completed wildcard batch usage event", + ); + expect(aggregate.get("provider_model_version")).toBe("provider-reported-batch-model"); + expect(aggregate.get("pricing_authority_id")).toBe(PRICING_AUTHORITY_ID); + expect(aggregate.get("resolved_pricing_model")).toBe("gpt-4o-2024-08-06"); + expect(aggregate.get("prompt_tokens")).toBe("13"); + expect(aggregate.get("completion_tokens")).toBe("7"); + expect(aggregate.get("cached_prompt_tokens")).toBe("3"); + + const calls = batchUpstream.receivedRequests.slice(before); + expect(calls).toHaveLength(2); + expect(calls[0]).toMatchObject({ method: "GET", path: "/v1/batches/batch-wildcard-completed" }); + expect(calls[1]).toMatchObject({ method: "GET", path: "/v1/files/file-wildcard-output/content" }); + }); +});