use futures::{Stream, StreamExt};
use serde_json::Value;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Provider {
Anthropic,
OpenAi,
Gemini,
}
impl Provider {
pub fn from_label(label: &str) -> Self {
match label {
"Anthropic" => Self::Anthropic,
"OpenAI" | "ChatGPT" => Self::OpenAi,
_ => Self::Gemini,
}
}
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct RealUsage {
pub model: String,
pub input_tokens: u64,
pub output_tokens: u64,
pub cache_read_tokens: u64,
pub cache_write_tokens: u64,
pub reasoning_tokens: u64,
pub provider_cost_usd: Option<f64>,
pub cohort: Option<super::holdout::Arm>,
pub wire: Option<Box<WireContext>>,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct WireContext {
pub provider: String,
pub person: Option<String>,
pub team: Option<String>,
pub project: Option<String>,
pub saved_tokens: u64,
pub uncompressed_input_tokens: u64,
pub is_local: bool,
pub routed_from: Option<String>,
pub counterfactual: Option<super::counterfactual::CounterfactualSlot>,
}
impl RealUsage {
fn is_meaningful(&self) -> bool {
!self.model.is_empty()
|| self.input_tokens > 0
|| self.output_tokens > 0
|| self.cache_read_tokens > 0
|| self.cache_write_tokens > 0
|| self.provider_cost_usd.is_some()
}
}
const MAX_LINE_BYTES: usize = 1 << 20;
pub struct Scanner {
provider: Provider,
url_model: Option<String>,
cohort: Option<super::holdout::Arm>,
wire: Option<Box<WireContext>>,
header_cost: Option<f64>,
buf: Vec<u8>,
usage: RealUsage,
}
impl Scanner {
pub fn new(provider: Provider, url_model: Option<String>) -> Self {
Self {
provider,
url_model,
cohort: None,
wire: None,
header_cost: None,
buf: Vec::new(),
usage: RealUsage::default(),
}
}
#[must_use]
pub fn with_cohort(mut self, cohort: Option<super::holdout::Arm>) -> Self {
self.cohort = cohort;
self
}
#[must_use]
pub fn with_wire_context(mut self, wire: Option<Box<WireContext>>) -> Self {
self.wire = wire;
self
}
#[must_use]
pub fn with_header_cost(mut self, cost: Option<f64>) -> Self {
self.header_cost = cost.filter(|c| c.is_finite() && *c >= 0.0);
self
}
pub fn feed(&mut self, chunk: &[u8]) {
self.buf.extend_from_slice(chunk);
while let Some(nl) = self.buf.iter().position(|&b| b == b'\n') {
let mut line: Vec<u8> = self.buf.drain(..=nl).collect();
line.pop(); if line.last() == Some(&b'\r') {
line.pop();
}
self.scan_line(&line);
}
if self.buf.len() > MAX_LINE_BYTES {
self.buf.clear();
}
}
pub fn feed_body(&mut self, body: &[u8]) {
if let Ok(v) = serde_json::from_slice::<Value>(body) {
self.absorb(&v);
}
}
pub fn finalize(mut self) -> Option<RealUsage> {
if !self.buf.is_empty() {
let line = std::mem::take(&mut self.buf);
self.scan_line(&line);
}
if self.usage.provider_cost_usd.is_none() {
self.usage.provider_cost_usd = self.header_cost;
}
if self.usage.is_meaningful() {
self.usage.cohort = self.cohort;
self.usage.wire = self.wire;
Some(self.usage)
} else {
None
}
}
fn scan_line(&mut self, line: &[u8]) {
let Ok(text) = std::str::from_utf8(line) else {
return;
};
let trimmed = text.trim();
if trimmed.is_empty() {
return;
}
if !self.line_might_be_relevant(trimmed) {
return;
}
let json_str = if let Some(rest) = trimmed.strip_prefix("data:") {
let r = rest.trim();
if r.is_empty() || r == "[DONE]" {
return;
}
r
} else if trimmed.starts_with('{') {
trimmed
.trim_start_matches([',', '['])
.trim_end_matches([',', ']'])
.trim()
} else {
return;
};
if let Ok(v) = serde_json::from_str::<Value>(json_str) {
self.absorb(&v);
}
}
fn line_might_be_relevant(&self, s: &str) -> bool {
match self.provider {
Provider::Anthropic | Provider::OpenAi => s.contains("usage"),
Provider::Gemini => s.contains("usageMetadata"),
}
}
fn absorb(&mut self, v: &Value) {
match self.provider {
Provider::Anthropic => absorb_anthropic(&mut self.usage, v),
Provider::OpenAi => absorb_openai(&mut self.usage, v),
Provider::Gemini => absorb_gemini(&mut self.usage, v, self.url_model.as_deref()),
}
}
}
fn absorb_anthropic(u: &mut RealUsage, v: &Value) {
let msg = v.get("message").unwrap_or(v);
if let Some(model) = msg.get("model").and_then(Value::as_str)
&& !model.is_empty()
{
u.model = model.to_string();
}
let Some(usage) = msg.get("usage").or_else(|| v.get("usage")) else {
return;
};
if let Some(n) = usage.get("input_tokens").and_then(Value::as_u64) {
u.input_tokens = n;
}
if let Some(n) = usage.get("cache_read_input_tokens").and_then(Value::as_u64) {
u.cache_read_tokens = n;
}
if let Some(n) = usage
.get("cache_creation_input_tokens")
.and_then(Value::as_u64)
{
u.cache_write_tokens = n;
}
if let Some(n) = usage.get("output_tokens").and_then(Value::as_u64)
&& n > 0
{
u.output_tokens = n;
}
}
fn absorb_openai(u: &mut RealUsage, v: &Value) {
let root = v.get("response").unwrap_or(v);
if let Some(model) = root.get("model").and_then(Value::as_str)
&& !model.is_empty()
{
u.model = model.to_string();
}
let Some(usage) = root.get("usage") else {
return;
};
if usage.is_null() {
return;
}
if let Some(cost) = usage.get("cost").and_then(Value::as_f64) {
let upstream = usage
.get("cost_details")
.and_then(|d| d.get("upstream_inference_cost"))
.and_then(Value::as_f64)
.unwrap_or(0.0);
let byok_upstream = if upstream > 0.0 && upstream != cost {
upstream
} else {
0.0
};
u.provider_cost_usd = Some(cost + byok_upstream);
}
let total_input = usage
.get("input_tokens")
.or_else(|| usage.get("prompt_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
let total_output = usage
.get("output_tokens")
.or_else(|| usage.get("completion_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
let input_details = usage
.get("input_tokens_details")
.or_else(|| usage.get("prompt_tokens_details"));
let cached = input_details
.and_then(|d| d.get("cached_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
let cache_write = input_details
.and_then(|d| d.get("cache_write_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
let reasoning = usage
.get("output_tokens_details")
.or_else(|| usage.get("completion_tokens_details"))
.and_then(|d| d.get("reasoning_tokens"))
.and_then(Value::as_u64)
.unwrap_or(0);
if total_input == 0 && total_output == 0 {
return;
}
u.input_tokens = total_input
.saturating_sub(cached)
.saturating_sub(cache_write);
u.cache_read_tokens = cached;
u.cache_write_tokens = cache_write;
u.output_tokens = total_output;
u.reasoning_tokens = reasoning;
}
fn absorb_gemini(u: &mut RealUsage, v: &Value, url_model: Option<&str>) {
if let Some(mv) = v.get("modelVersion").and_then(Value::as_str)
&& !mv.is_empty()
{
u.model = mv.to_string();
} else if u.model.is_empty()
&& let Some(m) = url_model
&& !m.is_empty()
{
u.model = m.to_string();
}
let Some(um) = v.get("usageMetadata") else {
return;
};
let prompt = um
.get("promptTokenCount")
.and_then(Value::as_u64)
.unwrap_or(0);
let candidates = um
.get("candidatesTokenCount")
.and_then(Value::as_u64)
.unwrap_or(0);
let cached = um
.get("cachedContentTokenCount")
.and_then(Value::as_u64)
.unwrap_or(0);
let thoughts = um
.get("thoughtsTokenCount")
.and_then(Value::as_u64)
.unwrap_or(0);
if prompt == 0 && candidates == 0 && thoughts == 0 {
return;
}
u.input_tokens = prompt.saturating_sub(cached);
u.cache_read_tokens = cached;
u.cache_write_tokens = 0;
u.output_tokens = candidates + thoughts;
u.reasoning_tokens = thoughts;
}
pub fn gemini_model_from_path(path: &str) -> Option<String> {
let after = path.rsplit_once("/models/").map(|(_, m)| m)?;
let model = after.split(':').next().unwrap_or(after).trim();
if model.is_empty() {
None
} else {
Some(model.to_string())
}
}
pub fn tee_stream<S, B, E>(
inner: S,
scanner: Scanner,
) -> impl Stream<Item = Result<B, E>> + Send + 'static
where
S: Stream<Item = Result<B, E>> + Send + Unpin + 'static,
B: AsRef<[u8]> + Send + 'static,
E: Send + 'static,
{
futures::stream::unfold(
(inner, Some(scanner)),
|(mut inner, mut scanner)| async move {
match inner.next().await {
Some(Ok(chunk)) => {
if let Some(s) = scanner.as_mut() {
s.feed(chunk.as_ref());
}
Some((Ok(chunk), (inner, scanner)))
}
Some(err) => Some((err, (inner, scanner))),
None => {
if let Some(s) = scanner.take()
&& let Some(usage) = s.finalize()
{
super::usage_meter::record(&usage);
}
None
}
}
},
)
}
#[cfg(test)]
mod tests {
use super::*;
fn feed_lines(
provider: Provider,
url_model: Option<&str>,
lines: &[&str],
) -> Option<RealUsage> {
let mut s = Scanner::new(provider, url_model.map(str::to_string));
for line in lines {
s.feed(line.as_bytes());
s.feed(b"\n");
}
s.finalize()
}
#[test]
fn anthropic_merges_message_start_and_delta() {
let u = feed_lines(
Provider::Anthropic,
None,
&[
r#"data: {"type":"message_start","message":{"model":"claude-opus-4-5-20251101","usage":{"input_tokens":100,"cache_read_input_tokens":2000,"cache_creation_input_tokens":50,"output_tokens":1}}}"#,
r#"data: {"type":"content_block_delta","index":0,"delta":{"text":"hello"}}"#,
r#"data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"output_tokens":73}}"#,
"data: {\"type\":\"message_stop\"}",
],
)
.expect("usage");
assert_eq!(u.model, "claude-opus-4-5-20251101");
assert_eq!(u.input_tokens, 100);
assert_eq!(u.cache_read_tokens, 2000);
assert_eq!(u.cache_write_tokens, 50);
assert_eq!(u.output_tokens, 73);
}
#[test]
fn anthropic_non_streaming_body() {
let mut s = Scanner::new(Provider::Anthropic, None);
s.feed_body(
br#"{"model":"claude-sonnet-4-5","usage":{"input_tokens":24,"output_tokens":18,"cache_creation_input_tokens":0,"cache_read_input_tokens":0}}"#,
);
let u = s.finalize().expect("usage");
assert_eq!(u.model, "claude-sonnet-4-5");
assert_eq!(u.input_tokens, 24);
assert_eq!(u.output_tokens, 18);
}
#[test]
fn scanner_stamps_wire_context_onto_usage() {
let wire = Box::new(WireContext {
provider: "Anthropic".into(),
person: Some("yves".into()),
team: None,
project: Some("billing".into()),
saved_tokens: 42,
uncompressed_input_tokens: 500,
is_local: false,
routed_from: None,
counterfactual: None,
});
let mut s = Scanner::new(Provider::Anthropic, None).with_wire_context(Some(wire.clone()));
s.feed_body(
br#"{"model":"claude-sonnet-4-5","usage":{"input_tokens":24,"output_tokens":18}}"#,
);
let u = s.finalize().expect("usage");
assert_eq!(u.wire, Some(wire));
}
#[test]
fn openai_responses_completed_event() {
let u = feed_lines(
Provider::OpenAi,
None,
&[
r#"data: {"type":"response.created","response":{"model":"gpt-5.4","usage":null}}"#,
r#"data: {"type":"response.completed","response":{"model":"gpt-5.4","usage":{"input_tokens":1289,"input_tokens_details":{"cached_tokens":289},"output_tokens":685,"output_tokens_details":{"reasoning_tokens":640},"total_tokens":1974}}}"#,
],
)
.expect("usage");
assert_eq!(u.model, "gpt-5.4");
assert_eq!(u.input_tokens, 1000); assert_eq!(u.cache_read_tokens, 289);
assert_eq!(u.output_tokens, 685);
assert_eq!(u.reasoning_tokens, 640);
}
#[test]
fn openai_chat_final_usage_chunk() {
let u = feed_lines(
Provider::OpenAi,
None,
&[
r#"data: {"choices":[{"delta":{"content":"hi"}}],"model":"gpt-5.4-mini"}"#,
r#"data: {"choices":[],"model":"gpt-5.4-mini","usage":{"prompt_tokens":500,"prompt_tokens_details":{"cached_tokens":100},"completion_tokens":40,"total_tokens":540}}"#,
"data: [DONE]",
],
)
.expect("usage");
assert_eq!(u.model, "gpt-5.4-mini");
assert_eq!(u.input_tokens, 400);
assert_eq!(u.cache_read_tokens, 100);
assert_eq!(u.output_tokens, 40);
assert_eq!(u.provider_cost_usd, None, "OpenAI reports no usage.cost");
}
#[test]
fn openrouter_cost_and_cache_writes_are_measured() {
let u = feed_lines(
Provider::OpenAi,
None,
&[
r#"data: {"choices":[{"delta":{"content":"hi"}}],"model":"deepseek/deepseek-v4-flash-20260423"}"#,
r#"data: {"choices":[],"model":"deepseek/deepseek-v4-flash-20260423","usage":{"prompt_tokens":700,"prompt_tokens_details":{"cached_tokens":150,"cache_write_tokens":50},"completion_tokens":40,"cost":0.0123,"cost_details":{"upstream_inference_cost":null},"total_tokens":740}}"#,
"data: [DONE]",
],
)
.expect("usage");
assert_eq!(u.input_tokens, 500, "700 - 150 cached - 50 cache-write");
assert_eq!(u.cache_read_tokens, 150);
assert_eq!(u.cache_write_tokens, 50);
assert_eq!(u.output_tokens, 40);
let cost = u.provider_cost_usd.expect("measured cost");
assert!((cost - 0.0123).abs() < 1e-12);
}
#[test]
fn openrouter_byok_adds_upstream_inference_cost() {
let mut s = Scanner::new(Provider::OpenAi, None);
s.feed_body(
br#"{"model":"anthropic/claude-sonnet-5","usage":{"prompt_tokens":100,"completion_tokens":10,"cost":0.05,"cost_details":{"upstream_inference_cost":0.95}}}"#,
);
let u = s.finalize().expect("usage");
let cost = u.provider_cost_usd.expect("measured cost");
assert!(
(cost - 1.0).abs() < 1e-12,
"OpenRouter fee + BYOK upstream bill"
);
}
#[test]
fn non_byok_upstream_equal_to_cost_is_not_doubled() {
let mut s = Scanner::new(Provider::OpenAi, None);
s.feed_body(
br#"{"model":"deepseek/deepseek-v4-flash","usage":{"prompt_tokens":50,"completion_tokens":5,"cost":0.000001568,"cost_details":{"upstream_inference_cost":0.000001568}}}"#,
);
let u = s.finalize().expect("usage");
let cost = u.provider_cost_usd.expect("measured cost");
assert!(
(cost - 0.000001568).abs() < 1e-15,
"#746: must book cost once, not 2x; got {cost}"
);
}
#[test]
fn gateway_header_cost_fills_when_body_has_none() {
let mut s = Scanner::new(Provider::OpenAi, None).with_header_cost(Some(0.0042));
s.feed_body(
br#"{"model":"azure-gpt-4o","usage":{"prompt_tokens":100,"completion_tokens":10}}"#,
);
let u = s.finalize().expect("usage");
assert_eq!(u.provider_cost_usd, Some(0.0042), "header is measured");
}
#[test]
fn body_reported_cost_beats_the_header_figure() {
let mut s = Scanner::new(Provider::OpenAi, None).with_header_cost(Some(9.99));
s.feed_body(
br#"{"model":"deepseek/deepseek-v4-flash","usage":{"prompt_tokens":100,"completion_tokens":10,"cost":0.05}}"#,
);
let u = s.finalize().expect("usage");
assert_eq!(u.provider_cost_usd, Some(0.05), "body wins over header");
}
#[test]
fn openrouter_free_model_reports_zero_cost_as_measured() {
let mut s = Scanner::new(Provider::OpenAi, None);
s.feed_body(
br#"{"model":"poolside/laguna-xs-2.1:free","usage":{"prompt_tokens":80,"completion_tokens":20,"cost":0}}"#,
);
let u = s.finalize().expect("usage");
assert_eq!(u.provider_cost_usd, Some(0.0), "free is a price, not a gap");
}
#[test]
fn gemini_usage_metadata_with_url_model() {
let u = feed_lines(
Provider::Gemini,
Some("gemini-2.5-pro"),
&[
r#"data: {"candidates":[{"content":{"parts":[{"text":"hi"}]}}],"usageMetadata":{"promptTokenCount":25,"candidatesTokenCount":7,"thoughtsTokenCount":39,"totalTokenCount":71}}"#,
],
)
.expect("usage");
assert_eq!(u.model, "gemini-2.5-pro");
assert_eq!(u.input_tokens, 25);
assert_eq!(u.output_tokens, 46); assert_eq!(u.reasoning_tokens, 39);
}
#[test]
fn gemini_prefers_model_version_over_url() {
let u = feed_lines(
Provider::Gemini,
Some("gemini-2.5-pro"),
&[
r#"data: {"modelVersion":"gemini-2.5-pro-002","usageMetadata":{"promptTokenCount":10,"candidatesTokenCount":5,"cachedContentTokenCount":4,"totalTokenCount":15}}"#,
],
)
.expect("usage");
assert_eq!(u.model, "gemini-2.5-pro-002");
assert_eq!(u.input_tokens, 6); assert_eq!(u.cache_read_tokens, 4);
assert_eq!(u.output_tokens, 5);
}
#[test]
fn split_chunks_reassemble_across_feed_calls() {
let mut s = Scanner::new(Provider::Anthropic, None);
let line = r#"data: {"type":"message_delta","delta":{},"usage":{"output_tokens":42}}"#;
let bytes = format!("{line}\n");
let (a, b) = bytes.as_bytes().split_at(20);
s.feed(a);
s.feed(b);
let u = s.finalize().expect("usage");
assert_eq!(u.output_tokens, 42);
}
#[test]
fn no_usage_yields_none() {
let out = feed_lines(
Provider::OpenAi,
None,
&[r#"data: {"choices":[{"delta":{"content":"hi"}}],"model":"gpt-5.4"}"#],
);
assert!(out.is_none(), "content-only stream reports no usage");
}
#[test]
fn final_event_without_trailing_newline() {
let mut s = Scanner::new(Provider::Anthropic, None);
s.feed(br#"data: {"type":"message_delta","delta":{},"usage":{"output_tokens":7}}"#);
let u = s.finalize().expect("flushes trailing partial line");
assert_eq!(u.output_tokens, 7);
}
#[test]
fn gemini_model_from_path_extracts() {
assert_eq!(
gemini_model_from_path("/v1beta/models/gemini-2.5-pro:streamGenerateContent")
.as_deref(),
Some("gemini-2.5-pro")
);
assert_eq!(
gemini_model_from_path("/v1beta/models/gemini-2.5-flash:generateContent").as_deref(),
Some("gemini-2.5-flash")
);
assert_eq!(gemini_model_from_path("/v1/chat/completions"), None);
}
#[tokio::test]
async fn tee_stream_passes_bytes_through_and_records() {
let chunks: Vec<Result<Vec<u8>, std::convert::Infallible>> = vec![
Ok(b"data: {\"type\":\"message_start\",\"message\":{\"model\":\"claude-haiku-4.5\",\"usage\":{\"input_tokens\":5,\"output_tokens\":1}}}\n".to_vec()),
Ok(b"data: {\"type\":\"message_delta\",\"delta\":{},\"usage\":{\"output_tokens\":9}}\n".to_vec()),
];
let inner = futures::stream::iter(chunks);
let scanner = Scanner::new(Provider::Anthropic, None);
let teed = tee_stream(inner, scanner);
let collected: Vec<_> = teed.collect().await;
assert_eq!(collected.len(), 2);
assert!(collected.iter().all(std::result::Result::is_ok));
}
}