1use std::collections::HashMap;
2
3use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5use uuid::Uuid;
6
7use crate::error::{Error, Result};
8use crate::event::{Event, EventId};
9use crate::event_key::EventKey;
10use crate::metadata::Metadata;
11use crate::namespace::Namespace;
12use crate::organization::OrganizationId;
13use crate::payload::{ContentType, Payload};
14use crate::topic::Topic;
15
16#[derive(Debug, Clone, Serialize, Deserialize)]
22pub struct SerializedEvent {
23 pub id: Uuid,
24 pub organization: String,
25 pub namespace: String,
26 pub topic: String,
27 pub key: String,
28 pub payload: SerializedPayload,
29 pub metadata: HashMap<String, String>,
30 pub timestamp: DateTime<Utc>,
31 pub version: u64,
32 #[serde(default)]
33 pub parent_id: Option<Uuid>,
34 #[serde(default)]
35 pub correlation_id: Option<String>,
36 #[serde(default)]
37 pub causation_id: Option<String>,
38}
39
40#[derive(Debug, Clone, Serialize, Deserialize)]
52#[serde(tag = "content_type", content = "data")]
53pub enum SerializedPayload {
54 #[serde(rename = "application/json")]
55 Json(serde_json::Value),
56 #[serde(rename = "text/plain")]
57 PlainText(String),
58 #[serde(rename = "application/octet-stream")]
59 Binary(#[serde(with = "base64_bytes")] Vec<u8>),
60}
61
62impl SerializedPayload {
63 pub fn content_type(&self) -> ContentType {
64 match self {
65 Self::Json(_) => ContentType::Json,
66 Self::PlainText(_) => ContentType::PlainText,
67 Self::Binary(_) => ContentType::Binary,
68 }
69 }
70
71 pub fn from_payload(payload: &Payload) -> Result<Self> {
72 match payload.content_type() {
73 ContentType::Json => {
74 let value = serde_json::from_slice(payload.data())
75 .map_err(|e| Error::Serialization(e.to_string()))?;
76 Ok(Self::Json(value))
77 }
78 ContentType::PlainText => {
79 let text = std::str::from_utf8(payload.data())
80 .map_err(|e| Error::Serialization(e.to_string()))?;
81 Ok(Self::PlainText(text.to_owned()))
82 }
83 ContentType::Binary => Ok(Self::Binary(payload.data().to_vec())),
84 }
85 }
86
87 pub fn into_payload(self) -> Result<Payload> {
88 match self {
89 Self::Json(value) => {
90 let bytes =
91 serde_json::to_vec(&value).map_err(|e| Error::Serialization(e.to_string()))?;
92 Ok(Payload::from_raw(bytes, ContentType::Json))
93 }
94 Self::PlainText(text) => Ok(Payload::from_string(text)),
95 Self::Binary(bytes) => Ok(Payload::from_bytes(bytes)),
96 }
97 }
98}
99
100mod base64_bytes {
101 use base64::Engine;
102 use base64::engine::general_purpose::STANDARD as BASE64;
103 use serde::{Deserialize, Deserializer, Serializer};
104
105 pub(super) fn serialize<S: Serializer>(bytes: &[u8], s: S) -> Result<S::Ok, S::Error> {
106 s.serialize_str(&BASE64.encode(bytes))
107 }
108
109 pub(super) fn deserialize<'de, D: Deserializer<'de>>(d: D) -> Result<Vec<u8>, D::Error> {
110 let encoded = String::deserialize(d)?;
111 BASE64.decode(encoded).map_err(serde::de::Error::custom)
112 }
113}
114
115impl SerializedEvent {
116 pub fn from_event(event: &Event<Payload>) -> Result<Self> {
117 Ok(Self {
118 id: *event.id().as_uuid(),
119 organization: event.organization().to_string(),
120 namespace: event.namespace().to_string(),
121 topic: event.topic().to_string(),
122 key: event.key().to_string(),
123 payload: SerializedPayload::from_payload(event.payload())?,
124 metadata: event
125 .metadata()
126 .as_map()
127 .iter()
128 .map(|(k, v)| (k.clone(), v.clone()))
129 .collect(),
130 timestamp: event.timestamp(),
131 version: event.version(),
132 parent_id: event.parent_id().map(|id| *id.as_uuid()),
133 correlation_id: event.correlation_id().map(|id| id.to_string()),
134 causation_id: event.causation_id().map(|id| id.to_string()),
135 })
136 }
137
138 pub fn to_event(&self) -> Result<Event<Payload>> {
139 let key = EventKey::new(&self.key)?;
140 let payload = self.payload.clone().into_payload()?;
141 let parent_id = self.parent_id.map(EventId::from_uuid);
142 let correlation_id = self
143 .correlation_id
144 .as_deref()
145 .map(EventKey::new)
146 .transpose()?;
147 let causation_id = self
148 .causation_id
149 .as_deref()
150 .map(EventKey::new)
151 .transpose()?;
152
153 Event::new(
154 EventId::from_uuid(self.id),
155 OrganizationId::new(&self.organization)?,
156 Namespace::new(&self.namespace)?,
157 Topic::new(&self.topic)?,
158 key,
159 payload,
160 Metadata::try_from(self.metadata.clone())?,
161 self.timestamp,
162 self.version,
163 parent_id,
164 correlation_id,
165 causation_id,
166 )
167 }
168
169 pub fn to_json_value(&self) -> serde_json::Value {
170 serde_json::to_value(self).expect("SerializedEvent must serialize to JSON")
171 }
172
173 pub fn from_json_value(value: serde_json::Value) -> Result<Self> {
174 serde_json::from_value(value).map_err(|e| Error::Serialization(e.to_string()))
175 }
176
177 pub fn to_json_string(&self) -> Result<String> {
178 serde_json::to_string(self).map_err(|e| Error::Serialization(e.to_string()))
179 }
180
181 pub fn from_json_str(s: &str) -> Result<Self> {
182 serde_json::from_str(s).map_err(|e| Error::Serialization(e.to_string()))
183 }
184
185 pub fn from_json_slice(bytes: &[u8]) -> Result<Self> {
186 let value =
187 serde_json::from_slice(bytes).map_err(|e| Error::Serialization(e.to_string()))?;
188 Self::from_json_value(value)
189 }
190}
191
192#[cfg(test)]
193mod tests {
194 use super::*;
195
196 fn event_with_key(payload: Payload, key: &str) -> Event {
197 Event::builder("acme", "/x", "thing.happened", key, payload)
198 .unwrap()
199 .build()
200 .expect("valid event")
201 }
202
203 #[test]
204 fn roundtrip() {
205 let payload = Payload::from_json(&serde_json::json!({"key": "value"})).unwrap();
206 let event = Event::builder("acme", "/task", "task.created", "task-123", payload)
207 .unwrap()
208 .build()
209 .unwrap();
210
211 let serialized = SerializedEvent::from_event(&event).unwrap();
212 assert_eq!(serialized.topic, "task.created");
213 assert_eq!(serialized.namespace, "/task");
214 assert_eq!(serialized.organization, "acme");
215 assert_eq!(serialized.key.as_str(), "task-123");
216
217 let restored = serialized.to_event().unwrap();
218 assert_eq!(restored.topic().as_str(), "task.created");
219 assert_eq!(restored.id(), event.id());
220 assert_eq!(restored.key().as_str(), "task-123");
221 }
222
223 #[test]
224 fn field_order_matches_event() {
225 let event = Event::builder(
226 "acme",
227 "/x",
228 "thing.happened",
229 "k",
230 Payload::from_string("p"),
231 )
232 .unwrap()
233 .build()
234 .unwrap();
235 let serialized = SerializedEvent::from_event(&event).unwrap();
236 let json = serialized.to_json_string().unwrap();
237 let id_pos = json.find("\"id\"").unwrap();
238 let org_pos = json.find("\"organization\"").unwrap();
239 let ns_pos = json.find("\"namespace\"").unwrap();
240 let topic_pos = json.find("\"topic\"").unwrap();
241 let key_pos = json.find("\"key\"").unwrap();
242 let payload_pos = json.find("\"payload\"").unwrap();
243 let metadata_pos = json.find("\"metadata\"").unwrap();
244 let timestamp_pos = json.find("\"timestamp\"").unwrap();
245 let version_pos = json.find("\"version\"").unwrap();
246 assert!(id_pos < org_pos);
247 assert!(org_pos < ns_pos);
248 assert!(ns_pos < topic_pos);
249 assert!(topic_pos < key_pos);
250 assert!(key_pos < payload_pos);
251 assert!(payload_pos < metadata_pos);
252 assert!(metadata_pos < timestamp_pos);
253 assert!(timestamp_pos < version_pos);
254 }
255
256 #[test]
257 fn plain_text_payload_is_a_raw_json_string() {
258 let event = event_with_key(Payload::from_string("hello world"), "k");
259 let serialized = SerializedEvent::from_event(&event).unwrap();
260 match &serialized.payload {
261 SerializedPayload::PlainText(text) => assert_eq!(text, "hello world"),
262 other => panic!("expected PlainText, got {other:?}"),
263 }
264 let json = serialized.to_json_string().unwrap();
265 assert!(json.contains("\"content_type\":\"text/plain\""));
266 assert!(json.contains("\"data\":\"hello world\""));
267
268 let restored = serialized.to_event().unwrap();
269 assert_eq!(restored.payload().data(), b"hello world");
270 assert_eq!(restored.payload().content_type(), ContentType::PlainText);
271 }
272
273 #[test]
274 fn binary_payload_is_base64_in_data_field() {
275 let bytes = vec![0xff, 0x00, 0x01, 0xfe, 0x80, 0x7f, 0x10];
276 let event = event_with_key(Payload::from_bytes(bytes.clone()), "k");
277 let serialized = SerializedEvent::from_event(&event).unwrap();
278 match &serialized.payload {
279 SerializedPayload::Binary(b) => assert_eq!(b, &bytes),
280 other => panic!("expected Binary, got {other:?}"),
281 }
282 let json = serialized.to_json_string().unwrap();
283 assert!(json.contains("\"content_type\":\"application/octet-stream\""));
284
285 let restored = serialized.to_event().unwrap();
286 assert_eq!(restored.payload().data(), bytes.as_slice());
287 assert_eq!(restored.payload().content_type(), ContentType::Binary);
288 }
289
290 #[test]
291 fn json_payload_stays_human_readable() {
292 let event = event_with_key(
293 Payload::from_json(&serde_json::json!({"k": "v"})).unwrap(),
294 "k",
295 );
296 let serialized = SerializedEvent::from_event(&event).unwrap();
297 match &serialized.payload {
298 SerializedPayload::Json(value) => {
299 assert_eq!(value, &serde_json::json!({"k": "v"}));
300 }
301 other => panic!("expected Json, got {other:?}"),
302 }
303 let json = serialized.to_json_string().unwrap();
304 assert!(json.contains("\"content_type\":\"application/json\""));
305 assert!(json.contains("\"data\":{\"k\":\"v\"}"));
306 }
307
308 #[test]
309 fn json_value_round_trip() {
310 let event = Event::builder(
311 "acme",
312 "/task",
313 "task.created",
314 "task-123",
315 Payload::from_json(&serde_json::json!({"key": "value"})).unwrap(),
316 )
317 .unwrap()
318 .build()
319 .unwrap();
320
321 let serialized = SerializedEvent::from_event(&event).unwrap();
322 let value = serialized.to_json_value();
323 let parsed = SerializedEvent::from_json_value(value).unwrap();
324
325 assert_eq!(parsed.id, serialized.id);
326 assert_eq!(parsed.topic, serialized.topic);
327 assert_eq!(parsed.payload.content_type(), ContentType::Json);
328 }
329
330 #[test]
331 fn json_string_round_trip() {
332 let event = Event::builder(
333 "acme",
334 "/task",
335 "task.created",
336 "task-123",
337 Payload::from_json(&serde_json::json!({"key": "value"})).unwrap(),
338 )
339 .unwrap()
340 .build()
341 .unwrap();
342
343 let serialized = SerializedEvent::from_event(&event).unwrap();
344 let s = serialized.to_json_string().unwrap();
345 let parsed = SerializedEvent::from_json_str(&s).unwrap();
346
347 assert_eq!(parsed.id, serialized.id);
348 assert_eq!(parsed.namespace, serialized.namespace);
349 assert_eq!(parsed.organization, serialized.organization);
350 }
351
352 #[test]
353 fn from_json_slice_roundtrip() {
354 let event = Event::builder(
355 "acme",
356 "/task",
357 "task.created",
358 "task-123",
359 Payload::from_json(&serde_json::json!({"key": "value"})).unwrap(),
360 )
361 .unwrap()
362 .build()
363 .unwrap();
364
365 let serialized = SerializedEvent::from_event(&event).unwrap();
366 let bytes = serialized.to_json_string().unwrap().into_bytes();
367 let parsed = SerializedEvent::from_json_slice(&bytes).unwrap();
368
369 assert_eq!(parsed.id, serialized.id);
370 assert_eq!(parsed.topic, serialized.topic);
371 assert_eq!(parsed.payload.content_type(), ContentType::Json);
372 }
373
374 #[test]
375 fn json_string_round_trip_with_binary() {
376 let bytes = vec![0xde, 0xad, 0xbe, 0xef];
377 let event = event_with_key(Payload::from_bytes(bytes.clone()), "b1");
378
379 let serialized = SerializedEvent::from_event(&event).unwrap();
380 let s = serialized.to_json_string().unwrap();
381 let parsed = SerializedEvent::from_json_str(&s).unwrap();
382 let restored = parsed.to_event().unwrap();
383
384 assert_eq!(restored.payload().data(), bytes.as_slice());
385 }
386
387 #[test]
388 fn lineage_fields_roundtrip() {
389 let parent_id = EventId::new();
390 let event = Event::builder(
391 "acme",
392 "/x",
393 "thing.happened",
394 "k",
395 Payload::from_string("p"),
396 )
397 .unwrap()
398 .parent_id(parent_id)
399 .correlation_id("corr")
400 .unwrap()
401 .causation_id("cause")
402 .unwrap()
403 .build()
404 .unwrap();
405
406 let serialized = SerializedEvent::from_event(&event).unwrap();
407 assert_eq!(serialized.key.as_str(), "k");
408 assert_eq!(serialized.parent_id, Some(*parent_id.as_uuid()));
409 assert_eq!(serialized.correlation_id.as_deref(), Some("corr"));
410 assert_eq!(serialized.causation_id.as_deref(), Some("cause"));
411
412 let restored = serialized.to_event().unwrap();
413 assert_eq!(restored.key().as_str(), "k");
414 assert_eq!(restored.parent_id(), Some(parent_id));
415 assert_eq!(
416 restored.correlation_id().map(EventKey::as_str),
417 Some("corr")
418 );
419 assert_eq!(restored.causation_id().map(EventKey::as_str), Some("cause"));
420 }
421
422 #[test]
423 fn serialized_event_rejects_missing_key() {
424 let value = serde_json::json!({
425 "id": uuid::Uuid::now_v7(),
426 "organization": "acme",
427 "namespace": "/task",
428 "topic": "task.created",
429 "payload": {"content_type": "text/plain", "data": "hello"},
430 "metadata": {},
431 "timestamp": chrono::Utc::now(),
432 "version": 1
433 });
434
435 let err = SerializedEvent::from_json_value(value).unwrap_err();
436 assert!(matches!(err, Error::Serialization(_)));
437 }
438}