1use chrono::{DateTime, Utc};
4
5use super::{error::EntityHydrationError, traits::*};
6
7pub type LastPersisted<'a, E> = std::slice::Iter<'a, PersistedEvent<E>>;
9
10pub 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
24pub struct PersistedEvent<E: EsEvent> {
30 pub entity_id: <E as EsEvent>::EntityId,
32 pub recorded_at: DateTime<Utc>,
34 pub sequence: usize,
36 pub event: E,
38 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
69pub struct EntityEvents<T: EsEvent> {
74 pub entity_id: <T as EsEvent>::EntityId,
76 persisted_events: Vec<PersistedEvent<T>>,
78 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 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 pub fn id(&self) -> &<T as EsEvent>::EntityId {
119 &self.entity_id
120 }
121
122 pub fn entity_first_persisted_at(&self) -> Option<DateTime<Utc>> {
124 self.persisted_events.first().map(|e| e.recorded_at)
125 }
126
127 pub fn entity_last_modified_at(&self) -> Option<DateTime<Utc>> {
129 self.persisted_events.last().map(|e| e.recorded_at)
130 }
131
132 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 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 pub fn any_new(&self) -> bool {
158 !self.new_events.is_empty()
159 }
160
161 pub fn len_persisted(&self) -> usize {
163 self.persisted_events.len()
164 }
165
166 pub fn iter_persisted(&self) -> impl DoubleEndedIterator<Item = &PersistedEvent<T>> + Clone {
168 self.persisted_events.iter()
169 }
170
171 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 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 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 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 #[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 assert_eq!(events.last_persisted(1).count(), 1);
491 assert_eq!(events.last_persisted(10).count(), 1);
493 assert_eq!(events.last_persisted(0).count(), 0);
495 }
496}