Skip to main content

synapse/ledger/
event.rs

1//! talos-aligned published usage event (camelCase JSON). Decouples the wire
2//! format from the DB-shaped `UsageEntry`. Mirrors the conventions of
3//! `talos-core/src/tool_events.rs` (`ToolEvent`): `rename_all = "camelCase"`,
4//! tenancy as `namespace`, a `type` discriminator, RFC3339 timestamp.
5
6use chrono::{DateTime, Utc};
7use serde::Serialize;
8
9use crate::ledger::UsageEntry;
10
11/// One terminal usage event, published to a topic per request.
12#[derive(Debug, Clone, Serialize)]
13#[serde(rename_all = "camelCase")]
14pub struct UsageEvent {
15    /// Tenancy. talos calls this `namespace`; synapse maps its `tenant` onto it.
16    pub namespace: String,
17    #[serde(skip_serializing_if = "Option::is_none")]
18    pub workspace: Option<String>,
19    pub request_id: String,
20    pub timestamp: DateTime<Utc>,
21    #[serde(rename = "type")]
22    pub event_type: &'static str,
23    pub route: String,
24    pub provider: String,
25    pub model: String,
26    pub lane: String,
27    pub input_tokens: u64,
28    pub output_tokens: u64,
29    pub cost_usd: f64,
30    pub status: String,
31    /// Lane discriminator: "chat" or "embedding".
32    pub op: String,
33}
34
35impl From<&UsageEntry> for UsageEvent {
36    fn from(e: &UsageEntry) -> Self {
37        Self {
38            namespace: e.tenant.clone(),
39            workspace: e.workspace.clone(),
40            request_id: e.request_id.clone(),
41            timestamp: e.ts,
42            event_type: "usage",
43            route: e.route.clone(),
44            provider: e.provider.clone(),
45            model: e.model.clone(),
46            lane: e.lane.clone(),
47            input_tokens: e.input_tokens,
48            output_tokens: e.output_tokens,
49            cost_usd: e.cost_usd,
50            status: e.status.clone(),
51            op: e.op.clone(),
52        }
53    }
54}
55
56impl UsageEvent {
57    /// Broker message attributes for subscription filtering (talos-aligned keys
58    /// plus synapse-useful `provider`/`status`). Returned as ordered pairs so
59    /// both backends can convert to their SDK-specific attribute type.
60    pub fn attributes(&self) -> Vec<(&'static str, String)> {
61        vec![
62            ("namespace", self.namespace.clone()),
63            ("requestId", self.request_id.clone()),
64            ("type", self.event_type.to_string()),
65            ("provider", self.provider.clone()),
66            ("status", self.status.clone()),
67        ]
68    }
69}
70
71#[cfg(test)]
72mod tests {
73    use super::*;
74    use chrono::TimeZone;
75
76    fn entry() -> UsageEntry {
77        UsageEntry {
78            ts: Utc.with_ymd_and_hms(2026, 6, 10, 15, 30, 45).unwrap(),
79            tenant: "acme".into(),
80            workspace: None,
81            route: "gemini-pro".into(),
82            provider: "vertex".into(),
83            model: "gemini-3-pro".into(),
84            lane: "standard".into(),
85            input_tokens: 3,
86            output_tokens: 5,
87            cost_usd: 0.001,
88            request_id: "req-1".into(),
89            status: "ok".into(),
90            op: "chat".into(),
91        }
92    }
93
94    #[test]
95    fn serializes_talos_aligned_camelcase() {
96        let v = serde_json::to_value(UsageEvent::from(&entry())).unwrap();
97        assert_eq!(v["namespace"], "acme");
98        assert_eq!(v["type"], "usage");
99        assert_eq!(v["requestId"], "req-1");
100        assert_eq!(v["inputTokens"], 3);
101        assert_eq!(v["outputTokens"], 5);
102        assert_eq!(v["costUsd"], 0.001);
103        assert_eq!(v["lane"], "standard");
104        assert_eq!(v["op"], "chat");
105        assert!(v.get("workspace").is_none());
106        assert!(v.get("request_id").is_none());
107        assert!(v.get("input_tokens").is_none());
108    }
109
110    #[test]
111    fn serializes_embedding_op() {
112        let mut e = entry();
113        e.op = "embedding".into();
114        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
115        assert_eq!(v["op"], "embedding");
116    }
117
118    #[test]
119    fn attributes_have_talos_keys() {
120        let attrs = UsageEvent::from(&entry()).attributes();
121        let keys: Vec<&str> = attrs.iter().map(|(k, _)| *k).collect();
122        assert_eq!(
123            keys,
124            vec!["namespace", "requestId", "type", "provider", "status"]
125        );
126        assert_eq!(attrs[0].1, "acme");
127    }
128}