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