eventuary_core/
payload_codec.rs1use serde::Serialize;
2use serde::de::DeserializeOwned;
3
4use crate::error::Result;
5use crate::event::Event;
6use crate::payload::Payload;
7
8pub trait PayloadCodec<P>: Clone + Send + Sync + 'static {
9 fn encode(&self, payload: &P) -> Result<Payload>;
10 fn decode(&self, payload: &Payload) -> Result<P>;
11}
12
13pub trait EventCodec<P>: Clone + Send + Sync + 'static {
14 fn encode(&self, event: &Event<P>) -> Result<Event<Payload>>;
15 fn decode(&self, event: Event<Payload>) -> Result<Event<P>>;
16}
17
18#[derive(Debug, Clone)]
19pub struct PayloadEventCodec<C> {
20 payload_codec: C,
21}
22
23impl<C> PayloadEventCodec<C> {
24 pub fn new(payload_codec: C) -> Self {
25 Self { payload_codec }
26 }
27
28 pub fn payload_codec(&self) -> &C {
29 &self.payload_codec
30 }
31}
32
33impl<P, C> EventCodec<P> for PayloadEventCodec<C>
34where
35 C: PayloadCodec<P>,
36{
37 fn encode(&self, event: &Event<P>) -> Result<Event<Payload>> {
38 event.encode_payload(&self.payload_codec)
39 }
40
41 fn decode(&self, event: Event<Payload>) -> Result<Event<P>> {
42 event.decode_payload(&self.payload_codec)
43 }
44}
45
46#[derive(Debug, Clone, Copy, Default)]
47pub struct JsonPayloadCodec;
48
49impl<P> PayloadCodec<P> for JsonPayloadCodec
50where
51 P: Serialize + DeserializeOwned,
52{
53 fn encode(&self, payload: &P) -> Result<Payload> {
54 Payload::from_json(payload)
55 }
56
57 fn decode(&self, payload: &Payload) -> Result<P> {
58 payload.to_json()
59 }
60}
61
62#[derive(Debug, Clone, Copy, Default)]
63pub struct PayloadPassthroughCodec;
64
65impl PayloadCodec<Payload> for PayloadPassthroughCodec {
66 fn encode(&self, payload: &Payload) -> Result<Payload> {
67 Ok(payload.clone())
68 }
69
70 fn decode(&self, payload: &Payload) -> Result<Payload> {
71 Ok(payload.clone())
72 }
73}
74
75#[cfg(test)]
76mod tests {
77 use super::*;
78
79 use crate::error::Error;
80
81 #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
82 struct UserUpdated {
83 user_id: String,
84 email: String,
85 }
86
87 #[test]
88 fn json_codec_roundtrips_typed_payload() {
89 let codec = JsonPayloadCodec;
90 let typed = UserUpdated {
91 user_id: "u-1".to_owned(),
92 email: "a@example.com".to_owned(),
93 };
94
95 let payload = codec.encode(&typed).unwrap();
96 assert_eq!(payload.content_type(), crate::ContentType::Json);
97
98 let decoded: UserUpdated = codec.decode(&payload).unwrap();
99 assert_eq!(decoded, typed);
100 }
101
102 #[test]
103 fn json_codec_rejects_non_json_payload() {
104 let codec = JsonPayloadCodec;
105 let payload = Payload::from_string("not-json");
106 let decoded: Result<serde_json::Value> = codec.decode(&payload);
107 let err = decoded.unwrap_err();
108 assert!(matches!(err, Error::InvalidPayload(_)));
109 }
110
111 #[test]
112 fn payload_passthrough_clones_payload() {
113 let codec = PayloadPassthroughCodec;
114 let payload = Payload::from_string("hello");
115 let encoded = codec.encode(&payload).unwrap();
116 let decoded = codec.decode(&encoded).unwrap();
117 assert_eq!(decoded.data(), b"hello");
118 assert_eq!(decoded.content_type(), crate::ContentType::PlainText);
119 }
120
121 #[test]
122 fn payload_event_codec_encodes_and_decodes_event() {
123 let codec = PayloadEventCodec::new(JsonPayloadCodec);
124 let event: Event<UserUpdated> = Event::create(
125 "acme",
126 "/users",
127 "user.updated",
128 "thing-1",
129 UserUpdated {
130 user_id: "u-1".to_owned(),
131 email: "a@example.com".to_owned(),
132 },
133 )
134 .unwrap();
135
136 let id = event.id();
137 let encoded = codec.encode(&event).unwrap();
138 assert_eq!(encoded.id(), id);
139 assert_eq!(encoded.payload().content_type(), crate::ContentType::Json);
140
141 let decoded: Event<UserUpdated> = codec.decode(encoded).unwrap();
142 assert_eq!(decoded.id(), id);
143 assert_eq!(decoded.payload().email, "a@example.com");
144 }
145}