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    /// Conversation / agent thread (`x-synapse-thread`).
26    #[serde(skip_serializing_if = "Option::is_none")]
27    pub thread_id: Option<String>,
28    /// Chat / work message within the thread (`x-synapse-message`).
29    #[serde(skip_serializing_if = "Option::is_none")]
30    pub message_id: Option<String>,
31    pub request_id: String,
32    pub timestamp: DateTime<Utc>,
33    #[serde(rename = "type")]
34    pub event_type: &'static str,
35    pub route: String,
36    pub provider: String,
37    pub model: String,
38    pub lane: String,
39    pub input_tokens: u64,
40    pub output_tokens: u64,
41    pub cost_usd: f64,
42    pub status: String,
43    /// Lane discriminator: "chat" or "embedding".
44    pub op: String,
45}
46
47impl From<&UsageEntry> for UsageEvent {
48    fn from(e: &UsageEntry) -> Self {
49        Self {
50            namespace: e.tenant.clone(),
51            workspace: e.workspace.clone(),
52            user: e.user.clone(),
53            thread_id: e.thread.clone(),
54            message_id: e.message.clone(),
55            request_id: e.request_id.clone(),
56            timestamp: e.ts,
57            event_type: "usage",
58            route: e.route.clone(),
59            provider: e.provider.clone(),
60            model: e.model.clone(),
61            lane: e.lane.clone(),
62            input_tokens: e.input_tokens,
63            output_tokens: e.output_tokens,
64            cost_usd: e.cost_usd,
65            status: e.status.clone(),
66            op: e.op.clone(),
67        }
68    }
69}
70
71impl UsageEvent {
72    /// Broker message attributes for subscription filtering (talos-aligned keys
73    /// plus synapse-useful `provider`/`status`). Returned as ordered pairs so
74    /// both backends can convert to their SDK-specific attribute type.
75    pub fn attributes(&self) -> Vec<(&'static str, String)> {
76        vec![
77            ("EventType", LEDGER_EVENT_TYPE.to_string()),
78            ("namespace", self.namespace.clone()),
79            ("requestId", self.request_id.clone()),
80            ("type", self.event_type.to_string()),
81            ("provider", self.provider.clone()),
82            ("status", self.status.clone()),
83        ]
84    }
85}
86
87#[cfg(test)]
88mod tests {
89    use super::*;
90    use chrono::TimeZone;
91
92    fn entry() -> UsageEntry {
93        UsageEntry {
94            ts: Utc.with_ymd_and_hms(2026, 6, 10, 15, 30, 45).unwrap(),
95            tenant: "acme".into(),
96            workspace: None,
97            user: None,
98            thread: None,
99            message: None,
100            route: "gemini-pro".into(),
101            provider: "vertex".into(),
102            model: "gemini-3-pro".into(),
103            lane: "standard".into(),
104            input_tokens: 3,
105            output_tokens: 5,
106            cost_usd: 0.001,
107            request_id: "req-1".into(),
108            status: "ok".into(),
109            op: "chat".into(),
110        }
111    }
112
113    #[test]
114    fn serializes_talos_aligned_camelcase() {
115        let v = serde_json::to_value(UsageEvent::from(&entry())).unwrap();
116        assert_eq!(v["namespace"], "acme");
117        assert_eq!(v["type"], "usage");
118        assert_eq!(v["requestId"], "req-1");
119        assert_eq!(v["inputTokens"], 3);
120        assert_eq!(v["outputTokens"], 5);
121        assert_eq!(v["costUsd"], 0.001);
122        assert_eq!(v["lane"], "standard");
123        assert_eq!(v["op"], "chat");
124        assert!(v.get("workspace").is_none());
125        assert!(v.get("user").is_none());
126        assert!(v.get("request_id").is_none());
127        assert!(v.get("input_tokens").is_none());
128    }
129
130    #[test]
131    fn serializes_thread_and_message_when_present() {
132        let mut e = entry();
133        e.thread = Some("thread-9".into());
134        e.message = Some("msg-7".into());
135        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
136        assert_eq!(v["threadId"], "thread-9");
137        assert_eq!(v["messageId"], "msg-7");
138    }
139
140    #[test]
141    fn serializes_user_when_present() {
142        let mut e = entry();
143        e.user = Some("user-42".into());
144        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
145        assert_eq!(v["user"], "user-42");
146    }
147
148    #[test]
149    fn serializes_embedding_op() {
150        let mut e = entry();
151        e.op = "embedding".into();
152        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
153        assert_eq!(v["op"], "embedding");
154    }
155
156    #[test]
157    fn attributes_have_talos_keys() {
158        let attrs = UsageEvent::from(&entry()).attributes();
159        let keys: Vec<&str> = attrs.iter().map(|(k, _)| *k).collect();
160        assert_eq!(
161            keys,
162            vec![
163                "EventType",
164                "namespace",
165                "requestId",
166                "type",
167                "provider",
168                "status"
169            ]
170        );
171        assert_eq!(attrs[0].1, LEDGER_EVENT_TYPE);
172        assert_eq!(attrs[1].1, "acme");
173    }
174}