Skip to main content

es_entity/
events.rs

1//! Manage events and related operations for event-sourcing.
2
3use chrono::{DateTime, Utc};
4
5use super::{error::EntityHydrationError, traits::*};
6
7/// An alias for iterator over the persisted events
8pub type LastPersisted<'a, E> = std::slice::Iter<'a, PersistedEvent<E>>;
9
10/// Represent the events in raw deserialized format when loaded from database
11///
12/// Events in the database are stored as JSON blobs and loaded initially as `GenericEvents<Id>` where `Id`
13/// belongs to the entity the events is a part of. Acts a bridge between database model and
14/// domain model when later converted to the `PersistedEvent` type internally
15pub struct GenericEvent<Id> {
16    pub entity_id: Id,
17    pub sequence: i32,
18    pub event: serde_json::Value,
19    pub context: Option<crate::ContextData>,
20    pub recorded_at: DateTime<Utc>,
21    pub forgettable_payload: Option<serde_json::Value>,
22}
23
24/// Strongly-typed event wrapper with metadata for successfully stored events.
25///
26/// Contains the event data along with persistence metadata (sequence, timestamp, entity_id).
27/// All `new_events` from [`EntityEvents`] are converted to this structure once persisted to construct
28/// entities, enabling event sourcing operations and other database operations.
29pub struct PersistedEvent<E: EsEvent> {
30    /// The identifier of the entity which the event is used to construct
31    pub entity_id: <E as EsEvent>::EntityId,
32    /// The timestamp which marks event persistence
33    pub recorded_at: DateTime<Utc>,
34    /// The sequence number of the event in the event stream
35    pub sequence: usize,
36    /// The event itself
37    pub event: E,
38    /// The context when the event was persisted
39    /// It is only popluated if 'event_context' set on EsEvent
40    pub context: Option<crate::ContextData>,
41}
42
43impl<E: Clone + EsEvent> Clone for PersistedEvent<E> {
44    fn clone(&self) -> Self {
45        PersistedEvent {
46            entity_id: self.entity_id.clone(),
47            recorded_at: self.recorded_at,
48            sequence: self.sequence,
49            event: self.event.clone(),
50            context: self.context.clone(),
51        }
52    }
53}
54
55pub struct EventWithContext<E: EsEvent> {
56    pub event: E,
57    pub context: Option<crate::ContextData>,
58}
59
60impl<E: Clone + EsEvent> Clone for EventWithContext<E> {
61    fn clone(&self) -> Self {
62        EventWithContext {
63            event: self.event.clone(),
64            context: self.context.clone(),
65        }
66    }
67}
68
69/// A [`Vec`] wrapper that manages event-stream of an entity with helpers for event-sourcing operations
70///
71/// Provides event sourcing operations for loading, appending, and persisting events in chronological
72/// sequence. Required field for all event-sourced entities to maintain their state change history.
73pub struct EntityEvents<T: EsEvent> {
74    /// The entity's id
75    pub entity_id: <T as EsEvent>::EntityId,
76    /// Events that have been persisted in database and marked
77    persisted_events: Vec<PersistedEvent<T>>,
78    /// New events that are yet to be persisted to track state changes
79    new_events: Vec<EventWithContext<T>>,
80}
81
82impl<T: Clone + EsEvent> Clone for EntityEvents<T> {
83    fn clone(&self) -> Self {
84        Self {
85            entity_id: self.entity_id.clone(),
86            persisted_events: self.persisted_events.clone(),
87            new_events: self.new_events.clone(),
88        }
89    }
90}
91
92impl<T> EntityEvents<T>
93where
94    T: EsEvent,
95{
96    /// Initializes a new `EntityEvents` instance with the given entity ID and initial events which is returned by [`IntoEvents`] method
97    pub fn init(id: <T as EsEvent>::EntityId, initial_events: impl IntoIterator<Item = T>) -> Self {
98        let context = if <T as EsEvent>::event_context() {
99            Some(crate::EventContext::data_for_storing())
100        } else {
101            None
102        };
103        let new_events = initial_events
104            .into_iter()
105            .map(|event| EventWithContext {
106                event,
107                context: context.clone(),
108            })
109            .collect();
110        Self {
111            entity_id: id,
112            persisted_events: Vec::new(),
113            new_events,
114        }
115    }
116
117    /// Returns a reference to the entity's identifier
118    pub fn id(&self) -> &<T as EsEvent>::EntityId {
119        &self.entity_id
120    }
121
122    /// Returns the timestamp of the first persisted event, indicating when the entity was created
123    pub fn entity_first_persisted_at(&self) -> Option<DateTime<Utc>> {
124        self.persisted_events.first().map(|e| e.recorded_at)
125    }
126
127    /// Returns the timestamp of the last persisted event, indicating when the entity was last modified
128    pub fn entity_last_modified_at(&self) -> Option<DateTime<Utc>> {
129        self.persisted_events.last().map(|e| e.recorded_at)
130    }
131
132    /// Appends a single new event to the entity's event stream to be persisted later
133    pub fn push(&mut self, event: T) {
134        let context = if <T as EsEvent>::event_context() {
135            Some(crate::EventContext::data_for_storing())
136        } else {
137            None
138        };
139        self.new_events.push(EventWithContext { event, context });
140    }
141
142    /// Appends multiple new events to the entity's event stream to be persisted later
143    pub fn extend(&mut self, events: impl IntoIterator<Item = T>) {
144        let context = if <T as EsEvent>::event_context() {
145            Some(crate::EventContext::data_for_storing())
146        } else {
147            None
148        };
149        self.new_events
150            .extend(events.into_iter().map(|event| EventWithContext {
151                event,
152                context: context.clone(),
153            }));
154    }
155
156    /// Returns true if there are any unpersisted events waiting to be saved
157    pub fn any_new(&self) -> bool {
158        !self.new_events.is_empty()
159    }
160
161    /// Returns the count of persisted events
162    pub fn len_persisted(&self) -> usize {
163        self.persisted_events.len()
164    }
165
166    /// Returns an iterator over all persisted events
167    pub fn iter_persisted(&self) -> impl DoubleEndedIterator<Item = &PersistedEvent<T>> + Clone {
168        self.persisted_events.iter()
169    }
170
171    /// Returns an iterator over the last `n` persisted events
172    ///
173    /// If fewer than `n` events have been persisted, all persisted events are
174    /// returned instead of panicking.
175    pub fn last_persisted(&self, n: usize) -> LastPersisted<'_, T> {
176        let start = self.persisted_events.len().saturating_sub(n);
177        self.persisted_events[start..].iter()
178    }
179
180    /// Returns an iterator over all events (both persisted and new) in chronological order
181    pub fn iter_all(&self) -> impl DoubleEndedIterator<Item = &T> + Clone {
182        self.persisted_events
183            .iter()
184            .map(|e| &e.event)
185            .chain(self.new_events.iter().map(|e| &e.event))
186    }
187
188    /// Loads and reconstructs the first entity from a stream of GenericEvents, marking events as `persisted`.
189    ///
190    /// Returns `Ok(None)` if no events are present, `Ok(Some(entity))` on success.
191    pub fn load_first<E: EsEntity<Event = T>>(
192        events: impl IntoIterator<Item = GenericEvent<<T as EsEvent>::EntityId>>,
193    ) -> Result<Option<E>, EntityHydrationError> {
194        let mut current_id = None;
195        let mut current = None;
196        for e in events {
197            if current_id.is_none() {
198                current_id = Some(e.entity_id.clone());
199                current = Some(Self {
200                    entity_id: e.entity_id.clone(),
201                    persisted_events: Vec::new(),
202                    new_events: Vec::new(),
203                });
204            }
205            if current_id.as_ref() != Some(&e.entity_id) {
206                break;
207            }
208            let cur = current.as_mut().expect("Could not get current");
209            let mut event_json = e.event;
210            if let Some(payload) = e.forgettable_payload {
211                crate::forgettable::inject_forgettable_payload(&mut event_json, payload);
212            }
213            cur.persisted_events.push(PersistedEvent {
214                entity_id: e.entity_id,
215                recorded_at: e.recorded_at,
216                sequence: e.sequence as usize,
217                event: serde_json::from_value(event_json)?,
218                context: e.context,
219            });
220        }
221        if let Some(current) = current {
222            Ok(Some(E::try_from_events(current)?))
223        } else {
224            Ok(None)
225        }
226    }
227
228    /// Loads and reconstructs up to `n` entities from a stream of GenericEvents.
229    /// Assumes the events are grouped by `id` and ordered by `sequence` per `id`.
230    ///
231    /// Returns both the entities and a flag indicating whether more entities were available in the stream.
232    pub fn load_n<E: EsEntity<Event = T>>(
233        events: impl IntoIterator<Item = GenericEvent<<T as EsEvent>::EntityId>>,
234        n: usize,
235    ) -> Result<(Vec<E>, bool), EntityHydrationError> {
236        let mut ret: Vec<E> = Vec::new();
237        let mut current_id = None;
238        let mut current = None;
239        for e in events {
240            if current_id.as_ref() != Some(&e.entity_id) {
241                if let Some(current) = current.take() {
242                    ret.push(E::try_from_events(current)?);
243                    if ret.len() == n {
244                        return Ok((ret, true));
245                    }
246                }
247
248                current_id = Some(e.entity_id.clone());
249                current = Some(Self {
250                    entity_id: e.entity_id.clone(),
251                    persisted_events: Vec::new(),
252                    new_events: Vec::new(),
253                });
254            }
255            let cur = current.as_mut().expect("Could not get current");
256            let mut event_json = e.event;
257            if let Some(payload) = e.forgettable_payload {
258                crate::forgettable::inject_forgettable_payload(&mut event_json, payload);
259            }
260            cur.persisted_events.push(PersistedEvent {
261                entity_id: e.entity_id,
262                recorded_at: e.recorded_at,
263                sequence: e.sequence as usize,
264                event: serde_json::from_value(event_json)?,
265                context: e.context,
266            });
267        }
268        if let Some(current) = current.take() {
269            ret.push(E::try_from_events(current)?);
270        }
271        Ok((ret, false))
272    }
273
274    #[doc(hidden)]
275    pub fn iter_new_events(&self) -> impl Iterator<Item = &EventWithContext<T>> {
276        self.new_events.iter()
277    }
278
279    #[doc(hidden)]
280    pub fn mark_new_events_persisted_at(
281        &mut self,
282        recorded_at: chrono::DateTime<chrono::Utc>,
283    ) -> usize {
284        let n = self.new_events.len();
285        let offset = self.persisted_events.len() + 1;
286        self.persisted_events
287            .extend(
288                self.new_events
289                    .drain(..)
290                    .enumerate()
291                    .map(|(i, event)| PersistedEvent {
292                        entity_id: self.entity_id.clone(),
293                        recorded_at,
294                        sequence: i + offset,
295                        event: event.event,
296                        context: event.context,
297                    }),
298            );
299        n
300    }
301
302    #[doc(hidden)]
303    pub fn new_event_types(&self) -> Vec<String> {
304        self.new_events
305            .iter()
306            .map(|event| event.event.event_type().to_string())
307            .collect()
308    }
309
310    #[doc(hidden)]
311    pub fn serialize_new_events(&self) -> Vec<serde_json::Value> {
312        self.new_events
313            .iter()
314            .map(|event| serde_json::to_value(&event.event).expect("Failed to serialize event"))
315            .collect()
316    }
317
318    /// Forgets all forgettable payloads in persisted events and returns the taken events.
319    ///
320    /// Applies `forget_fn` to each persisted event, then takes ownership of the event
321    /// stream, leaving `self` as an empty shell. The returned `EntityEvents` can be passed
322    /// to `TryFromEvents::try_from_events` to rebuild the entity with forgotten fields.
323    #[doc(hidden)]
324    pub fn forget_and_take(&mut self, mut forget_fn: impl FnMut(&mut T)) -> Self {
325        for persisted in &mut self.persisted_events {
326            forget_fn(&mut persisted.event);
327        }
328        let entity_id = self.entity_id.clone();
329        std::mem::replace(
330            self,
331            Self {
332                entity_id,
333                persisted_events: Vec::new(),
334                new_events: Vec::new(),
335            },
336        )
337    }
338
339    #[doc(hidden)]
340    pub fn serialize_new_event_contexts(&self) -> Option<Vec<crate::ContextData>> {
341        if <T as EsEvent>::event_context() {
342            let contexts = self
343                .new_events
344                .iter()
345                .map(|event| event.context.clone().expect("Missing context"))
346                .collect();
347
348            Some(contexts)
349        } else {
350            None
351        }
352    }
353}
354
355#[cfg(test)]
356mod tests {
357    use super::*;
358    use uuid::Uuid;
359
360    #[derive(Debug, serde::Serialize, serde::Deserialize)]
361    enum DummyEntityEvent {
362        Created(String),
363    }
364
365    impl EsEvent for DummyEntityEvent {
366        type EntityId = Uuid;
367        fn event_context() -> bool {
368            true
369        }
370        fn event_type(&self) -> &'static str {
371            match self {
372                Self::Created(_) => "created",
373            }
374        }
375    }
376
377    struct DummyEntity {
378        name: String,
379
380        events: EntityEvents<DummyEntityEvent>,
381    }
382
383    impl EsEntity for DummyEntity {
384        type Event = DummyEntityEvent;
385        type New = NewDummyEntity;
386
387        fn events_mut(&mut self) -> &mut EntityEvents<DummyEntityEvent> {
388            &mut self.events
389        }
390        fn events(&self) -> &EntityEvents<DummyEntityEvent> {
391            &self.events
392        }
393    }
394
395    impl TryFromEvents<DummyEntityEvent> for DummyEntity {
396        fn try_from_events(
397            events: EntityEvents<DummyEntityEvent>,
398        ) -> Result<Self, EntityHydrationError> {
399            let name = events
400                .iter_persisted()
401                .map(|e| match &e.event {
402                    DummyEntityEvent::Created(name) => name.clone(),
403                })
404                .next()
405                .expect("Could not find name");
406            Ok(Self { name, events })
407        }
408    }
409
410    struct NewDummyEntity {}
411
412    impl IntoEvents<DummyEntityEvent> for NewDummyEntity {
413        fn into_events(self) -> EntityEvents<DummyEntityEvent> {
414            EntityEvents::init(
415                Uuid::parse_str("00000000-0000-0000-0000-000000000000").unwrap(),
416                vec![DummyEntityEvent::Created("".to_owned())],
417            )
418        }
419    }
420
421    #[test]
422    fn load_zero_events() {
423        let generic_events = vec![];
424        let res = EntityEvents::load_first::<DummyEntity>(generic_events);
425        assert!(matches!(res, Ok(None)));
426    }
427
428    #[test]
429    fn load_first() {
430        let generic_events = vec![GenericEvent {
431            entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000001").unwrap(),
432            sequence: 1,
433            event: serde_json::to_value(DummyEntityEvent::Created("dummy-name".to_owned()))
434                .expect("Could not serialize"),
435            context: None,
436            recorded_at: chrono::Utc::now(),
437            forgettable_payload: None,
438        }];
439        let entity: DummyEntity = EntityEvents::load_first(generic_events)
440            .expect("Could not load")
441            .expect("No entity found");
442        assert!(entity.name == "dummy-name");
443    }
444
445    #[test]
446    fn load_n() {
447        let generic_events = vec![
448            GenericEvent {
449                entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000002").unwrap(),
450                sequence: 1,
451                event: serde_json::to_value(DummyEntityEvent::Created("dummy-name".to_owned()))
452                    .expect("Could not serialize"),
453                context: None,
454                recorded_at: chrono::Utc::now(),
455                forgettable_payload: None,
456            },
457            GenericEvent {
458                entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000003").unwrap(),
459                sequence: 1,
460                event: serde_json::to_value(DummyEntityEvent::Created("other-name".to_owned()))
461                    .expect("Could not serialize"),
462                context: None,
463                recorded_at: chrono::Utc::now(),
464                forgettable_payload: None,
465            },
466        ];
467        let (entity, more): (Vec<DummyEntity>, _) =
468            EntityEvents::load_n(generic_events, 2).expect("Could not load");
469        assert!(!more);
470        assert_eq!(entity.len(), 2);
471    }
472
473    #[test]
474    fn last_persisted_does_not_panic_when_n_exceeds_len() {
475        let generic_events = vec![GenericEvent {
476            entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000004").unwrap(),
477            sequence: 1,
478            event: serde_json::to_value(DummyEntityEvent::Created("dummy".to_owned()))
479                .expect("Could not serialize"),
480            context: None,
481            recorded_at: chrono::Utc::now(),
482            forgettable_payload: None,
483        }];
484        let entity: DummyEntity = EntityEvents::load_first(generic_events)
485            .expect("Could not load")
486            .expect("No entity found");
487        let events = entity.events();
488
489        // n == len works
490        assert_eq!(events.last_persisted(1).count(), 1);
491        // n > len clamps to all persisted events instead of underflowing
492        assert_eq!(events.last_persisted(10).count(), 1);
493        // n == 0 yields nothing
494        assert_eq!(events.last_persisted(0).count(), 0);
495    }
496}