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};
pub const ECON_TYPE: &str = "ai.openlatch.economics.usage";
pub const PROVIDER: &str = "anthropic";
pub const DEFAULT_SOURCE: &str = "claude-code";
#[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>,
}
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,
}
}
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()
.unwrap_or_else(|| DEFAULT_SOURCE.to_string())
}
}
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 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": obs.model.clone().unwrap_or_default(),
"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,
"ai.openlatch.session.assurance": 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 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,
});
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,
"boundary: cloud channel full — economics event dropped"
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::boundary::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,
}
}
#[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::boundary::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");
}
}