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> {
173 let start = self.persisted_events.len() - n;
174 self.persisted_events[start..].iter()
175 }
176
177 pub fn iter_all(&self) -> impl DoubleEndedIterator<Item = &T> + Clone {
179 self.persisted_events
180 .iter()
181 .map(|e| &e.event)
182 .chain(self.new_events.iter().map(|e| &e.event))
183 }
184
185 pub fn load_first<E: EsEntity<Event = T>>(
189 events: impl IntoIterator<Item = GenericEvent<<T as EsEvent>::EntityId>>,
190 ) -> Result<Option<E>, EntityHydrationError> {
191 let mut current_id = None;
192 let mut current = None;
193 for e in events {
194 if current_id.is_none() {
195 current_id = Some(e.entity_id.clone());
196 current = Some(Self {
197 entity_id: e.entity_id.clone(),
198 persisted_events: Vec::new(),
199 new_events: Vec::new(),
200 });
201 }
202 if current_id.as_ref() != Some(&e.entity_id) {
203 break;
204 }
205 let cur = current.as_mut().expect("Could not get current");
206 let mut event_json = e.event;
207 if let Some(payload) = e.forgettable_payload {
208 crate::forgettable::inject_forgettable_payload(&mut event_json, payload);
209 }
210 cur.persisted_events.push(PersistedEvent {
211 entity_id: e.entity_id,
212 recorded_at: e.recorded_at,
213 sequence: e.sequence as usize,
214 event: serde_json::from_value(event_json)?,
215 context: e.context,
216 });
217 }
218 if let Some(current) = current {
219 Ok(Some(E::try_from_events(current)?))
220 } else {
221 Ok(None)
222 }
223 }
224
225 pub fn load_n<E: EsEntity<Event = T>>(
230 events: impl IntoIterator<Item = GenericEvent<<T as EsEvent>::EntityId>>,
231 n: usize,
232 ) -> Result<(Vec<E>, bool), EntityHydrationError> {
233 let mut ret: Vec<E> = Vec::new();
234 let mut current_id = None;
235 let mut current = None;
236 for e in events {
237 if current_id.as_ref() != Some(&e.entity_id) {
238 if let Some(current) = current.take() {
239 ret.push(E::try_from_events(current)?);
240 if ret.len() == n {
241 return Ok((ret, true));
242 }
243 }
244
245 current_id = Some(e.entity_id.clone());
246 current = Some(Self {
247 entity_id: e.entity_id.clone(),
248 persisted_events: Vec::new(),
249 new_events: Vec::new(),
250 });
251 }
252 let cur = current.as_mut().expect("Could not get current");
253 let mut event_json = e.event;
254 if let Some(payload) = e.forgettable_payload {
255 crate::forgettable::inject_forgettable_payload(&mut event_json, payload);
256 }
257 cur.persisted_events.push(PersistedEvent {
258 entity_id: e.entity_id,
259 recorded_at: e.recorded_at,
260 sequence: e.sequence as usize,
261 event: serde_json::from_value(event_json)?,
262 context: e.context,
263 });
264 }
265 if let Some(current) = current.take() {
266 ret.push(E::try_from_events(current)?);
267 }
268 Ok((ret, false))
269 }
270
271 #[doc(hidden)]
272 pub fn iter_new_events(&self) -> impl Iterator<Item = &EventWithContext<T>> {
273 self.new_events.iter()
274 }
275
276 #[doc(hidden)]
277 pub fn mark_new_events_persisted_at(
278 &mut self,
279 recorded_at: chrono::DateTime<chrono::Utc>,
280 ) -> usize {
281 let n = self.new_events.len();
282 let offset = self.persisted_events.len() + 1;
283 self.persisted_events
284 .extend(
285 self.new_events
286 .drain(..)
287 .enumerate()
288 .map(|(i, event)| PersistedEvent {
289 entity_id: self.entity_id.clone(),
290 recorded_at,
291 sequence: i + offset,
292 event: event.event,
293 context: event.context,
294 }),
295 );
296 n
297 }
298
299 #[doc(hidden)]
300 pub fn new_event_types(&self) -> Vec<String> {
301 self.new_events
302 .iter()
303 .map(|event| event.event.event_type().to_string())
304 .collect()
305 }
306
307 #[doc(hidden)]
308 pub fn serialize_new_events(&self) -> Vec<serde_json::Value> {
309 self.new_events
310 .iter()
311 .map(|event| serde_json::to_value(&event.event).expect("Failed to serialize event"))
312 .collect()
313 }
314
315 #[doc(hidden)]
321 pub fn forget_and_take(&mut self, mut forget_fn: impl FnMut(&mut T)) -> Self {
322 for persisted in &mut self.persisted_events {
323 forget_fn(&mut persisted.event);
324 }
325 let entity_id = self.entity_id.clone();
326 std::mem::replace(
327 self,
328 Self {
329 entity_id,
330 persisted_events: Vec::new(),
331 new_events: Vec::new(),
332 },
333 )
334 }
335
336 #[doc(hidden)]
337 pub fn serialize_new_event_contexts(&self) -> Option<Vec<crate::ContextData>> {
338 if <T as EsEvent>::event_context() {
339 let contexts = self
340 .new_events
341 .iter()
342 .map(|event| event.context.clone().expect("Missing context"))
343 .collect();
344
345 Some(contexts)
346 } else {
347 None
348 }
349 }
350}
351
352#[cfg(test)]
353mod tests {
354 use super::*;
355 use uuid::Uuid;
356
357 #[derive(Debug, serde::Serialize, serde::Deserialize)]
358 enum DummyEntityEvent {
359 Created(String),
360 }
361
362 impl EsEvent for DummyEntityEvent {
363 type EntityId = Uuid;
364 fn event_context() -> bool {
365 true
366 }
367 fn event_type(&self) -> &'static str {
368 match self {
369 Self::Created(_) => "created",
370 }
371 }
372 }
373
374 struct DummyEntity {
375 name: String,
376
377 events: EntityEvents<DummyEntityEvent>,
378 }
379
380 impl EsEntity for DummyEntity {
381 type Event = DummyEntityEvent;
382 type New = NewDummyEntity;
383
384 fn events_mut(&mut self) -> &mut EntityEvents<DummyEntityEvent> {
385 &mut self.events
386 }
387 fn events(&self) -> &EntityEvents<DummyEntityEvent> {
388 &self.events
389 }
390 }
391
392 impl TryFromEvents<DummyEntityEvent> for DummyEntity {
393 fn try_from_events(
394 events: EntityEvents<DummyEntityEvent>,
395 ) -> Result<Self, EntityHydrationError> {
396 let name = events
397 .iter_persisted()
398 .map(|e| match &e.event {
399 DummyEntityEvent::Created(name) => name.clone(),
400 })
401 .next()
402 .expect("Could not find name");
403 Ok(Self { name, events })
404 }
405 }
406
407 struct NewDummyEntity {}
408
409 impl IntoEvents<DummyEntityEvent> for NewDummyEntity {
410 fn into_events(self) -> EntityEvents<DummyEntityEvent> {
411 EntityEvents::init(
412 Uuid::parse_str("00000000-0000-0000-0000-000000000000").unwrap(),
413 vec![DummyEntityEvent::Created("".to_owned())],
414 )
415 }
416 }
417
418 #[test]
419 fn load_zero_events() {
420 let generic_events = vec![];
421 let res = EntityEvents::load_first::<DummyEntity>(generic_events);
422 assert!(matches!(res, Ok(None)));
423 }
424
425 #[test]
426 fn load_first() {
427 let generic_events = vec![GenericEvent {
428 entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000001").unwrap(),
429 sequence: 1,
430 event: serde_json::to_value(DummyEntityEvent::Created("dummy-name".to_owned()))
431 .expect("Could not serialize"),
432 context: None,
433 recorded_at: chrono::Utc::now(),
434 forgettable_payload: None,
435 }];
436 let entity: DummyEntity = EntityEvents::load_first(generic_events)
437 .expect("Could not load")
438 .expect("No entity found");
439 assert!(entity.name == "dummy-name");
440 }
441
442 #[test]
443 fn load_n() {
444 let generic_events = vec![
445 GenericEvent {
446 entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000002").unwrap(),
447 sequence: 1,
448 event: serde_json::to_value(DummyEntityEvent::Created("dummy-name".to_owned()))
449 .expect("Could not serialize"),
450 context: None,
451 recorded_at: chrono::Utc::now(),
452 forgettable_payload: None,
453 },
454 GenericEvent {
455 entity_id: Uuid::parse_str("00000000-0000-0000-0000-000000000003").unwrap(),
456 sequence: 1,
457 event: serde_json::to_value(DummyEntityEvent::Created("other-name".to_owned()))
458 .expect("Could not serialize"),
459 context: None,
460 recorded_at: chrono::Utc::now(),
461 forgettable_payload: None,
462 },
463 ];
464 let (entity, more): (Vec<DummyEntity>, _) =
465 EntityEvents::load_n(generic_events, 2).expect("Could not load");
466 assert!(!more);
467 assert_eq!(entity.len(), 2);
468 }
469}