Skip to main content

eventuary_core/
serialization.rs

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/// Wire-format representation of an [`Event`]. Field order matches
17/// `Event` so the JSON shape is predictable and self-documenting.
18///
19/// `id` and `parent_id` carry `Uuid` directly so backends with native
20/// UUID columns can bind/fetch without a String round-trip.
21#[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/// Wire-format representation of a [`Payload`].
41///
42/// Each variant carries the payload bytes in the natural shape for its
43/// content type:
44/// - `Json` carries the parsed `serde_json::Value` so the wire format
45///   stays human-readable and `jq`-friendly.
46/// - `PlainText` carries the text as a raw JSON string.
47/// - `Binary` carries the raw bytes base64-encoded.
48///
49/// The wire tag is the content-type string, so the JSON shape is
50/// `{"content_type": "...", "data": ...}` regardless of variant.
51#[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}