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    /// Caller-supplied classification of the work (`x-synapse-user-task-type`).
46    #[serde(skip_serializing_if = "Option::is_none")]
47    pub user_task_type: Option<String>,
48    /// Resolved classification of the work the gateway performed. Always set
49    /// (defaults to `"simple"`), so unlike `user_task_type` it is never skipped.
50    pub ai_task_type: String,
51}
52
53impl From<&UsageEntry> for UsageEvent {
54    fn from(e: &UsageEntry) -> Self {
55        Self {
56            namespace: e.tenant.clone(),
57            workspace: e.workspace.clone(),
58            user: e.user.clone(),
59            thread_id: e.thread.clone(),
60            message_id: e.message.clone(),
61            request_id: e.request_id.clone(),
62            timestamp: e.ts,
63            event_type: "usage",
64            route: e.route.clone(),
65            provider: e.provider.clone(),
66            model: e.model.clone(),
67            lane: e.lane.clone(),
68            input_tokens: e.input_tokens,
69            output_tokens: e.output_tokens,
70            cost_usd: e.cost_usd,
71            status: e.status.clone(),
72            op: e.op.clone(),
73            user_task_type: e.user_task_type.clone(),
74            ai_task_type: e.ai_task_type.clone(),
75        }
76    }
77}
78
79impl UsageEvent {
80    /// Broker message attributes for subscription filtering (talos-aligned keys
81    /// plus synapse-useful `provider`/`status`). Returned as ordered pairs so
82    /// both backends can convert to their SDK-specific attribute type.
83    pub fn attributes(&self) -> Vec<(&'static str, String)> {
84        vec![
85            ("EventType", LEDGER_EVENT_TYPE.to_string()),
86            ("namespace", self.namespace.clone()),
87            ("requestId", self.request_id.clone()),
88            ("type", self.event_type.to_string()),
89            ("provider", self.provider.clone()),
90            ("status", self.status.clone()),
91        ]
92    }
93}
94
95#[cfg(test)]
96mod tests {
97    use super::*;
98    use chrono::TimeZone;
99
100    fn entry() -> UsageEntry {
101        UsageEntry {
102            ts: Utc.with_ymd_and_hms(2026, 6, 10, 15, 30, 45).unwrap(),
103            tenant: "acme".into(),
104            workspace: None,
105            user: None,
106            thread: None,
107            message: None,
108            route: "gemini-pro".into(),
109            provider: "vertex".into(),
110            model: "gemini-3-pro".into(),
111            lane: "standard".into(),
112            input_tokens: 3,
113            output_tokens: 5,
114            cost_usd: 0.001,
115            request_id: "req-1".into(),
116            status: "ok".into(),
117            op: "chat".into(),
118            user_task_type: None,
119            ai_task_type: "simple".into(),
120        }
121    }
122
123    #[test]
124    fn serializes_talos_aligned_camelcase() {
125        let v = serde_json::to_value(UsageEvent::from(&entry())).unwrap();
126        assert_eq!(v["namespace"], "acme");
127        assert_eq!(v["type"], "usage");
128        assert_eq!(v["requestId"], "req-1");
129        assert_eq!(v["inputTokens"], 3);
130        assert_eq!(v["outputTokens"], 5);
131        assert_eq!(v["costUsd"], 0.001);
132        assert_eq!(v["lane"], "standard");
133        assert_eq!(v["op"], "chat");
134        assert!(v.get("workspace").is_none());
135        assert!(v.get("user").is_none());
136        assert!(v.get("request_id").is_none());
137        assert!(v.get("input_tokens").is_none());
138        assert!(v.get("userTaskType").is_none());
139        // Always resolved, so unlike userTaskType it is always serialized.
140        assert_eq!(v["aiTaskType"], "simple");
141        assert!(v.get("ai_task_type").is_none());
142    }
143
144    #[test]
145    fn serializes_the_resolved_ai_task_type() {
146        let mut e = entry();
147        e.ai_task_type = "conversation".into();
148        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
149        assert_eq!(v["aiTaskType"], "conversation");
150    }
151
152    #[test]
153    fn serializes_user_task_type_when_present() {
154        let mut e = entry();
155        e.user_task_type = Some("summarisation".into());
156        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
157        assert_eq!(v["userTaskType"], "summarisation");
158        assert!(v.get("user_task_type").is_none());
159    }
160
161    #[test]
162    fn serializes_thread_and_message_when_present() {
163        let mut e = entry();
164        e.thread = Some("thread-9".into());
165        e.message = Some("msg-7".into());
166        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
167        assert_eq!(v["threadId"], "thread-9");
168        assert_eq!(v["messageId"], "msg-7");
169    }
170
171    #[test]
172    fn serializes_user_when_present() {
173        let mut e = entry();
174        e.user = Some("user-42".into());
175        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
176        assert_eq!(v["user"], "user-42");
177    }
178
179    #[test]
180    fn serializes_embedding_op() {
181        let mut e = entry();
182        e.op = "embedding".into();
183        let v = serde_json::to_value(UsageEvent::from(&e)).unwrap();
184        assert_eq!(v["op"], "embedding");
185    }
186
187    #[test]
188    fn attributes_have_talos_keys() {
189        let attrs = UsageEvent::from(&entry()).attributes();
190        let keys: Vec<&str> = attrs.iter().map(|(k, _)| *k).collect();
191        assert_eq!(
192            keys,
193            vec![
194                "EventType",
195                "namespace",
196                "requestId",
197                "type",
198                "provider",
199                "status"
200            ]
201        );
202        assert_eq!(attrs[0].1, LEDGER_EVENT_TYPE);
203        assert_eq!(attrs[1].1, "acme");
204    }
205}