use tokio::sync::mpsc::Sender;
use crate::cloud::CloudEvent;
use crate::privacy::PrivacyFilter;
use super::billing::BillingMode;
use super::capture::{CaptureGap, CostBasis, PricingInputs, Usage};
use super::churn::ChurnFinding;
use super::session::Resolved;
use super::transforms::{TransformDecision, TRANSFORM_WIRE_KEYS};
use super::wire_format::WireFormat;
pub const ECON_TYPE: &str = "ai.openlatch.economics.usage";
#[derive(Clone, Debug)]
pub struct Observation {
pub measured: bool,
pub event_id: String,
pub occurred_at: String,
pub model: Option<String>,
pub model_known: bool,
pub billing: BillingMode,
pub install_id: String,
pub session: Resolved,
pub pricing: PricingInputs,
pub churn: Option<ChurnFinding>,
pub request_body_len: usize,
pub has_breakpoint: bool,
pub transform: Option<TransformDecision>,
pub optimize: Option<OptimizeEvidence>,
pub wire_format: WireFormat,
pub attributable_agent: Option<&'static str>,
}
impl Observation {
pub fn none() -> Self {
Observation {
measured: false,
event_id: String::new(),
occurred_at: String::new(),
model: None,
model_known: false,
billing: BillingMode::Unknown,
install_id: String::new(),
session: Resolved::unknown(),
pricing: PricingInputs::default(),
churn: None,
request_body_len: 0,
has_breakpoint: false,
transform: None,
optimize: None,
wire_format: WireFormat::Unknown,
attributable_agent: None,
}
}
fn resolved_agent_id(&self) -> String {
self.session
.agent_id
.clone()
.unwrap_or_else(|| self.install_id.clone())
}
fn resolved_source(&self) -> String {
self.session
.source
.clone()
.or_else(|| self.attributable_agent.map(str::to_string))
.unwrap_or_else(|| "unknown".to_string())
}
}
#[derive(Clone, Debug)]
pub struct OptimizeEvidence {
pub context: crate::zone_eval::OptimizeContext,
pub enforced: bool,
pub result: &'static str,
pub second_optimize_degraded: bool,
}
const GATEWAY_VENDOR_PREFIXES: &[(&str, &str)] = &[
("anthropic", "anthropic"),
("openai", "openai"),
("google", "gemini"),
("gemini", "gemini"),
("x-ai", "xai"),
("xai", "xai"),
("mistralai", "mistral"),
("mistral", "mistral"),
("deepseek", "deepseek"),
("cohere", "cohere"),
];
pub fn priced_identity(format_provider: &str, model: &str) -> (String, String) {
let model = model.trim();
let (provider, model) = match model.split_once('/') {
Some((vendor, rest)) if !rest.is_empty() => {
match GATEWAY_VENDOR_PREFIXES
.iter()
.find(|(prefix, _)| prefix.eq_ignore_ascii_case(vendor))
{
Some((_, provider)) => ((*provider).to_string(), rest.to_string()),
None if format_provider == "google" && vendor == "models" => {
("gemini".to_string(), rest.to_string())
}
None => (format_provider.to_string(), model.to_string()),
}
}
_ => (format_provider.to_string(), model.to_string()),
};
let provider = match provider.as_str() {
"google" | "gcp.gemini" | "gcp.vertex_ai" | "vertex_ai" => "gemini".to_string(),
_ => provider,
};
let model = if provider == "anthropic" {
model.replace('.', "-")
} else {
model
};
(provider, model)
}
pub fn assemble_data(
obs: &Observation,
usage: &Usage,
basis: CostBasis,
gap: Option<CaptureGap>,
cache_preserved: bool,
) -> serde_json::Value {
let agent_id = obs.resolved_agent_id();
let source = obs.resolved_source();
let (provider, model) = priced_identity(
obs.wire_format.provider(),
obs.model.as_deref().unwrap_or_default(),
);
let mut data = serde_json::json!({
"gen_ai.usage.input_tokens": usage.input_tokens,
"gen_ai.usage.cache_creation.input_tokens": usage.cache_write,
"gen_ai.usage.cache_read.input_tokens": usage.cache_read,
"gen_ai.usage.output_tokens": usage.output_tokens,
"ai.openlatch.cache.ephemeral_5m_input_tokens": usage.eph_5m,
"ai.openlatch.cache.ephemeral_1h_input_tokens": usage.eph_1h,
"gen_ai.request.model": model,
"gen_ai.provider.name": provider,
"ai.openlatch.cost.basis": basis.as_str(),
"ai.openlatch.billing.mode": obs.billing.as_str(),
"ai.openlatch.request.batch": obs.pricing.batch,
"ai.openlatch.request.fast_mode": obs.pricing.fast_mode,
"ai.openlatch.request.inference_geo": obs.pricing.inference_geo,
"ai.openlatch.session.agent_id": agent_id,
"ai.openlatch.session.source": source,
"ai.openlatch.session.install_id": obs.install_id,
"ai.openlatch.session.agent_session_id": obs.session.session_id,
crate::model_relay::session::ASSURANCE_KEY: obs.session.assurance.as_str(),
"ai.openlatch.event.id": obs.event_id,
"ai.openlatch.event.occurred_at": obs.occurred_at,
"ai.openlatch.cache.preserved": cache_preserved,
});
data["ai.openlatch.capture.gap"] = match gap {
Some(g) => serde_json::Value::String(g.as_str().to_string()),
None => serde_json::Value::Null,
};
let (offset, layer, class, block_index, byte_len, finding_id) = match &obs.churn {
Some(f) => (
serde_json::json!(f.divergence_offset),
serde_json::json!(f.churn_layer.as_str()),
serde_json::json!(f.churn_class.as_str()),
serde_json::json!(f.churn_block_index),
serde_json::json!(f.churn_byte_len),
serde_json::json!(f.finding_id),
),
None => (
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
serde_json::Value::Null,
),
};
data["ai.openlatch.prefix.divergence_offset"] = offset;
data["ai.openlatch.prefix.churn_layer"] = layer;
data["ai.openlatch.prefix.churn_class"] = class;
data["ai.openlatch.prefix.churn_block_index"] = block_index;
data["ai.openlatch.prefix.churn_byte_len"] = byte_len;
data["ai.openlatch.prefix.finding_id"] = finding_id;
match &obs.transform {
Some(t) => {
if let serde_json::Value::Object(fields) = t.to_wire_object() {
for (key, value) in fields {
data[key] = value;
}
}
}
None => {
for key in TRANSFORM_WIRE_KEYS {
data[key] = serde_json::Value::Null;
}
}
}
data
}
pub fn assemble_event(
obs: &Observation,
mut data: serde_json::Value,
privacy: &PrivacyFilter,
) -> CloudEvent {
crate::privacy::filter_event_with(&mut data, privacy);
let source = obs.resolved_source();
let agent_id = obs.resolved_agent_id();
let mut envelope = serde_json::json!({
"specversion": "1.0",
"id": crate::envelope::new_event_id(),
"source": source,
"type": ECON_TYPE,
"time": crate::envelope::current_timestamp(),
"datacontenttype": "application/json",
"data": data,
});
if obs.wire_format != WireFormat::Unknown {
envelope["wireformat"] = serde_json::json!(obs.wire_format.as_str());
envelope["olmodelprovider"] = serde_json::json!(
priced_identity(
obs.wire_format.provider(),
obs.model.as_deref().unwrap_or_default()
)
.0
);
}
envelope["olsessionassurance"] = serde_json::json!(obs.session.assurance.as_str());
if let Some(model) = obs.model.as_deref().filter(|m| !m.trim().is_empty()) {
let (_, model) = priced_identity(obs.wire_format.provider(), model);
envelope["olmodelslug"] = serde_json::json!(model);
}
if let Some(optimize) = &obs.optimize {
envelope["oloptimize"] = serde_json::to_value(&optimize.context).unwrap_or_default();
envelope["olenforced"] = serde_json::json!(i64::from(optimize.enforced));
envelope["olresult"] = serde_json::json!(optimize.result);
if optimize.second_optimize_degraded {
envelope["oloptimizedegraded"] = serde_json::json!("second_optimize");
}
}
if let Some(sid) = obs.session.session_id.as_deref() {
envelope["subject"] = serde_json::Value::String(sid.to_string());
}
CloudEvent { envelope, agent_id }
}
#[allow(clippy::too_many_arguments)]
pub fn build_and_emit(
obs: &Observation,
usage: &Usage,
basis: CostBasis,
gap: Option<CaptureGap>,
cache_preserved: bool,
privacy: &PrivacyFilter,
cloud_tx: Option<&Sender<CloudEvent>>,
) {
if !obs.measured {
return;
}
let Some(tx) = cloud_tx else {
return;
};
let data = assemble_data(obs, usage, basis, gap, cache_preserved);
let event = assemble_event(obs, data, privacy);
if tx.try_send(event).is_err() {
tracing::warn!(
code = crate::error::ERR_CLOUD_UNREACHABLE,
"model relay: cloud channel full — economics event dropped"
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::model_relay::session::Assurance;
fn obs() -> Observation {
Observation {
measured: true,
event_id: "0190aaaa-bbbb-cccc-dddd-eeeeeeeeeeee".to_string(),
occurred_at: "2026-07-23T12:00:00Z".to_string(),
model: Some("claude-opus-4-8".to_string()),
model_known: true,
billing: BillingMode::ApiKey,
install_id: "agt_install".to_string(),
session: Resolved {
agent_id: Some("agt_install".to_string()),
source: Some("claude-code".to_string()),
session_id: Some("sess_a".to_string()),
assurance: Assurance::Attested,
},
pricing: PricingInputs::default(),
churn: None,
request_body_len: 100,
has_breakpoint: true,
transform: None,
optimize: None,
wire_format: WireFormat::AnthropicMessages,
attributable_agent: None,
}
}
#[test]
fn the_row_carries_the_shared_facts_under_the_shared_names() {
let mut observation = obs();
observation.wire_format = WireFormat::GoogleGenerateContent;
observation.model = Some("gemini-3.1-flash-lite".to_string());
observation.session.assurance = Assurance::Inferred;
let ev = assemble_event(
&observation,
serde_json::json!({}),
&PrivacyFilter::new(&[]),
);
assert_eq!(ev.envelope["olsessionassurance"], "inferred");
assert_eq!(ev.envelope["olmodelslug"], "gemini-3.1-flash-lite");
assert_eq!(ev.envelope["olmodelprovider"], "gemini");
assert!(ev.envelope.get("olprompt").is_none());
observation.wire_format = WireFormat::Unknown;
let ev = assemble_event(
&observation,
serde_json::json!({}),
&PrivacyFilter::new(&[]),
);
assert!(ev.envelope.get("olmodelprovider").is_none());
assert_eq!(
ev.envelope["olsessionassurance"], "inferred",
"assurance does not depend on the route"
);
}
#[test]
fn the_row_carries_the_format_its_request_spoke() {
let mut observation = obs();
observation.wire_format = WireFormat::GoogleGenerateContent;
let ev = assemble_event(
&observation,
serde_json::json!({}),
&PrivacyFilter::new(&[]),
);
assert_eq!(ev.envelope["wireformat"], "google-generate-content");
observation.wire_format = WireFormat::Unknown;
let ev = assemble_event(
&observation,
serde_json::json!({}),
&PrivacyFilter::new(&[]),
);
assert!(
ev.envelope.get("wireformat").is_none(),
"an unobserved format is absent, never asserted"
);
}
#[test]
fn the_envelope_carries_the_session_id_as_subject() {
let ev = assemble_event(&obs(), serde_json::json!({}), &PrivacyFilter::new(&[]));
assert_eq!(
ev.envelope["subject"], "sess_a",
"an economics event with no subject belongs to no session"
);
}
#[test]
fn an_unattributed_request_names_no_subject() {
let mut o = obs();
o.session.session_id = None;
let ev = assemble_event(&o, serde_json::json!({}), &PrivacyFilter::new(&[]));
assert!(
ev.envelope.get("subject").is_none(),
"no session resolved must mean no subject, not a fabricated one"
);
}
#[test]
fn provider_name_follows_the_format() {
let usage = Usage::default();
let mut anthropic = obs();
anthropic.wire_format = WireFormat::AnthropicMessages;
let data = assemble_data(&anthropic, &usage, CostBasis::ProviderReported, None, true);
assert_eq!(data["gen_ai.provider.name"], "anthropic");
let mut responses = obs();
responses.wire_format = WireFormat::OpenAiResponses;
let data = assemble_data(&responses, &usage, CostBasis::ProviderReported, None, true);
assert_eq!(data["gen_ai.provider.name"], "openai");
}
#[test]
fn every_usage_row_names_its_provider() {
let formats = [
(
WireFormat::AnthropicMessages,
"claude-opus-4-8",
"anthropic",
),
(WireFormat::OpenAiResponses, "gpt-5", "openai"),
(
WireFormat::OpenAiChatCompletions,
"qwen2.5-coder:14b",
"openai-compatible",
),
(
WireFormat::GoogleGenerateContent,
"gemini-2.5-pro",
"gemini",
),
(WireFormat::OllamaNative, "llama3:8b", "ollama"),
(WireFormat::Unknown, "", "unknown"),
];
for (format, model, provider) in formats {
for model in [Some(model.to_string()), None] {
let mut o = obs();
o.wire_format = format;
o.model = model.clone();
let data = assemble_data(
&o,
&Usage::default(),
CostBasis::ProviderReported,
None,
true,
);
assert_eq!(
data["gen_ai.provider.name"], provider,
"{format:?} model={model:?}"
);
}
}
let mut gemini = obs();
gemini.wire_format = WireFormat::GoogleGenerateContent;
gemini.model = Some("models/gemini-2.5-pro".to_string());
let data = assemble_data(
&gemini,
&Usage::default(),
CostBasis::ProviderReported,
None,
true,
);
assert_eq!(data["gen_ai.provider.name"], "gemini");
assert_eq!(data["gen_ai.request.model"], "gemini-2.5-pro");
for (provider, model, want_model) in [
("gcp.gemini", "gemini-2.5-pro", "gemini-2.5-pro"),
("gcp.vertex_ai", "gemini-2.5-pro", "gemini-2.5-pro"),
("vertex_ai", "gemini-2.5-flash", "gemini-2.5-flash"),
(
"openai-compatible",
"gemini/gemini-2.5-pro",
"gemini-2.5-pro",
),
] {
assert_eq!(
priced_identity(provider, model),
("gemini".to_string(), want_model.to_string()),
"{provider} {model}"
);
}
assert_eq!(
priced_identity("openai-compatible", "mystery/model-x"),
(
"openai-compatible".to_string(),
"mystery/model-x".to_string()
)
);
}
#[test]
fn model_ids_are_normalized_to_the_vendors_own() {
let cases = [
(
"openai-compatible",
"anthropic/claude-sonnet-4.5",
"anthropic",
"claude-sonnet-4-5",
),
("openai-compatible", "openai/gpt-4.1", "openai", "gpt-4.1"),
(
"openai-compatible",
"google/gemini-2.5-pro",
"gemini",
"gemini-2.5-pro",
),
("openai-compatible", "x-ai/grok-4", "xai", "grok-4"),
("google", "gemini-2.5-pro", "gemini", "gemini-2.5-pro"),
(
"google",
"models/gemini-2.5-flash",
"gemini",
"gemini-2.5-flash",
),
(
"anthropic",
"claude-opus-4-8",
"anthropic",
"claude-opus-4-8",
),
(
"openai-compatible",
"qwen2.5-coder:14b",
"openai-compatible",
"qwen2.5-coder:14b",
),
(
"openai-compatible",
"meta-llama/llama-3.3-70b",
"openai-compatible",
"meta-llama/llama-3.3-70b",
),
("ollama", "llama3:8b", "ollama", "llama3:8b"),
];
for (format, model, provider, normalized) in cases {
assert_eq!(
priced_identity(format, model),
(provider.to_string(), normalized.to_string()),
"{format} {model}"
);
}
let mut routed = obs();
routed.wire_format = WireFormat::OpenAiChatCompletions;
routed.model = Some("anthropic/claude-sonnet-4.5".to_string());
let data = assemble_data(
&routed,
&Usage::default(),
CostBasis::ProviderReported,
None,
true,
);
assert_eq!(data["gen_ai.provider.name"], "anthropic");
assert_eq!(data["gen_ai.request.model"], "claude-sonnet-4-5");
}
#[test]
fn assembled_data_has_exact_contract_fields_and_no_usd() {
let usage = Usage {
input_tokens: 50,
cache_read: 100_000,
cache_write: 0,
eph_5m: 0,
eph_1h: 0,
output_tokens: 12,
};
let data = assemble_data(&obs(), &usage, CostBasis::ProviderReported, None, true);
assert_eq!(data["gen_ai.usage.input_tokens"], 50);
assert_eq!(data["gen_ai.usage.cache_read.input_tokens"], 100_000);
assert_eq!(data["gen_ai.provider.name"], "anthropic");
assert_eq!(data["ai.openlatch.cost.basis"], "provider_reported");
assert_eq!(data["ai.openlatch.billing.mode"], "api_key");
assert_eq!(data["ai.openlatch.session.assurance"], "attested");
assert_eq!(data["ai.openlatch.request.batch"], false);
assert!(data["ai.openlatch.request.inference_geo"].is_null());
assert!(data["ai.openlatch.capture.gap"].is_null());
assert!(data["ai.openlatch.prefix.finding_id"].is_null());
assert!(data["ai.openlatch.transform.rule_id"].is_null());
assert!(data["ai.openlatch.transform.outcome"].is_null());
assert!(data["ai.openlatch.transform.tokens_net"].is_null());
assert!(data.get("ai.openlatch.transform.finding_id").is_none());
let total = data["gen_ai.usage.input_tokens"].as_u64().unwrap()
+ data["gen_ai.usage.cache_creation.input_tokens"]
.as_u64()
.unwrap()
+ data["gen_ai.usage.cache_read.input_tokens"]
.as_u64()
.unwrap();
assert_eq!(total, 100_050);
let s = data.to_string();
assert!(!s.contains("cost_input"));
assert!(!s.contains("cost_total"));
assert!(!s.contains("pricebook"));
assert!(!s.contains("ai.openlatch.cost.input"));
}
#[test]
fn a_matching_transform_populates_the_nullable_tuple() {
use crate::model_relay::transforms::evaluate_would_have;
let mut messages = vec![
serde_json::json!({ "role": "user", "content": "a".repeat(400) }),
serde_json::json!({ "role": "assistant", "content": "a".repeat(400) }),
];
for _ in 0..6 {
messages.push(serde_json::json!({ "role": "user", "content": "hi" }));
}
let body = serde_json::json!({ "model": "claude-opus-4-8", "messages": messages });
let decision = evaluate_would_have(&body).expect("a matching L-1 rule");
let mut obs = obs();
obs.transform = Some(decision);
let usage = Usage {
input_tokens: 10,
..Usage::default()
};
let data = assemble_data(&obs, &usage, CostBasis::ProviderReported, None, true);
assert_eq!(data["ai.openlatch.transform.rule_id"], "OL-ECO-001");
assert_eq!(data["ai.openlatch.transform.lever"], "history_trim");
assert_eq!(data["ai.openlatch.transform.outcome"], "skipped_stage");
assert_eq!(data["ai.openlatch.transform.ladder_stage"], "observe");
assert_eq!(data["ai.openlatch.transform.rule_version"], 1);
assert_eq!(data["ai.openlatch.transform.bundle_revision"], 0);
assert_eq!(data["ai.openlatch.transform.write_multiplier"], 1.25);
assert!(
data["ai.openlatch.transform.tokens_gross"]
.as_u64()
.unwrap()
> 0
);
assert!(data["ai.openlatch.transform.tokens_net"].as_f64().unwrap() > 0.0);
assert_ne!(data["ai.openlatch.transform.outcome"], "applied");
assert!(data.get("ai.openlatch.transform.finding_id").is_none());
}
#[test]
fn no_op_observation_emits_nothing() {
let filter = PrivacyFilter::new(&[]);
let (tx, mut rx) = tokio::sync::mpsc::channel(1);
build_and_emit(
&Observation::none(),
&Usage::default(),
CostBasis::TokenizerEstimated,
None,
false,
&filter,
Some(&tx),
);
assert!(rx.try_recv().is_err(), "no event for an unmeasured request");
}
#[test]
fn unresolved_source_is_unknown() {
let mut obs = obs();
obs.session.source = None;
obs.attributable_agent = None;
let data = assemble_data(
&obs,
&Usage::default(),
CostBasis::ProviderReported,
None,
true,
);
assert_eq!(
data["ai.openlatch.session.source"], "unknown",
"an unattributed request must say so, not name an agent"
);
let event = assemble_event(&obs, data, &PrivacyFilter::new(&[]));
assert_eq!(
event.envelope["source"], "unknown",
"the CloudEvents `source` carries the same claim as the data field"
);
}
#[test]
fn the_sole_wired_speaker_names_an_otherwise_unattributed_request() {
let mut obs = obs();
obs.session.source = None;
obs.wire_format = WireFormat::AnthropicMessages;
obs.attributable_agent = Some("claude-code");
let data = assemble_data(
&obs,
&Usage::default(),
CostBasis::ProviderReported,
None,
true,
);
assert_eq!(
data["ai.openlatch.session.source"], "claude-code",
"the only agent wired to this protocol on this install is the author"
);
let event = assemble_event(&obs, data, &PrivacyFilter::new(&[]));
assert_eq!(
event.envelope["source"], "claude-code",
"the CloudEvents `source` carries the same claim as the data field"
);
}
#[test]
fn an_observed_session_outranks_the_wiring_deduction() {
let mut obs = obs();
obs.session.source = Some("codex-cli".to_string());
obs.attributable_agent = Some("claude-code");
let data = assemble_data(
&obs,
&Usage::default(),
CostBasis::ProviderReported,
None,
true,
);
assert_eq!(
data["ai.openlatch.session.source"], "codex-cli",
"a hook that saw the session outranks an inference from our own wiring"
);
}
}