1use async_trait::async_trait;
32use std::collections::HashMap;
33use std::marker::PhantomData;
34use std::sync::Arc;
35
36use backbone_messaging::crud_event::{CrudEvent, CrudEventPublisher, EventMetadata, NoOpCrudEventPublisher};
37
38use crate::persistence::traits::{CrudRepository, PersistentEntity, RepositoryError};
39
40#[derive(Debug, thiserror::Error)]
44pub enum ServiceError {
45 #[error("not found")]
46 NotFound,
47
48 #[error("already exists: {0}")]
49 AlreadyExists(String),
50
51 #[error("validation failed: {0}")]
52 Validation(String),
53
54 #[error("repository error: {0}")]
55 Repository(#[from] RepositoryError),
56
57 #[error("internal error: {0}")]
58 Internal(String),
59}
60
61pub type ServiceResult<T> = Result<T, ServiceError>;
62
63pub trait FromCreateDto<DTO>: Sized {
70 fn from_create_dto(dto: DTO) -> ServiceResult<Self>;
71}
72
73pub trait ApplyUpdateDto<DTO> {
75 fn apply_update(self, dto: DTO) -> ServiceResult<Self>
76 where
77 Self: Sized;
78}
79
80#[async_trait]
86pub trait ServiceLifecycle<E: PersistentEntity>: Send + Sync {
87 async fn before_create(&self, entity: &mut E) -> ServiceResult<()> {
88 let _ = entity;
89 Ok(())
90 }
91
92 async fn after_create(&self, entity: &E) -> ServiceResult<()> {
93 let _ = entity;
94 Ok(())
95 }
96
97 async fn before_update(&self, entity: &mut E) -> ServiceResult<()> {
98 let _ = entity;
99 Ok(())
100 }
101
102 async fn after_update(&self, entity: &E) -> ServiceResult<()> {
103 let _ = entity;
104 Ok(())
105 }
106
107 async fn before_delete(&self, entity: &E) -> ServiceResult<()> {
108 let _ = entity;
109 Ok(())
110 }
111
112 async fn after_delete(&self, id: &str) -> ServiceResult<()> {
113 let _ = id;
114 Ok(())
115 }
116}
117
118pub struct NoOpLifecycle<E> {
120 _phantom: PhantomData<E>,
121}
122
123impl<E> NoOpLifecycle<E> {
124 pub fn new() -> Self {
125 Self {
126 _phantom: PhantomData,
127 }
128 }
129}
130
131impl<E> Default for NoOpLifecycle<E> {
132 fn default() -> Self {
133 Self::new()
134 }
135}
136
137#[async_trait]
138impl<E: PersistentEntity> ServiceLifecycle<E> for NoOpLifecycle<E> {}
139
140pub struct GenericCrudService<E, C, U, R>
154where
155 E: PersistentEntity + Clone,
156 R: CrudRepository<E>,
157{
158 repository: Arc<R>,
159 lifecycle: Arc<dyn ServiceLifecycle<E>>,
160 event_publisher: Arc<dyn CrudEventPublisher<E>>,
161 _phantom: PhantomData<(C, U)>,
162}
163
164impl<E, C, U, R> GenericCrudService<E, C, U, R>
165where
166 E: PersistentEntity + Clone + FromCreateDto<C> + ApplyUpdateDto<U>,
167 C: Send + Sync + 'static,
168 U: Send + Sync + 'static,
169 R: CrudRepository<E>,
170{
171 pub fn new(
173 repository: Arc<R>,
174 lifecycle: Arc<dyn ServiceLifecycle<E>>,
175 event_publisher: Arc<dyn CrudEventPublisher<E>>,
176 ) -> Self {
177 Self {
178 repository,
179 lifecycle,
180 event_publisher,
181 _phantom: PhantomData,
182 }
183 }
184
185 pub fn with_lifecycle(repository: Arc<R>, lifecycle: Arc<dyn ServiceLifecycle<E>>) -> Self {
187 Self::new(repository, lifecycle, NoOpCrudEventPublisher::arc())
188 }
189
190 pub fn with_repository(repository: Arc<R>) -> Self {
192 Self::new(
193 repository,
194 Arc::new(NoOpLifecycle::new()),
195 NoOpCrudEventPublisher::arc(),
196 )
197 }
198
199 pub fn with_event_publisher(mut self, publisher: Arc<dyn CrudEventPublisher<E>>) -> Self {
201 self.event_publisher = publisher;
202 self
203 }
204
205 pub fn repository(&self) -> &R {
207 &self.repository
208 }
209
210 pub async fn list(
213 &self,
214 page: u32,
215 limit: u32,
216 filters: HashMap<String, String>,
217 ) -> ServiceResult<(Vec<E>, u64)> {
218 self.repository
219 .list_filtered(page, limit, filters)
220 .await
221 .map_err(ServiceError::Repository)
222 }
223
224 pub async fn list_with_info(
227 &self,
228 page: u32,
229 limit: u32,
230 filters: HashMap<String, String>,
231 ) -> ServiceResult<(Vec<E>, backbone_orm::repository::PaginationInfo)> {
232 self.repository
233 .list_filtered_with_info(page, limit, filters)
234 .await
235 .map_err(ServiceError::Repository)
236 }
237
238 pub async fn aggregate(
240 &self,
241 spec: &backbone_orm::repository::AggregateSpec,
242 filters: HashMap<String, String>,
243 ) -> ServiceResult<backbone_orm::repository::AggregateResult> {
244 self.repository
245 .aggregate_filtered(spec, filters)
246 .await
247 .map_err(ServiceError::Repository)
248 }
249
250 pub async fn create(&self, dto: C) -> ServiceResult<E> {
251 let mut entity = E::from_create_dto(dto)?;
252 self.lifecycle.before_create(&mut entity).await?;
253 let saved = self
254 .repository
255 .create(entity)
256 .await
257 .map_err(ServiceError::Repository)?;
258 self.lifecycle.after_create(&saved).await?;
259
260 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
261 let _ = self
262 .event_publisher
263 .publish(CrudEvent::Created { entity: saved.clone(), metadata: meta })
264 .await;
265
266 Ok(saved)
267 }
268
269 pub async fn get_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
270 self.repository
271 .find_by_id(id)
272 .await
273 .map_err(ServiceError::Repository)
274 }
275
276 pub async fn update(&self, id: &str, dto: U) -> ServiceResult<Option<E>> {
277 let existing = self
278 .repository
279 .find_by_id(id)
280 .await
281 .map_err(ServiceError::Repository)?;
282 let Some(before) = existing else {
283 return Ok(None);
284 };
285 let mut updated = before.clone().apply_update(dto)?;
286 self.lifecycle.before_update(&mut updated).await?;
287 let saved = self
288 .repository
289 .update(updated)
290 .await
291 .map_err(ServiceError::Repository)?;
292 self.lifecycle.after_update(&saved).await?;
293
294 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
295 let _ = self
296 .event_publisher
297 .publish(CrudEvent::Updated { before, after: saved.clone(), metadata: meta })
298 .await;
299
300 Ok(Some(saved))
301 }
302
303 pub async fn soft_delete(&self, id: &str) -> ServiceResult<bool> {
304 let existing = self
305 .repository
306 .find_by_id(id)
307 .await
308 .map_err(ServiceError::Repository)?;
309 let Some(entity) = existing else {
310 return Ok(false);
311 };
312 self.lifecycle.before_delete(&entity).await?;
313 let deleted = self
314 .repository
315 .soft_delete(id)
316 .await
317 .map_err(ServiceError::Repository)?;
318 self.lifecycle.after_delete(id).await?;
319
320 if deleted {
321 let meta = EventMetadata::new(id, std::any::type_name::<E>());
322 let _ = self
323 .event_publisher
324 .publish(CrudEvent::SoftDeleted { entity, metadata: meta })
325 .await;
326 }
327
328 Ok(deleted)
329 }
330
331 pub async fn restore(&self, id: &str) -> ServiceResult<Option<E>> {
332 let restored = self
333 .repository
334 .restore(id)
335 .await
336 .map_err(ServiceError::Repository)?;
337
338 if let Some(ref entity) = restored {
339 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
340 let _ = self
341 .event_publisher
342 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
343 .await;
344 }
345
346 Ok(restored)
347 }
348
349 pub async fn hard_delete(&self, id: &str) -> ServiceResult<bool> {
350 let deleted = self
351 .repository
352 .hard_delete(id)
353 .await
354 .map_err(ServiceError::Repository)?;
355
356 if deleted {
357 let meta = EventMetadata::new(id, std::any::type_name::<E>());
358 let _ = self
359 .event_publisher
360 .publish(CrudEvent::HardDeleted {
361 entity_id: id.to_string(),
362 metadata: meta,
363 })
364 .await;
365 }
366
367 Ok(deleted)
368 }
369
370 pub async fn list_deleted(
371 &self,
372 page: u32,
373 limit: u32,
374 ) -> ServiceResult<(Vec<E>, u64)> {
375 self.repository
376 .list_deleted(page, limit)
377 .await
378 .map_err(ServiceError::Repository)
379 }
380
381 pub async fn empty_trash(&self) -> ServiceResult<u64> {
382 self.repository
383 .empty_trash()
384 .await
385 .map_err(ServiceError::Repository)
386 }
387
388 pub async fn count(&self) -> ServiceResult<u64> {
389 self.repository
390 .count()
391 .await
392 .map_err(ServiceError::Repository)
393 }
394
395 pub async fn count_active(&self) -> ServiceResult<u64> {
396 self.repository
397 .count()
398 .await
399 .map_err(ServiceError::Repository)
400 }
401
402 pub async fn find_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
404 self.get_by_id(id).await
405 }
406
407 pub async fn permanent_delete(&self, id: &str) -> ServiceResult<bool> {
409 self.hard_delete(id).await
410 }
411
412 pub async fn count_deleted(&self) -> ServiceResult<u64> {
414 self.repository
415 .count_deleted()
416 .await
417 .map_err(ServiceError::Repository)
418 }
419
420 pub async fn get_deleted_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
422 self.repository
423 .find_by_id_including_deleted(id)
424 .await
425 .map_err(ServiceError::Repository)
426 }
427
428 pub async fn list_deleted_filtered(
430 &self,
431 page: u32,
432 limit: u32,
433 _filters: HashMap<String, String>,
434 ) -> ServiceResult<(Vec<E>, u64)> {
435 self.list_deleted(page, limit).await
436 }
437
438 pub async fn upsert(&self, dto: C) -> ServiceResult<E> {
440 self.create(dto).await
441 }
442
443 pub async fn partial_update(
447 &self,
448 id: &str,
449 fields: HashMap<String, serde_json::Value>,
450 ) -> ServiceResult<Option<E>>
451 where
452 E: serde::Serialize + serde::de::DeserializeOwned,
453 {
454 let existing = self
455 .repository
456 .find_by_id(id)
457 .await
458 .map_err(ServiceError::Repository)?;
459 let Some(entity) = existing else {
460 return Ok(None);
461 };
462 let mut map = serde_json::to_value(&entity)
464 .map_err(|e| ServiceError::Internal(e.to_string()))?;
465 if let serde_json::Value::Object(ref mut obj) = map {
466 for (key, value) in fields {
467 obj.insert(key, value);
468 }
469 }
470 let mut patched: E = serde_json::from_value(map)
471 .map_err(|e| ServiceError::Internal(e.to_string()))?;
472 self.lifecycle.before_update(&mut patched).await?;
473 let saved = self
474 .repository
475 .update(patched)
476 .await
477 .map_err(ServiceError::Repository)?;
478 self.lifecycle.after_update(&saved).await?;
479 Ok(Some(saved))
480 }
481
482 pub async fn bulk_create(&self, dtos: Vec<C>) -> ServiceResult<Vec<E>>
483 where
484 E: FromCreateDto<C>,
485 {
486 let mut results = Vec::with_capacity(dtos.len());
487 for dto in dtos {
488 let entity = self.create(dto).await?;
489 results.push(entity);
490 }
491 Ok(results)
492 }
493
494 pub async fn bulk_soft_delete(&self, ids: Vec<String>) -> ServiceResult<u64> {
503 check_batch_size(ids.len())?;
504 let ids = dedup_ids(ids);
505 let mut entities = Vec::with_capacity(ids.len());
510 for id in &ids {
511 match self
512 .repository
513 .find_by_id(id)
514 .await
515 .map_err(ServiceError::Repository)?
516 {
517 Some(e) => entities.push(e),
518 None => return Err(ServiceError::Validation(format!("id '{id}' not found"))),
519 }
520 }
521 for entity in &entities {
522 self.lifecycle.before_delete(entity).await?;
523 }
524 let affected = self
525 .repository
526 .bulk_soft_delete(&ids)
527 .await
528 .map_err(ServiceError::Repository)?;
529 for (id, entity) in ids.iter().zip(entities.into_iter()) {
530 self.lifecycle.after_delete(id).await?;
531 let meta = EventMetadata::new(id, std::any::type_name::<E>());
532 let _ = self
533 .event_publisher
534 .publish(CrudEvent::SoftDeleted { entity, metadata: meta })
535 .await;
536 }
537 Ok(affected)
538 }
539
540 pub async fn bulk_restore(&self, ids: Vec<String>) -> ServiceResult<Vec<E>> {
542 check_batch_size(ids.len())?;
543 let ids = dedup_ids(ids);
544 let restored = self
545 .repository
546 .bulk_restore(&ids)
547 .await
548 .map_err(ServiceError::Repository)?;
549 for entity in &restored {
550 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
551 let _ = self
552 .event_publisher
553 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
554 .await;
555 }
556 Ok(restored)
557 }
558
559 pub async fn restore_all(&self) -> ServiceResult<u64> {
562 let restored = self
563 .repository
564 .restore_all()
565 .await
566 .map_err(ServiceError::Repository)?;
567 for entity in &restored {
568 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
569 let _ = self
570 .event_publisher
571 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
572 .await;
573 }
574 Ok(restored.len() as u64)
575 }
576
577 pub async fn bulk_permanent_delete(&self, ids: Vec<String>) -> ServiceResult<u64> {
579 check_batch_size(ids.len())?;
580 let ids = dedup_ids(ids);
581 let affected = self
586 .repository
587 .bulk_hard_delete(&ids)
588 .await
589 .map_err(ServiceError::Repository)?;
590 for id in &ids {
591 let meta = EventMetadata::new(id, std::any::type_name::<E>());
592 let _ = self
593 .event_publisher
594 .publish(CrudEvent::HardDeleted {
595 entity_id: id.to_string(),
596 metadata: meta,
597 })
598 .await;
599 }
600 Ok(affected)
601 }
602
603 pub async fn bulk_update(&self, items: Vec<(String, U)>) -> ServiceResult<Vec<E>> {
605 check_batch_size(items.len())?;
606 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
607 return Err(ServiceError::Validation(format!(
608 "duplicate id '{dup}' in bulk update"
609 )));
610 }
611 let mut befores = Vec::with_capacity(items.len());
612 let mut prepared = Vec::with_capacity(items.len());
613 for (id, dto) in items {
614 let Some(before) = self
615 .repository
616 .find_by_id(&id)
617 .await
618 .map_err(ServiceError::Repository)?
619 else {
620 return Err(ServiceError::Validation(format!("id '{id}' not found")));
621 };
622 let mut updated = before.clone().apply_update(dto)?;
623 self.lifecycle.before_update(&mut updated).await?;
624 befores.push(before);
625 prepared.push(updated);
626 }
627 let saved = self
628 .repository
629 .bulk_update(prepared)
630 .await
631 .map_err(ServiceError::Repository)?;
632 self.publish_bulk_updates(befores, &saved).await?;
633 Ok(saved)
634 }
635
636 async fn publish_bulk_updates(&self, befores: Vec<E>, saved: &[E]) -> ServiceResult<()> {
640 for (before, after) in befores.into_iter().zip(saved.iter()) {
641 self.lifecycle.after_update(after).await?;
642 let meta = EventMetadata::new(after.entity_id(), std::any::type_name::<E>());
643 let _ = self
644 .event_publisher
645 .publish(CrudEvent::Updated {
646 before,
647 after: after.clone(),
648 metadata: meta,
649 })
650 .await;
651 }
652 Ok(())
653 }
654
655 pub async fn bulk_partial_update(
657 &self,
658 items: Vec<(String, HashMap<String, serde_json::Value>)>,
659 ) -> ServiceResult<Vec<E>>
660 where
661 E: serde::Serialize + serde::de::DeserializeOwned,
662 {
663 check_batch_size(items.len())?;
664 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
665 return Err(ServiceError::Validation(format!(
666 "duplicate id '{dup}' in bulk partial update"
667 )));
668 }
669 let mut befores = Vec::with_capacity(items.len());
670 let mut prepared = Vec::with_capacity(items.len());
671 for (id, fields) in items {
672 let Some(before) = self
673 .repository
674 .find_by_id(&id)
675 .await
676 .map_err(ServiceError::Repository)?
677 else {
678 return Err(ServiceError::Validation(format!("id '{id}' not found")));
679 };
680 let mut map = serde_json::to_value(&before)
681 .map_err(|e| ServiceError::Internal(e.to_string()))?;
682 if let serde_json::Value::Object(ref mut obj) = map {
683 for (key, value) in fields {
684 obj.insert(key, value);
685 }
686 }
687 let mut patched: E = serde_json::from_value(map)
688 .map_err(|e| ServiceError::Internal(e.to_string()))?;
689 self.lifecycle.before_update(&mut patched).await?;
690 befores.push(before);
691 prepared.push(patched);
692 }
693 let saved = self
694 .repository
695 .bulk_update(prepared)
696 .await
697 .map_err(ServiceError::Repository)?;
698 self.publish_bulk_updates(befores, &saved).await?;
699 Ok(saved)
700 }
701}
702
703pub const MAX_BATCH_SIZE: usize = 1000;
707
708fn check_batch_size(count: usize) -> Result<(), ServiceError> {
710 if count > MAX_BATCH_SIZE {
711 return Err(ServiceError::Validation(format!(
712 "Batch too large: {count} items exceeds the maximum of {MAX_BATCH_SIZE}."
713 )));
714 }
715 Ok(())
716}
717
718fn dedup_ids(ids: Vec<String>) -> Vec<String> {
722 let mut seen = std::collections::HashSet::new();
723 ids.into_iter().filter(|id| seen.insert(id.clone())).collect()
724}
725
726fn first_duplicate_id(ids: impl IntoIterator<Item = String>) -> Option<String> {
729 let mut seen = std::collections::HashSet::new();
730 for id in ids {
731 if !seen.insert(id.clone()) {
732 return Some(id);
733 }
734 }
735 None
736}
737
738#[async_trait::async_trait]
749impl<E, C, U, R> crate::http::CrudService<E, C, U> for GenericCrudService<E, C, U, R>
750where
751 E: PersistentEntity
752 + Clone
753 + serde::Serialize
754 + serde::de::DeserializeOwned
755 + FromCreateDto<C>
756 + ApplyUpdateDto<U>
757 + Send
758 + Sync
759 + 'static,
760 C: Send + Sync + 'static,
761 U: Send + Sync + 'static,
762 R: CrudRepository<E> + Send + Sync + 'static,
763{
764 type Error = ServiceError;
765
766 fn entity_name() -> &'static str {
767 std::any::type_name::<E>()
768 }
769
770 async fn fetch_related_json(&self, table: &str, ids: &[String]) -> Vec<serde_json::Value> {
771 self.repository().fetch_related_json(table, ids).await
772 }
773
774 async fn list(
775 &self,
776 page: u32,
777 limit: u32,
778 filters: HashMap<String, String>,
779 ) -> Result<(Vec<E>, u64), ServiceError> {
780 self.list(page, limit, filters).await
781 }
782
783 async fn list_with_info(
784 &self,
785 page: u32,
786 limit: u32,
787 filters: HashMap<String, String>,
788 ) -> Result<(Vec<E>, backbone_orm::repository::PaginationInfo), ServiceError> {
789 self.list_with_info(page, limit, filters).await
790 }
791
792 async fn aggregate(
793 &self,
794 spec: &backbone_orm::repository::AggregateSpec,
795 filters: HashMap<String, String>,
796 ) -> Result<backbone_orm::repository::AggregateResult, ServiceError> {
797 self.aggregate(spec, filters).await
798 }
799
800 fn table_name(&self) -> Option<&str> {
801 self.repository.table_name()
802 }
803
804 async fn create(&self, dto: C) -> Result<E, ServiceError> {
805 self.create(dto).await
806 }
807
808 async fn get_by_id(&self, id: &str) -> Result<Option<E>, ServiceError> {
809 self.get_by_id(id).await
810 }
811
812 async fn update(&self, id: &str, dto: U) -> Result<Option<E>, ServiceError> {
813 self.update(id, dto).await
814 }
815
816 async fn partial_update(
817 &self,
818 id: &str,
819 fields: HashMap<String, serde_json::Value>,
820 ) -> Result<Option<E>, ServiceError> {
821 self.partial_update(id, fields).await
822 }
823
824 async fn soft_delete(&self, id: &str) -> Result<bool, ServiceError> {
825 self.soft_delete(id).await
826 }
827
828 async fn bulk_create(&self, items: Vec<C>) -> Result<Vec<E>, ServiceError> {
829 self.bulk_create(items).await
830 }
831
832 async fn upsert(&self, dto: C) -> Result<E, ServiceError> {
833 self.upsert(dto).await
834 }
835
836 async fn list_deleted(&self, page: u32, limit: u32) -> Result<(Vec<E>, u64), ServiceError> {
837 self.list_deleted(page, limit).await
838 }
839
840 async fn restore(&self, id: &str) -> Result<Option<E>, ServiceError> {
841 self.restore(id).await
842 }
843
844 async fn empty_trash(&self) -> Result<u64, ServiceError> {
845 self.empty_trash().await
846 }
847
848 async fn get_deleted_by_id(&self, id: &str) -> Result<Option<E>, ServiceError> {
849 self.get_deleted_by_id(id).await
850 }
851
852 async fn permanent_delete(&self, id: &str) -> Result<bool, ServiceError> {
853 self.permanent_delete(id).await
854 }
855
856 async fn list_deleted_filtered(
857 &self,
858 page: u32,
859 limit: u32,
860 filters: HashMap<String, String>,
861 ) -> Result<(Vec<E>, u64), ServiceError> {
862 self.list_deleted_filtered(page, limit, filters).await
863 }
864
865 async fn count_active(&self) -> Result<u64, ServiceError> {
866 self.count_active().await
867 }
868
869 async fn count_deleted(&self) -> Result<u64, ServiceError> {
870 self.count_deleted().await
871 }
872
873 async fn bulk_soft_delete(&self, ids: Vec<String>) -> Result<u64, ServiceError> {
874 self.bulk_soft_delete(ids).await
875 }
876
877 async fn bulk_restore(&self, ids: Vec<String>) -> Result<Vec<E>, ServiceError> {
878 self.bulk_restore(ids).await
879 }
880
881 async fn bulk_permanent_delete(&self, ids: Vec<String>) -> Result<u64, ServiceError> {
882 self.bulk_permanent_delete(ids).await
883 }
884
885 async fn restore_all(&self) -> Result<u64, ServiceError> {
886 self.restore_all().await
887 }
888
889 async fn bulk_update(&self, items: Vec<(String, U)>) -> Result<Vec<E>, ServiceError> {
890 self.bulk_update(items).await
891 }
892
893 async fn bulk_partial_update(
894 &self,
895 items: Vec<(String, HashMap<String, serde_json::Value>)>,
896 ) -> Result<Vec<E>, ServiceError> {
897 self.bulk_partial_update(items).await
898 }
899}
900
901impl<E, C, U, R> Clone for GenericCrudService<E, C, U, R>
903where
904 E: PersistentEntity + Clone,
905 R: CrudRepository<E>,
906{
907 fn clone(&self) -> Self {
908 Self {
909 repository: self.repository.clone(),
910 lifecycle: self.lifecycle.clone(),
911 event_publisher: self.event_publisher.clone(),
912 _phantom: PhantomData,
913 }
914 }
915}
916
917#[cfg(test)]
918mod tests {
919 use super::*;
920 use chrono::{DateTime, Utc};
921 use serde::{Deserialize, Serialize};
922
923 #[derive(Debug, Clone, Serialize, Deserialize)]
926 struct Widget {
927 id: String,
928 name: String,
929 #[serde(skip_serializing_if = "Option::is_none")]
930 deleted_at: Option<DateTime<Utc>>,
931 created_at: DateTime<Utc>,
932 updated_at: DateTime<Utc>,
933 }
934
935 impl crate::persistence::traits::PersistentEntity for Widget {
936 fn entity_id(&self) -> String {
937 self.id.clone()
938 }
939 fn set_entity_id(&mut self, id: String) {
940 self.id = id;
941 }
942 fn created_at(&self) -> Option<DateTime<Utc>> {
943 Some(self.created_at)
944 }
945 fn set_created_at(&mut self, ts: DateTime<Utc>) {
946 self.created_at = ts;
947 }
948 fn updated_at(&self) -> Option<DateTime<Utc>> {
949 Some(self.updated_at)
950 }
951 fn set_updated_at(&mut self, ts: DateTime<Utc>) {
952 self.updated_at = ts;
953 }
954 fn deleted_at(&self) -> Option<DateTime<Utc>> {
955 self.deleted_at
956 }
957 fn set_deleted_at(&mut self, ts: Option<DateTime<Utc>>) {
958 self.deleted_at = ts;
959 }
960 }
961
962 struct CreateWidgetDto {
963 name: String,
964 }
965
966 struct UpdateWidgetDto {
967 name: String,
968 }
969
970 impl FromCreateDto<CreateWidgetDto> for Widget {
971 fn from_create_dto(dto: CreateWidgetDto) -> ServiceResult<Self> {
972 Ok(Widget {
973 id: uuid::Uuid::new_v4().to_string(),
974 name: dto.name,
975 deleted_at: None,
976 created_at: Utc::now(),
977 updated_at: Utc::now(),
978 })
979 }
980 }
981
982 impl ApplyUpdateDto<UpdateWidgetDto> for Widget {
983 fn apply_update(mut self, dto: UpdateWidgetDto) -> ServiceResult<Self> {
984 self.name = dto.name;
985 self.updated_at = Utc::now();
986 Ok(self)
987 }
988 }
989
990 struct InMemoryWidgetRepo {
992 data: tokio::sync::Mutex<Vec<Widget>>,
993 }
994
995 impl InMemoryWidgetRepo {
996 fn new() -> Self {
997 Self {
998 data: tokio::sync::Mutex::new(Vec::new()),
999 }
1000 }
1001 }
1002
1003 #[async_trait]
1004 impl crate::persistence::traits::CrudRepository<Widget> for InMemoryWidgetRepo {
1005 async fn create(&self, entity: Widget) -> Result<Widget, RepositoryError> {
1006 self.data.lock().await.push(entity.clone());
1007 Ok(entity)
1008 }
1009
1010 async fn find_by_id(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1011 Ok(self
1012 .data
1013 .lock()
1014 .await
1015 .iter()
1016 .find(|w| w.id == id && w.deleted_at.is_none())
1017 .cloned())
1018 }
1019
1020 async fn update(&self, entity: Widget) -> Result<Widget, RepositoryError> {
1021 let mut data = self.data.lock().await;
1022 data.retain(|w| w.id != entity.id);
1023 data.push(entity.clone());
1024 Ok(entity)
1025 }
1026
1027 async fn soft_delete(&self, id: &str) -> Result<bool, RepositoryError> {
1028 let mut data = self.data.lock().await;
1029 if let Some(w) = data.iter_mut().find(|w| w.id == id) {
1030 w.deleted_at = Some(Utc::now());
1031 return Ok(true);
1032 }
1033 Ok(false)
1034 }
1035
1036 async fn find_by_id_including_deleted(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1037 Ok(self.data.lock().await.iter().find(|w| w.id == id).cloned())
1038 }
1039
1040 async fn restore(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1041 let mut data = self.data.lock().await;
1042 if let Some(w) = data.iter_mut().find(|w| w.id == id) {
1043 w.deleted_at = None;
1044 return Ok(Some(w.clone()));
1045 }
1046 Ok(None)
1047 }
1048
1049 async fn hard_delete(&self, id: &str) -> Result<bool, RepositoryError> {
1050 let mut data = self.data.lock().await;
1051 let before = data.len();
1052 data.retain(|w| w.id != id);
1053 Ok(data.len() < before)
1054 }
1055
1056 async fn list(
1057 &self,
1058 page: u32,
1059 limit: u32,
1060 ) -> Result<(Vec<Widget>, u64), RepositoryError> {
1061 let data = self.data.lock().await;
1062 let active: Vec<_> = data.iter().filter(|w| w.deleted_at.is_none()).cloned().collect();
1063 let total = active.len() as u64;
1064 let offset = ((page.saturating_sub(1)) as usize) * (limit as usize);
1065 let limit = limit as usize;
1066 let page = active.into_iter().skip(offset).take(limit).collect();
1067 Ok((page, total))
1068 }
1069
1070 async fn list_deleted(
1071 &self,
1072 page: u32,
1073 limit: u32,
1074 ) -> Result<(Vec<Widget>, u64), RepositoryError> {
1075 let data = self.data.lock().await;
1076 let deleted: Vec<_> = data.iter().filter(|w| w.deleted_at.is_some()).cloned().collect();
1077 let total = deleted.len() as u64;
1078 let offset = ((page.saturating_sub(1)) as usize) * (limit as usize);
1079 let limit = limit as usize;
1080 let page = deleted.into_iter().skip(offset).take(limit).collect();
1081 Ok((page, total))
1082 }
1083
1084 async fn count(&self) -> Result<u64, RepositoryError> {
1085 let data = self.data.lock().await;
1086 Ok(data.iter().filter(|w| w.deleted_at.is_none()).count() as u64)
1087 }
1088
1089 async fn count_deleted(&self) -> Result<u64, RepositoryError> {
1090 let data = self.data.lock().await;
1091 Ok(data.iter().filter(|w| w.deleted_at.is_some()).count() as u64)
1092 }
1093
1094 async fn bulk_create(&self, entities: Vec<Widget>) -> Result<Vec<Widget>, RepositoryError> {
1095 let mut data = self.data.lock().await;
1096 data.extend(entities.clone());
1097 Ok(entities)
1098 }
1099
1100 async fn empty_trash(&self) -> Result<u64, RepositoryError> {
1101 let mut data = self.data.lock().await;
1102 let before = data.len();
1103 data.retain(|w| w.deleted_at.is_none());
1104 Ok((before - data.len()) as u64)
1105 }
1106 }
1107
1108 #[tokio::test]
1109 async fn create_and_get_roundtrip() {
1110 let repo = Arc::new(InMemoryWidgetRepo::new());
1111 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1112 GenericCrudService::with_repository(repo);
1113
1114 let widget = service
1115 .create(CreateWidgetDto {
1116 name: "sprocket".into(),
1117 })
1118 .await
1119 .unwrap();
1120
1121 let found = service.get_by_id(&widget.id).await.unwrap();
1122 assert!(found.is_some());
1123 assert_eq!(found.unwrap().name, "sprocket");
1124 }
1125
1126 #[tokio::test]
1127 async fn soft_delete_hides_from_list() {
1128 let repo = Arc::new(InMemoryWidgetRepo::new());
1129 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1130 GenericCrudService::with_repository(repo);
1131
1132 let w = service.create(CreateWidgetDto { name: "w".into() }).await.unwrap();
1133 service.soft_delete(&w.id).await.unwrap();
1134
1135 let (items, _) = service.list(1, 20, Default::default()).await.unwrap();
1136 assert!(items.is_empty());
1137
1138 let (deleted, _) = service.list_deleted(1, 20).await.unwrap();
1139 assert_eq!(deleted.len(), 1);
1140 }
1141
1142 #[tokio::test]
1143 async fn update_publishes_event_and_returns_new_entity() {
1144 let repo = Arc::new(InMemoryWidgetRepo::new());
1145 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1146 GenericCrudService::with_repository(repo);
1147
1148 let w = service.create(CreateWidgetDto { name: "old".into() }).await.unwrap();
1149 let updated = service
1150 .update(&w.id, UpdateWidgetDto { name: "new".into() })
1151 .await
1152 .unwrap();
1153 assert_eq!(updated.unwrap().name, "new");
1154 }
1155
1156 fn svc() -> GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> {
1159 GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()))
1160 }
1161
1162 #[tokio::test]
1163 async fn bulk_soft_delete_success() {
1164 let service = svc();
1165 let mut ids = Vec::new();
1166 for n in 0..3 {
1167 ids.push(service.create(CreateWidgetDto { name: format!("w{n}") }).await.unwrap().id);
1168 }
1169 let n = service.bulk_soft_delete(ids).await.unwrap();
1170 assert_eq!(n, 3);
1171 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1172 assert!(active.is_empty());
1173 }
1174
1175 #[tokio::test]
1176 async fn bulk_soft_delete_is_all_or_nothing_on_missing_id() {
1177 let service = svc();
1178 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1179 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1180
1181 let res = service
1183 .bulk_soft_delete(vec![a.id.clone(), b.id.clone(), "does-not-exist".into()])
1184 .await;
1185 assert!(res.is_err());
1186
1187 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1188 assert_eq!(active.len(), 2, "no rows should have been deleted");
1189 }
1190
1191 #[tokio::test]
1192 async fn bulk_restore_and_restore_all() {
1193 let service = svc();
1194 let mut ids = Vec::new();
1195 for n in 0..3 {
1196 ids.push(service.create(CreateWidgetDto { name: format!("w{n}") }).await.unwrap().id);
1197 }
1198 service.bulk_soft_delete(ids.clone()).await.unwrap();
1199
1200 let restored = service.bulk_restore(vec![ids[0].clone(), ids[1].clone()]).await.unwrap();
1202 assert_eq!(restored.len(), 2);
1203
1204 let n = service.restore_all().await.unwrap();
1206 assert_eq!(n, 1);
1207
1208 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1209 assert_eq!(active.len(), 3);
1210 }
1211
1212 #[tokio::test]
1213 async fn bulk_update_is_all_or_nothing_on_missing_id() {
1214 let service = svc();
1215 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1216
1217 let res = service
1218 .bulk_update(vec![
1219 (a.id.clone(), UpdateWidgetDto { name: "A2".into() }),
1220 ("missing".into(), UpdateWidgetDto { name: "X".into() }),
1221 ])
1222 .await;
1223 assert!(res.is_err());
1224
1225 let still = service.get_by_id(&a.id).await.unwrap().unwrap();
1227 assert_eq!(still.name, "a");
1228 }
1229
1230 #[tokio::test]
1231 async fn bulk_partial_update_success() {
1232 let service = svc();
1233 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1234 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1235
1236 let mut patch_a = std::collections::HashMap::new();
1237 patch_a.insert("name".to_string(), serde_json::json!("a2"));
1238 let mut patch_b = std::collections::HashMap::new();
1239 patch_b.insert("name".to_string(), serde_json::json!("b2"));
1240
1241 let saved = service
1242 .bulk_partial_update(vec![(a.id.clone(), patch_a), (b.id.clone(), patch_b)])
1243 .await
1244 .unwrap();
1245 assert_eq!(saved.len(), 2);
1246 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a2");
1247 assert_eq!(service.get_by_id(&b.id).await.unwrap().unwrap().name, "b2");
1248 }
1249
1250 #[tokio::test]
1251 async fn bulk_soft_delete_tolerates_duplicate_ids() {
1252 let service = svc();
1253 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1254 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1255
1256 let n = service
1258 .bulk_soft_delete(vec![a.id.clone(), a.id.clone(), b.id.clone()])
1259 .await
1260 .unwrap();
1261 assert_eq!(n, 2, "two distinct rows soft-deleted");
1262
1263 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1264 assert!(active.is_empty());
1265 }
1266
1267 #[tokio::test]
1268 async fn bulk_update_rejects_duplicate_ids() {
1269 let service = svc();
1270 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1271
1272 let res = service
1274 .bulk_update(vec![
1275 (a.id.clone(), UpdateWidgetDto { name: "first".into() }),
1276 (a.id.clone(), UpdateWidgetDto { name: "second".into() }),
1277 ])
1278 .await;
1279 assert!(res.is_err());
1280 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1281 }
1282
1283 #[tokio::test]
1284 async fn bulk_soft_delete_rejects_oversized_batch() {
1285 let service = svc();
1286 let ids: Vec<String> = (0..MAX_BATCH_SIZE + 1).map(|n| format!("id-{n}")).collect();
1287 assert!(service.bulk_soft_delete(ids).await.is_err());
1289 }
1290}