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("validation failed: {}", .0.iter().map(|v| v.message.as_str()).collect::<Vec<_>>().join("; "))]
58 Violations(Vec<crate::violation::Violation>),
59
60 #[error("repository error: {0}")]
61 Repository(#[from] RepositoryError),
62
63 #[error("internal error: {0}")]
64 Internal(String),
65}
66
67pub type ServiceResult<T> = Result<T, ServiceError>;
68
69pub trait FromCreateDto<DTO>: Sized {
76 fn from_create_dto(dto: DTO) -> ServiceResult<Self>;
77}
78
79pub trait ApplyUpdateDto<DTO> {
81 fn apply_update(self, dto: DTO) -> ServiceResult<Self>
82 where
83 Self: Sized;
84}
85
86#[async_trait]
92pub trait ServiceLifecycle<E: PersistentEntity>: Send + Sync {
93 async fn before_create(&self, entity: &mut E) -> ServiceResult<()> {
94 let _ = entity;
95 Ok(())
96 }
97
98 async fn after_create(&self, entity: &E) -> ServiceResult<()> {
99 let _ = entity;
100 Ok(())
101 }
102
103 async fn before_update(&self, entity: &mut E) -> ServiceResult<()> {
104 let _ = entity;
105 Ok(())
106 }
107
108 async fn after_update(&self, entity: &E) -> ServiceResult<()> {
109 let _ = entity;
110 Ok(())
111 }
112
113 async fn before_delete(&self, entity: &E) -> ServiceResult<()> {
114 let _ = entity;
115 Ok(())
116 }
117
118 async fn after_delete(&self, id: &str) -> ServiceResult<()> {
119 let _ = id;
120 Ok(())
121 }
122
123 async fn before_restore(&self, entity: &E) -> ServiceResult<()> {
125 let _ = entity;
126 Ok(())
127 }
128
129 async fn before_hard_delete(&self, entity: &E) -> ServiceResult<()> {
132 let _ = entity;
133 Ok(())
134 }
135
136 async fn before_restore_all(&self) -> ServiceResult<()> {
139 Ok(())
140 }
141
142 async fn before_empty_trash(&self) -> ServiceResult<()> {
145 Ok(())
146 }
147}
148
149pub struct NoOpLifecycle<E> {
151 _phantom: PhantomData<E>,
152}
153
154impl<E> NoOpLifecycle<E> {
155 pub fn new() -> Self {
156 Self {
157 _phantom: PhantomData,
158 }
159 }
160}
161
162impl<E> Default for NoOpLifecycle<E> {
163 fn default() -> Self {
164 Self::new()
165 }
166}
167
168#[async_trait]
169impl<E: PersistentEntity> ServiceLifecycle<E> for NoOpLifecycle<E> {}
170
171pub struct GenericCrudService<E, C, U, R>
185where
186 E: PersistentEntity + Clone,
187 R: CrudRepository<E>,
188{
189 repository: Arc<R>,
190 lifecycle: Arc<dyn ServiceLifecycle<E>>,
191 event_publisher: Arc<dyn CrudEventPublisher<E>>,
192 guard: Arc<dyn crate::write_guard::WriteGuard<E>>,
193 _phantom: PhantomData<(C, U)>,
194}
195
196impl<E, C, U, R> GenericCrudService<E, C, U, R>
197where
198 E: PersistentEntity + Clone + FromCreateDto<C> + ApplyUpdateDto<U>,
199 C: Send + Sync + 'static,
200 U: Send + Sync + 'static,
201 R: CrudRepository<E>,
202{
203 pub fn new(
205 repository: Arc<R>,
206 lifecycle: Arc<dyn ServiceLifecycle<E>>,
207 event_publisher: Arc<dyn CrudEventPublisher<E>>,
208 ) -> Self {
209 Self {
210 repository,
211 lifecycle,
212 event_publisher,
213 guard: Arc::new(crate::write_guard::AllowAll),
214 _phantom: PhantomData,
215 }
216 }
217
218 pub fn with_guard(mut self, guard: Arc<dyn crate::write_guard::WriteGuard<E>>) -> Self {
221 self.guard = guard;
222 self
223 }
224
225 async fn consult(
227 &self,
228 kind: crate::write_guard::WriteKind,
229 before: Option<&E>,
230 after: Option<&E>,
231 ) -> ServiceResult<()> {
232 let ctx = crate::write_guard::WriteCtx { kind, before, after };
233 let outcome = self.guard.check(&ctx).await;
234 for v in &outcome.shadow {
235 tracing::warn!(
236 target: "backbone::write_guard",
237 entity = std::any::type_name::<E>(),
238 kind = ?kind,
239 code = %v.code,
240 path = %v.path,
241 "a rule running in shadow would refuse this write: {}",
242 v.message
243 );
244 }
245 if outcome.refuse.is_empty() {
246 Ok(())
247 } else {
248 Err(ServiceError::Violations(outcome.refuse))
249 }
250 }
251
252 pub fn with_lifecycle(repository: Arc<R>, lifecycle: Arc<dyn ServiceLifecycle<E>>) -> Self {
254 Self::new(repository, lifecycle, NoOpCrudEventPublisher::arc())
255 }
256
257 pub fn with_repository(repository: Arc<R>) -> Self {
259 Self::new(
260 repository,
261 Arc::new(NoOpLifecycle::new()),
262 NoOpCrudEventPublisher::arc(),
263 )
264 }
265
266 pub fn with_event_publisher(mut self, publisher: Arc<dyn CrudEventPublisher<E>>) -> Self {
268 self.event_publisher = publisher;
269 self
270 }
271
272 pub fn repository(&self) -> &R {
274 &self.repository
275 }
276
277 pub async fn list(
280 &self,
281 page: u32,
282 limit: u32,
283 filters: HashMap<String, String>,
284 ) -> ServiceResult<(Vec<E>, u64)> {
285 self.repository
286 .list_filtered(page, limit, filters)
287 .await
288 .map_err(ServiceError::Repository)
289 }
290
291 pub async fn list_with_info(
294 &self,
295 page: u32,
296 limit: u32,
297 filters: HashMap<String, String>,
298 ) -> ServiceResult<(Vec<E>, backbone_orm::repository::PaginationInfo)> {
299 self.repository
300 .list_filtered_with_info(page, limit, filters)
301 .await
302 .map_err(ServiceError::Repository)
303 }
304
305 pub async fn aggregate(
307 &self,
308 spec: &backbone_orm::repository::AggregateSpec,
309 filters: HashMap<String, String>,
310 ) -> ServiceResult<backbone_orm::repository::AggregateResult> {
311 self.repository
312 .aggregate_filtered(spec, filters)
313 .await
314 .map_err(ServiceError::Repository)
315 }
316
317 pub async fn create(&self, dto: C) -> ServiceResult<E> {
318 let mut entity = E::from_create_dto(dto)?;
319 self.consult(crate::write_guard::WriteKind::Create, None, Some(&entity)).await?;
320 self.lifecycle.before_create(&mut entity).await?;
321 let saved = self
322 .repository
323 .create(entity)
324 .await
325 .map_err(ServiceError::Repository)?;
326 self.lifecycle.after_create(&saved).await?;
327
328 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
329 let _ = self
330 .event_publisher
331 .publish(CrudEvent::Created { entity: saved.clone(), metadata: meta })
332 .await;
333
334 Ok(saved)
335 }
336
337 pub async fn get_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
338 self.repository
339 .find_by_id(id)
340 .await
341 .map_err(ServiceError::Repository)
342 }
343
344 pub async fn update(&self, id: &str, dto: U) -> ServiceResult<Option<E>> {
345 let existing = self
346 .repository
347 .find_by_id(id)
348 .await
349 .map_err(ServiceError::Repository)?;
350 let Some(before) = existing else {
351 return Ok(None);
352 };
353 let mut updated = before.clone().apply_update(dto)?;
354 refuse_protected_changes::<E>(&as_object(&before)?, &as_object(&updated)?)?;
355 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&updated)).await?;
356 self.lifecycle.before_update(&mut updated).await?;
357 let saved = self
358 .repository
359 .update(updated)
360 .await
361 .map_err(ServiceError::Repository)?;
362 self.lifecycle.after_update(&saved).await?;
363
364 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
365 let _ = self
366 .event_publisher
367 .publish(CrudEvent::Updated { before, after: saved.clone(), metadata: meta })
368 .await;
369
370 Ok(Some(saved))
371 }
372
373 pub async fn soft_delete(&self, id: &str) -> ServiceResult<bool> {
374 let existing = self
375 .repository
376 .find_by_id(id)
377 .await
378 .map_err(ServiceError::Repository)?;
379 let Some(entity) = existing else {
380 return Ok(false);
381 };
382 self.consult(crate::write_guard::WriteKind::Delete, Some(&entity), None).await?;
383 self.lifecycle.before_delete(&entity).await?;
384 let deleted = self
385 .repository
386 .soft_delete(id)
387 .await
388 .map_err(ServiceError::Repository)?;
389 self.lifecycle.after_delete(id).await?;
390
391 if deleted {
392 let meta = EventMetadata::new(id, std::any::type_name::<E>());
393 let _ = self
394 .event_publisher
395 .publish(CrudEvent::SoftDeleted { entity, metadata: meta })
396 .await;
397 }
398
399 Ok(deleted)
400 }
401
402 pub async fn restore(&self, id: &str) -> ServiceResult<Option<E>> {
403 let Some(trashed) = self.get_deleted_by_id(id).await? else {
404 return Ok(None);
405 };
406 self.consult(crate::write_guard::WriteKind::Restore, Some(&trashed), None).await?;
407 self.lifecycle.before_restore(&trashed).await?;
408 let restored = self
409 .repository
410 .restore(id)
411 .await
412 .map_err(ServiceError::Repository)?;
413
414 if let Some(ref entity) = restored {
415 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
416 let _ = self
417 .event_publisher
418 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
419 .await;
420 }
421
422 Ok(restored)
423 }
424
425 pub async fn hard_delete(&self, id: &str) -> ServiceResult<bool> {
426 let Some(entity) = self.get_deleted_by_id(id).await? else {
427 return Ok(false);
428 };
429 self.consult(crate::write_guard::WriteKind::HardDelete, Some(&entity), None).await?;
430 self.lifecycle.before_hard_delete(&entity).await?;
431 let deleted = self
432 .repository
433 .hard_delete(id)
434 .await
435 .map_err(ServiceError::Repository)?;
436
437 if deleted {
438 let meta = EventMetadata::new(id, std::any::type_name::<E>());
439 let _ = self
440 .event_publisher
441 .publish(CrudEvent::HardDeleted {
442 entity_id: id.to_string(),
443 metadata: meta,
444 })
445 .await;
446 }
447
448 Ok(deleted)
449 }
450
451 pub async fn list_deleted(
452 &self,
453 page: u32,
454 limit: u32,
455 ) -> ServiceResult<(Vec<E>, u64)> {
456 self.repository
457 .list_deleted(page, limit)
458 .await
459 .map_err(ServiceError::Repository)
460 }
461
462 pub async fn empty_trash(&self) -> ServiceResult<u64> {
463 self.lifecycle.before_empty_trash().await?;
464 self.repository
465 .empty_trash()
466 .await
467 .map_err(ServiceError::Repository)
468 }
469
470 pub async fn count(&self) -> ServiceResult<u64> {
471 self.repository
472 .count()
473 .await
474 .map_err(ServiceError::Repository)
475 }
476
477 pub async fn count_active(&self) -> ServiceResult<u64> {
478 self.repository
479 .count()
480 .await
481 .map_err(ServiceError::Repository)
482 }
483
484 pub async fn find_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
486 self.get_by_id(id).await
487 }
488
489 pub async fn permanent_delete(&self, id: &str) -> ServiceResult<bool> {
491 self.hard_delete(id).await
492 }
493
494 pub async fn count_deleted(&self) -> ServiceResult<u64> {
496 self.repository
497 .count_deleted()
498 .await
499 .map_err(ServiceError::Repository)
500 }
501
502 pub async fn get_deleted_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
504 self.repository
505 .find_by_id_including_deleted(id)
506 .await
507 .map_err(ServiceError::Repository)
508 }
509
510 pub async fn list_deleted_filtered(
512 &self,
513 page: u32,
514 limit: u32,
515 _filters: HashMap<String, String>,
516 ) -> ServiceResult<(Vec<E>, u64)> {
517 self.list_deleted(page, limit).await
518 }
519
520 pub async fn upsert(&self, dto: C) -> ServiceResult<E> {
522 self.create(dto).await
523 }
524
525 pub async fn partial_update(
530 &self,
531 id: &str,
532 fields: HashMap<String, serde_json::Value>,
533 ) -> ServiceResult<Option<E>>
534 where
535 E: serde::Serialize + serde::de::DeserializeOwned,
536 {
537 let existing = self
538 .repository
539 .find_by_id(id)
540 .await
541 .map_err(ServiceError::Repository)?;
542 let Some(before) = existing else {
543 return Ok(None);
544 };
545 let mut patched = merge_patch(&before, fields)?;
546 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&patched)).await?;
547 self.lifecycle.before_update(&mut patched).await?;
548 let saved = self
549 .repository
550 .update(patched)
551 .await
552 .map_err(ServiceError::Repository)?;
553 self.lifecycle.after_update(&saved).await?;
554
555 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
556 let _ = self
557 .event_publisher
558 .publish(CrudEvent::Updated { before, after: saved.clone(), metadata: meta })
559 .await;
560
561 Ok(Some(saved))
562 }
563
564 pub async fn bulk_create(&self, dtos: Vec<C>) -> ServiceResult<Vec<E>>
565 where
566 E: FromCreateDto<C>,
567 {
568 let mut results = Vec::with_capacity(dtos.len());
569 for dto in dtos {
570 let entity = self.create(dto).await?;
571 results.push(entity);
572 }
573 Ok(results)
574 }
575
576 pub async fn bulk_soft_delete(&self, ids: Vec<String>) -> ServiceResult<u64> {
585 check_batch_size(ids.len())?;
586 let ids = dedup_ids(ids);
587 let mut entities = Vec::with_capacity(ids.len());
592 for id in &ids {
593 match self
594 .repository
595 .find_by_id(id)
596 .await
597 .map_err(ServiceError::Repository)?
598 {
599 Some(e) => entities.push(e),
600 None => return Err(ServiceError::Validation(format!("id '{id}' not found"))),
601 }
602 }
603 for entity in &entities {
604 self.consult(crate::write_guard::WriteKind::Delete, Some(entity), None).await?;
605 self.lifecycle.before_delete(entity).await?;
606 }
607 let affected = self
608 .repository
609 .bulk_soft_delete(&ids)
610 .await
611 .map_err(ServiceError::Repository)?;
612 for (id, entity) in ids.iter().zip(entities.into_iter()) {
613 self.lifecycle.after_delete(id).await?;
614 let meta = EventMetadata::new(id, std::any::type_name::<E>());
615 let _ = self
616 .event_publisher
617 .publish(CrudEvent::SoftDeleted { entity, metadata: meta })
618 .await;
619 }
620 Ok(affected)
621 }
622
623 pub async fn bulk_restore(&self, ids: Vec<String>) -> ServiceResult<Vec<E>> {
625 check_batch_size(ids.len())?;
626 let ids = dedup_ids(ids);
627 for id in &ids {
630 if let Some(trashed) = self.get_deleted_by_id(id).await? {
631 self.consult(crate::write_guard::WriteKind::Restore, Some(&trashed), None).await?;
632 self.lifecycle.before_restore(&trashed).await?;
633 }
634 }
635 let restored = self
636 .repository
637 .bulk_restore(&ids)
638 .await
639 .map_err(ServiceError::Repository)?;
640 for entity in &restored {
641 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
642 let _ = self
643 .event_publisher
644 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
645 .await;
646 }
647 Ok(restored)
648 }
649
650 pub async fn restore_all(&self) -> ServiceResult<u64> {
653 self.lifecycle.before_restore_all().await?;
654 let restored = self
655 .repository
656 .restore_all()
657 .await
658 .map_err(ServiceError::Repository)?;
659 for entity in &restored {
660 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
661 let _ = self
662 .event_publisher
663 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
664 .await;
665 }
666 Ok(restored.len() as u64)
667 }
668
669 pub async fn bulk_permanent_delete(&self, ids: Vec<String>) -> ServiceResult<u64> {
671 check_batch_size(ids.len())?;
672 let ids = dedup_ids(ids);
673 for id in &ids {
679 if let Some(trashed) = self.get_deleted_by_id(id).await? {
680 self.consult(crate::write_guard::WriteKind::HardDelete, Some(&trashed), None).await?;
681 self.lifecycle.before_hard_delete(&trashed).await?;
682 }
683 }
684 let affected = self
685 .repository
686 .bulk_hard_delete(&ids)
687 .await
688 .map_err(ServiceError::Repository)?;
689 for id in &ids {
690 let meta = EventMetadata::new(id, std::any::type_name::<E>());
691 let _ = self
692 .event_publisher
693 .publish(CrudEvent::HardDeleted {
694 entity_id: id.to_string(),
695 metadata: meta,
696 })
697 .await;
698 }
699 Ok(affected)
700 }
701
702 pub async fn bulk_update(&self, items: Vec<(String, U)>) -> ServiceResult<Vec<E>> {
704 check_batch_size(items.len())?;
705 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
706 return Err(ServiceError::Validation(format!(
707 "duplicate id '{dup}' in bulk update"
708 )));
709 }
710 let mut befores = Vec::with_capacity(items.len());
711 let mut prepared = Vec::with_capacity(items.len());
712 for (id, dto) in items {
713 let Some(before) = self
714 .repository
715 .find_by_id(&id)
716 .await
717 .map_err(ServiceError::Repository)?
718 else {
719 return Err(ServiceError::Validation(format!("id '{id}' not found")));
720 };
721 let mut updated = before.clone().apply_update(dto)?;
722 refuse_protected_changes::<E>(&as_object(&before)?, &as_object(&updated)?)?;
723 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&updated)).await?;
724 self.lifecycle.before_update(&mut updated).await?;
725 befores.push(before);
726 prepared.push(updated);
727 }
728 let saved = self
729 .repository
730 .bulk_update(prepared)
731 .await
732 .map_err(ServiceError::Repository)?;
733 self.publish_bulk_updates(befores, &saved).await?;
734 Ok(saved)
735 }
736
737 async fn publish_bulk_updates(&self, befores: Vec<E>, saved: &[E]) -> ServiceResult<()> {
741 for (before, after) in befores.into_iter().zip(saved.iter()) {
742 self.lifecycle.after_update(after).await?;
743 let meta = EventMetadata::new(after.entity_id(), std::any::type_name::<E>());
744 let _ = self
745 .event_publisher
746 .publish(CrudEvent::Updated {
747 before,
748 after: after.clone(),
749 metadata: meta,
750 })
751 .await;
752 }
753 Ok(())
754 }
755
756 pub async fn bulk_partial_update(
758 &self,
759 items: Vec<(String, HashMap<String, serde_json::Value>)>,
760 ) -> ServiceResult<Vec<E>>
761 where
762 E: serde::Serialize + serde::de::DeserializeOwned,
763 {
764 check_batch_size(items.len())?;
765 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
766 return Err(ServiceError::Validation(format!(
767 "duplicate id '{dup}' in bulk partial update"
768 )));
769 }
770 let mut befores = Vec::with_capacity(items.len());
771 let mut prepared = Vec::with_capacity(items.len());
772 for (id, fields) in items {
773 let Some(before) = self
774 .repository
775 .find_by_id(&id)
776 .await
777 .map_err(ServiceError::Repository)?
778 else {
779 return Err(ServiceError::Validation(format!("id '{id}' not found")));
780 };
781 let mut patched = merge_patch(&before, fields)?;
782 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&patched)).await?;
783 self.lifecycle.before_update(&mut patched).await?;
784 befores.push(before);
785 prepared.push(patched);
786 }
787 let saved = self
788 .repository
789 .bulk_update(prepared)
790 .await
791 .map_err(ServiceError::Repository)?;
792 self.publish_bulk_updates(befores, &saved).await?;
793 Ok(saved)
794 }
795}
796
797pub const ALWAYS_WRITE_PROTECTED: &[&str] = &["id", "metadata"];
804
805fn as_object(entity: &impl serde::Serialize) -> ServiceResult<serde_json::Map<String, serde_json::Value>> {
806 match serde_json::to_value(entity).map_err(|e| ServiceError::Internal(e.to_string()))? {
807 serde_json::Value::Object(map) => Ok(map),
808 _ => Err(ServiceError::Internal("an entity must serialize to a JSON object".into())),
809 }
810}
811
812fn refuse_protected_changes<E: PersistentEntity>(
816 before: &serde_json::Map<String, serde_json::Value>,
817 after: &serde_json::Map<String, serde_json::Value>,
818) -> ServiceResult<()> {
819 let refused: Vec<crate::violation::Violation> = ALWAYS_WRITE_PROTECTED
820 .iter()
821 .chain(E::write_protected_fields())
822 .filter(|key| before.get(**key) != after.get(**key))
823 .map(|key| {
824 crate::violation::Violation::new(
825 *key,
826 "field_not_writable",
827 format!(
828 "field_not_writable: `{key}` cannot be changed by a generic write; \
829 it changes only through the operation that owns it"
830 ),
831 )
832 })
833 .collect();
834 if refused.is_empty() {
835 Ok(())
836 } else {
837 Err(ServiceError::Violations(refused))
838 }
839}
840
841fn merge_patch<E: PersistentEntity>(
849 before: &E,
850 fields: HashMap<String, serde_json::Value>,
851) -> ServiceResult<E> {
852 let before_obj = as_object(before)?;
853 let mut merged = before_obj.clone();
854 let mut given = Vec::new();
855 for (key, value) in fields {
856 if !value.is_null() {
857 given.push(key.clone());
858 }
859 merged.insert(key, value);
860 }
861 let patched: E = serde_json::from_value(serde_json::Value::Object(merged))
862 .map_err(|e| ServiceError::Validation(e.to_string()))?;
863 let after_obj = as_object(&patched)?;
864
865 let mut unknown: Vec<String> = given.into_iter().filter(|k| !after_obj.contains_key(k)).collect();
866 if !unknown.is_empty() {
867 unknown.sort();
868 return Err(ServiceError::Violations(
869 unknown
870 .into_iter()
871 .map(|k| {
872 let message = format!("unknown_field: `{k}` is not a field of this record");
873 crate::violation::Violation::new(k, "unknown_field", message)
874 })
875 .collect(),
876 ));
877 }
878 refuse_protected_changes::<E>(&before_obj, &after_obj)?;
879 Ok(patched)
880}
881
882pub const MAX_BATCH_SIZE: usize = 1000;
886
887fn check_batch_size(count: usize) -> Result<(), ServiceError> {
889 if count > MAX_BATCH_SIZE {
890 return Err(ServiceError::Validation(format!(
891 "Batch too large: {count} items exceeds the maximum of {MAX_BATCH_SIZE}."
892 )));
893 }
894 Ok(())
895}
896
897fn dedup_ids(ids: Vec<String>) -> Vec<String> {
901 let mut seen = std::collections::HashSet::new();
902 ids.into_iter().filter(|id| seen.insert(id.clone())).collect()
903}
904
905fn first_duplicate_id(ids: impl IntoIterator<Item = String>) -> Option<String> {
908 let mut seen = std::collections::HashSet::new();
909 for id in ids {
910 if !seen.insert(id.clone()) {
911 return Some(id);
912 }
913 }
914 None
915}
916
917#[async_trait::async_trait]
928impl<E, C, U, R> crate::http::CrudService<E, C, U> for GenericCrudService<E, C, U, R>
929where
930 E: PersistentEntity
931 + Clone
932 + serde::Serialize
933 + serde::de::DeserializeOwned
934 + FromCreateDto<C>
935 + ApplyUpdateDto<U>
936 + Send
937 + Sync
938 + 'static,
939 C: Send + Sync + 'static,
940 U: Send + Sync + 'static,
941 R: CrudRepository<E> + Send + Sync + 'static,
942{
943 type Error = ServiceError;
944
945 fn violations_of(err: &ServiceError) -> Option<Vec<crate::violation::Violation>> {
946 match err {
947 ServiceError::Violations(v) => Some(v.clone()),
948 _ => None,
949 }
950 }
951
952 fn entity_name() -> &'static str {
953 std::any::type_name::<E>()
954 }
955
956 async fn fetch_related_json(&self, table: &str, ids: &[String]) -> Vec<serde_json::Value> {
957 self.repository().fetch_related_json(table, ids).await
958 }
959
960 async fn list(
961 &self,
962 page: u32,
963 limit: u32,
964 filters: HashMap<String, String>,
965 ) -> Result<(Vec<E>, u64), ServiceError> {
966 self.list(page, limit, filters).await
967 }
968
969 async fn list_with_info(
970 &self,
971 page: u32,
972 limit: u32,
973 filters: HashMap<String, String>,
974 ) -> Result<(Vec<E>, backbone_orm::repository::PaginationInfo), ServiceError> {
975 self.list_with_info(page, limit, filters).await
976 }
977
978 async fn aggregate(
979 &self,
980 spec: &backbone_orm::repository::AggregateSpec,
981 filters: HashMap<String, String>,
982 ) -> Result<backbone_orm::repository::AggregateResult, ServiceError> {
983 self.aggregate(spec, filters).await
984 }
985
986 fn table_name(&self) -> Option<&str> {
987 self.repository.table_name()
988 }
989
990 async fn create(&self, dto: C) -> Result<E, ServiceError> {
991 self.create(dto).await
992 }
993
994 async fn get_by_id(&self, id: &str) -> Result<Option<E>, ServiceError> {
995 self.get_by_id(id).await
996 }
997
998 async fn update(&self, id: &str, dto: U) -> Result<Option<E>, ServiceError> {
999 self.update(id, dto).await
1000 }
1001
1002 async fn partial_update(
1003 &self,
1004 id: &str,
1005 fields: HashMap<String, serde_json::Value>,
1006 ) -> Result<Option<E>, ServiceError> {
1007 self.partial_update(id, fields).await
1008 }
1009
1010 async fn soft_delete(&self, id: &str) -> Result<bool, ServiceError> {
1011 self.soft_delete(id).await
1012 }
1013
1014 async fn bulk_create(&self, items: Vec<C>) -> Result<Vec<E>, ServiceError> {
1015 self.bulk_create(items).await
1016 }
1017
1018 async fn upsert(&self, dto: C) -> Result<E, ServiceError> {
1019 self.upsert(dto).await
1020 }
1021
1022 async fn list_deleted(&self, page: u32, limit: u32) -> Result<(Vec<E>, u64), ServiceError> {
1023 self.list_deleted(page, limit).await
1024 }
1025
1026 async fn restore(&self, id: &str) -> Result<Option<E>, ServiceError> {
1027 self.restore(id).await
1028 }
1029
1030 async fn empty_trash(&self) -> Result<u64, ServiceError> {
1031 self.empty_trash().await
1032 }
1033
1034 async fn get_deleted_by_id(&self, id: &str) -> Result<Option<E>, ServiceError> {
1035 self.get_deleted_by_id(id).await
1036 }
1037
1038 async fn permanent_delete(&self, id: &str) -> Result<bool, ServiceError> {
1039 self.permanent_delete(id).await
1040 }
1041
1042 async fn list_deleted_filtered(
1043 &self,
1044 page: u32,
1045 limit: u32,
1046 filters: HashMap<String, String>,
1047 ) -> Result<(Vec<E>, u64), ServiceError> {
1048 self.list_deleted_filtered(page, limit, filters).await
1049 }
1050
1051 async fn count_active(&self) -> Result<u64, ServiceError> {
1052 self.count_active().await
1053 }
1054
1055 async fn count_deleted(&self) -> Result<u64, ServiceError> {
1056 self.count_deleted().await
1057 }
1058
1059 async fn bulk_soft_delete(&self, ids: Vec<String>) -> Result<u64, ServiceError> {
1060 self.bulk_soft_delete(ids).await
1061 }
1062
1063 async fn bulk_restore(&self, ids: Vec<String>) -> Result<Vec<E>, ServiceError> {
1064 self.bulk_restore(ids).await
1065 }
1066
1067 async fn bulk_permanent_delete(&self, ids: Vec<String>) -> Result<u64, ServiceError> {
1068 self.bulk_permanent_delete(ids).await
1069 }
1070
1071 async fn restore_all(&self) -> Result<u64, ServiceError> {
1072 self.restore_all().await
1073 }
1074
1075 async fn bulk_update(&self, items: Vec<(String, U)>) -> Result<Vec<E>, ServiceError> {
1076 self.bulk_update(items).await
1077 }
1078
1079 async fn bulk_partial_update(
1080 &self,
1081 items: Vec<(String, HashMap<String, serde_json::Value>)>,
1082 ) -> Result<Vec<E>, ServiceError> {
1083 self.bulk_partial_update(items).await
1084 }
1085}
1086
1087impl<E, C, U, R> Clone for GenericCrudService<E, C, U, R>
1089where
1090 E: PersistentEntity + Clone,
1091 R: CrudRepository<E>,
1092{
1093 fn clone(&self) -> Self {
1094 Self {
1095 repository: self.repository.clone(),
1096 lifecycle: self.lifecycle.clone(),
1097 event_publisher: self.event_publisher.clone(),
1098 guard: self.guard.clone(),
1099 _phantom: PhantomData,
1100 }
1101 }
1102}
1103
1104#[cfg(test)]
1105mod tests {
1106 use super::*;
1107 use chrono::{DateTime, Utc};
1108 use serde::{Deserialize, Serialize};
1109
1110 #[derive(Debug, Clone, Serialize, Deserialize)]
1113 struct Widget {
1114 id: String,
1115 name: String,
1116 #[serde(skip_serializing_if = "Option::is_none")]
1117 note: Option<String>,
1118 secret_hash: String,
1119 #[serde(skip_serializing_if = "Option::is_none")]
1120 deleted_at: Option<DateTime<Utc>>,
1121 created_at: DateTime<Utc>,
1122 updated_at: DateTime<Utc>,
1123 }
1124
1125 impl crate::persistence::traits::PersistentEntity for Widget {
1126 fn write_protected_fields() -> &'static [&'static str] {
1127 &["secret_hash"]
1128 }
1129 fn entity_id(&self) -> String {
1130 self.id.clone()
1131 }
1132 fn set_entity_id(&mut self, id: String) {
1133 self.id = id;
1134 }
1135 fn created_at(&self) -> Option<DateTime<Utc>> {
1136 Some(self.created_at)
1137 }
1138 fn set_created_at(&mut self, ts: DateTime<Utc>) {
1139 self.created_at = ts;
1140 }
1141 fn updated_at(&self) -> Option<DateTime<Utc>> {
1142 Some(self.updated_at)
1143 }
1144 fn set_updated_at(&mut self, ts: DateTime<Utc>) {
1145 self.updated_at = ts;
1146 }
1147 fn deleted_at(&self) -> Option<DateTime<Utc>> {
1148 self.deleted_at
1149 }
1150 fn set_deleted_at(&mut self, ts: Option<DateTime<Utc>>) {
1151 self.deleted_at = ts;
1152 }
1153 }
1154
1155 struct CreateWidgetDto {
1156 name: String,
1157 }
1158
1159 struct UpdateWidgetDto {
1160 name: String,
1161 }
1162
1163 impl FromCreateDto<CreateWidgetDto> for Widget {
1164 fn from_create_dto(dto: CreateWidgetDto) -> ServiceResult<Self> {
1165 Ok(Widget {
1166 id: uuid::Uuid::new_v4().to_string(),
1167 name: dto.name,
1168 note: None,
1169 secret_hash: "h0".into(),
1170 deleted_at: None,
1171 created_at: Utc::now(),
1172 updated_at: Utc::now(),
1173 })
1174 }
1175 }
1176
1177 impl ApplyUpdateDto<UpdateWidgetDto> for Widget {
1178 fn apply_update(mut self, dto: UpdateWidgetDto) -> ServiceResult<Self> {
1179 self.name = dto.name;
1180 self.updated_at = Utc::now();
1181 Ok(self)
1182 }
1183 }
1184
1185 struct InMemoryWidgetRepo {
1187 data: tokio::sync::Mutex<Vec<Widget>>,
1188 }
1189
1190 impl InMemoryWidgetRepo {
1191 fn new() -> Self {
1192 Self {
1193 data: tokio::sync::Mutex::new(Vec::new()),
1194 }
1195 }
1196 }
1197
1198 #[async_trait]
1199 impl crate::persistence::traits::CrudRepository<Widget> for InMemoryWidgetRepo {
1200 async fn create(&self, entity: Widget) -> Result<Widget, RepositoryError> {
1201 self.data.lock().await.push(entity.clone());
1202 Ok(entity)
1203 }
1204
1205 async fn find_by_id(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1206 Ok(self
1207 .data
1208 .lock()
1209 .await
1210 .iter()
1211 .find(|w| w.id == id && w.deleted_at.is_none())
1212 .cloned())
1213 }
1214
1215 async fn update(&self, entity: Widget) -> Result<Widget, RepositoryError> {
1216 let mut data = self.data.lock().await;
1217 data.retain(|w| w.id != entity.id);
1218 data.push(entity.clone());
1219 Ok(entity)
1220 }
1221
1222 async fn soft_delete(&self, id: &str) -> Result<bool, RepositoryError> {
1223 let mut data = self.data.lock().await;
1224 if let Some(w) = data.iter_mut().find(|w| w.id == id) {
1225 w.deleted_at = Some(Utc::now());
1226 return Ok(true);
1227 }
1228 Ok(false)
1229 }
1230
1231 async fn find_by_id_including_deleted(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1232 Ok(self.data.lock().await.iter().find(|w| w.id == id).cloned())
1233 }
1234
1235 async fn restore(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1236 let mut data = self.data.lock().await;
1237 if let Some(w) = data.iter_mut().find(|w| w.id == id) {
1238 w.deleted_at = None;
1239 return Ok(Some(w.clone()));
1240 }
1241 Ok(None)
1242 }
1243
1244 async fn hard_delete(&self, id: &str) -> Result<bool, RepositoryError> {
1245 let mut data = self.data.lock().await;
1246 let before = data.len();
1247 data.retain(|w| w.id != id);
1248 Ok(data.len() < before)
1249 }
1250
1251 async fn list(
1252 &self,
1253 page: u32,
1254 limit: u32,
1255 ) -> Result<(Vec<Widget>, u64), RepositoryError> {
1256 let data = self.data.lock().await;
1257 let active: Vec<_> = data.iter().filter(|w| w.deleted_at.is_none()).cloned().collect();
1258 let total = active.len() as u64;
1259 let offset = ((page.saturating_sub(1)) as usize) * (limit as usize);
1260 let limit = limit as usize;
1261 let page = active.into_iter().skip(offset).take(limit).collect();
1262 Ok((page, total))
1263 }
1264
1265 async fn list_deleted(
1266 &self,
1267 page: u32,
1268 limit: u32,
1269 ) -> Result<(Vec<Widget>, u64), RepositoryError> {
1270 let data = self.data.lock().await;
1271 let deleted: Vec<_> = data.iter().filter(|w| w.deleted_at.is_some()).cloned().collect();
1272 let total = deleted.len() as u64;
1273 let offset = ((page.saturating_sub(1)) as usize) * (limit as usize);
1274 let limit = limit as usize;
1275 let page = deleted.into_iter().skip(offset).take(limit).collect();
1276 Ok((page, total))
1277 }
1278
1279 async fn count(&self) -> Result<u64, RepositoryError> {
1280 let data = self.data.lock().await;
1281 Ok(data.iter().filter(|w| w.deleted_at.is_none()).count() as u64)
1282 }
1283
1284 async fn count_deleted(&self) -> Result<u64, RepositoryError> {
1285 let data = self.data.lock().await;
1286 Ok(data.iter().filter(|w| w.deleted_at.is_some()).count() as u64)
1287 }
1288
1289 async fn bulk_create(&self, entities: Vec<Widget>) -> Result<Vec<Widget>, RepositoryError> {
1290 let mut data = self.data.lock().await;
1291 data.extend(entities.clone());
1292 Ok(entities)
1293 }
1294
1295 async fn empty_trash(&self) -> Result<u64, RepositoryError> {
1296 let mut data = self.data.lock().await;
1297 let before = data.len();
1298 data.retain(|w| w.deleted_at.is_none());
1299 Ok((before - data.len()) as u64)
1300 }
1301 }
1302
1303 #[tokio::test]
1304 async fn create_and_get_roundtrip() {
1305 let repo = Arc::new(InMemoryWidgetRepo::new());
1306 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1307 GenericCrudService::with_repository(repo);
1308
1309 let widget = service
1310 .create(CreateWidgetDto {
1311 name: "sprocket".into(),
1312 })
1313 .await
1314 .unwrap();
1315
1316 let found = service.get_by_id(&widget.id).await.unwrap();
1317 assert!(found.is_some());
1318 assert_eq!(found.unwrap().name, "sprocket");
1319 }
1320
1321 #[tokio::test]
1322 async fn soft_delete_hides_from_list() {
1323 let repo = Arc::new(InMemoryWidgetRepo::new());
1324 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1325 GenericCrudService::with_repository(repo);
1326
1327 let w = service.create(CreateWidgetDto { name: "w".into() }).await.unwrap();
1328 service.soft_delete(&w.id).await.unwrap();
1329
1330 let (items, _) = service.list(1, 20, Default::default()).await.unwrap();
1331 assert!(items.is_empty());
1332
1333 let (deleted, _) = service.list_deleted(1, 20).await.unwrap();
1334 assert_eq!(deleted.len(), 1);
1335 }
1336
1337 #[tokio::test]
1338 async fn update_publishes_event_and_returns_new_entity() {
1339 let repo = Arc::new(InMemoryWidgetRepo::new());
1340 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1341 GenericCrudService::with_repository(repo);
1342
1343 let w = service.create(CreateWidgetDto { name: "old".into() }).await.unwrap();
1344 let updated = service
1345 .update(&w.id, UpdateWidgetDto { name: "new".into() })
1346 .await
1347 .unwrap();
1348 assert_eq!(updated.unwrap().name, "new");
1349 }
1350
1351 fn svc() -> GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> {
1354 GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()))
1355 }
1356
1357 #[tokio::test]
1358 async fn bulk_soft_delete_success() {
1359 let service = svc();
1360 let mut ids = Vec::new();
1361 for n in 0..3 {
1362 ids.push(service.create(CreateWidgetDto { name: format!("w{n}") }).await.unwrap().id);
1363 }
1364 let n = service.bulk_soft_delete(ids).await.unwrap();
1365 assert_eq!(n, 3);
1366 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1367 assert!(active.is_empty());
1368 }
1369
1370 #[tokio::test]
1371 async fn bulk_soft_delete_is_all_or_nothing_on_missing_id() {
1372 let service = svc();
1373 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1374 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1375
1376 let res = service
1378 .bulk_soft_delete(vec![a.id.clone(), b.id.clone(), "does-not-exist".into()])
1379 .await;
1380 assert!(res.is_err());
1381
1382 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1383 assert_eq!(active.len(), 2, "no rows should have been deleted");
1384 }
1385
1386 #[tokio::test]
1387 async fn bulk_restore_and_restore_all() {
1388 let service = svc();
1389 let mut ids = Vec::new();
1390 for n in 0..3 {
1391 ids.push(service.create(CreateWidgetDto { name: format!("w{n}") }).await.unwrap().id);
1392 }
1393 service.bulk_soft_delete(ids.clone()).await.unwrap();
1394
1395 let restored = service.bulk_restore(vec![ids[0].clone(), ids[1].clone()]).await.unwrap();
1397 assert_eq!(restored.len(), 2);
1398
1399 let n = service.restore_all().await.unwrap();
1401 assert_eq!(n, 1);
1402
1403 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1404 assert_eq!(active.len(), 3);
1405 }
1406
1407 #[tokio::test]
1408 async fn bulk_update_is_all_or_nothing_on_missing_id() {
1409 let service = svc();
1410 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1411
1412 let res = service
1413 .bulk_update(vec![
1414 (a.id.clone(), UpdateWidgetDto { name: "A2".into() }),
1415 ("missing".into(), UpdateWidgetDto { name: "X".into() }),
1416 ])
1417 .await;
1418 assert!(res.is_err());
1419
1420 let still = service.get_by_id(&a.id).await.unwrap().unwrap();
1422 assert_eq!(still.name, "a");
1423 }
1424
1425 #[tokio::test]
1426 async fn bulk_partial_update_success() {
1427 let service = svc();
1428 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1429 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1430
1431 let mut patch_a = std::collections::HashMap::new();
1432 patch_a.insert("name".to_string(), serde_json::json!("a2"));
1433 let mut patch_b = std::collections::HashMap::new();
1434 patch_b.insert("name".to_string(), serde_json::json!("b2"));
1435
1436 let saved = service
1437 .bulk_partial_update(vec![(a.id.clone(), patch_a), (b.id.clone(), patch_b)])
1438 .await
1439 .unwrap();
1440 assert_eq!(saved.len(), 2);
1441 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a2");
1442 assert_eq!(service.get_by_id(&b.id).await.unwrap().unwrap().name, "b2");
1443 }
1444
1445 #[tokio::test]
1446 async fn bulk_soft_delete_tolerates_duplicate_ids() {
1447 let service = svc();
1448 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1449 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1450
1451 let n = service
1453 .bulk_soft_delete(vec![a.id.clone(), a.id.clone(), b.id.clone()])
1454 .await
1455 .unwrap();
1456 assert_eq!(n, 2, "two distinct rows soft-deleted");
1457
1458 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1459 assert!(active.is_empty());
1460 }
1461
1462 #[tokio::test]
1463 async fn bulk_update_rejects_duplicate_ids() {
1464 let service = svc();
1465 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1466
1467 let res = service
1469 .bulk_update(vec![
1470 (a.id.clone(), UpdateWidgetDto { name: "first".into() }),
1471 (a.id.clone(), UpdateWidgetDto { name: "second".into() }),
1472 ])
1473 .await;
1474 assert!(res.is_err());
1475 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1476 }
1477
1478 #[tokio::test]
1479 async fn bulk_soft_delete_rejects_oversized_batch() {
1480 let service = svc();
1481 let ids: Vec<String> = (0..MAX_BATCH_SIZE + 1).map(|n| format!("id-{n}")).collect();
1482 assert!(service.bulk_soft_delete(ids).await.is_err());
1484 }
1485
1486 struct RekeyWidgetDto {
1491 name: String,
1492 secret_hash: String,
1493 }
1494
1495 impl ApplyUpdateDto<RekeyWidgetDto> for Widget {
1496 fn apply_update(mut self, dto: RekeyWidgetDto) -> ServiceResult<Self> {
1497 self.name = dto.name;
1498 self.secret_hash = dto.secret_hash;
1499 Ok(self)
1500 }
1501 }
1502
1503 type RekeyService = GenericCrudService<Widget, CreateWidgetDto, RekeyWidgetDto, InMemoryWidgetRepo>;
1504
1505 fn patch(pairs: &[(&str, serde_json::Value)]) -> HashMap<String, serde_json::Value> {
1506 pairs.iter().map(|(k, v)| (k.to_string(), v.clone())).collect()
1507 }
1508
1509 fn validation_message(res: ServiceResult<impl std::fmt::Debug>) -> String {
1512 match res {
1513 Err(ServiceError::Validation(m)) => m,
1514 Err(e @ ServiceError::Violations(_)) => e.to_string(),
1515 other => panic!("expected a validation error, got {other:?}"),
1516 }
1517 }
1518
1519 #[tokio::test]
1520 async fn patch_rejects_a_field_the_entity_does_not_have() {
1521 let service = svc();
1522 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1523
1524 let msg = validation_message(service.partial_update(&w.id, patch(&[("nmae", serde_json::json!("b"))])).await);
1525 assert!(msg.contains("unknown_field") && msg.contains("nmae"), "{msg}");
1526 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().name, "a");
1527 }
1528
1529 #[tokio::test]
1530 async fn patch_rejects_changing_the_id() {
1531 let service = svc();
1532 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1533
1534 let msg = validation_message(service.partial_update(&w.id, patch(&[("id", serde_json::json!("other"))])).await);
1535 assert!(msg.contains("field_not_writable") && msg.contains("id"), "{msg}");
1536 assert!(service.get_by_id(&w.id).await.unwrap().is_some());
1537 }
1538
1539 #[tokio::test]
1540 async fn patch_rejects_changing_a_field_the_entity_protects() {
1541 let service = svc();
1542 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1543
1544 let msg = validation_message(
1545 service.partial_update(&w.id, patch(&[("secret_hash", serde_json::json!("h1"))])).await,
1546 );
1547 assert!(msg.contains("field_not_writable") && msg.contains("secret_hash"), "{msg}");
1548 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().secret_hash, "h0");
1549 }
1550
1551 #[tokio::test]
1552 async fn patch_accepts_a_protected_field_echoed_unchanged() {
1553 let service = svc();
1554 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1555
1556 let saved = service
1557 .partial_update(&w.id, patch(&[("name", serde_json::json!("b")), ("secret_hash", serde_json::json!("h0"))]))
1558 .await
1559 .unwrap()
1560 .unwrap();
1561 assert_eq!(saved.name, "b");
1562 }
1563
1564 #[tokio::test]
1565 async fn patch_sets_and_clears_an_optional_field_that_is_not_serialized_when_empty() {
1566 let service = svc();
1567 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1568
1569 let set = service.partial_update(&w.id, patch(&[("note", serde_json::json!("hi"))])).await.unwrap().unwrap();
1570 assert_eq!(set.note.as_deref(), Some("hi"));
1571 let cleared = service.partial_update(&w.id, patch(&[("note", serde_json::Value::Null)])).await.unwrap().unwrap();
1572 assert_eq!(cleared.note, None);
1573 }
1574
1575 #[tokio::test]
1576 async fn put_rejects_changing_a_field_the_entity_protects() {
1577 let service: RekeyService = GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()));
1578 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1579
1580 let msg = validation_message(
1581 service.update(&w.id, RekeyWidgetDto { name: "b".into(), secret_hash: "h1".into() }).await,
1582 );
1583 assert!(msg.contains("field_not_writable") && msg.contains("secret_hash"), "{msg}");
1584
1585 let saved = service
1587 .update(&w.id, RekeyWidgetDto { name: "b".into(), secret_hash: "h0".into() })
1588 .await
1589 .unwrap()
1590 .unwrap();
1591 assert_eq!(saved.name, "b");
1592 }
1593
1594 #[tokio::test]
1595 async fn bulk_update_rejects_a_protected_change_before_writing_any_row() {
1596 let service: RekeyService = GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()));
1597 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1598 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1599
1600 let res = service
1601 .bulk_update(vec![
1602 (a.id.clone(), RekeyWidgetDto { name: "a2".into(), secret_hash: "h0".into() }),
1603 (b.id.clone(), RekeyWidgetDto { name: "b2".into(), secret_hash: "h1".into() }),
1604 ])
1605 .await;
1606 assert!(validation_message(res).contains("field_not_writable"));
1607 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1608 }
1609
1610 #[tokio::test]
1611 async fn bulk_partial_update_rejects_a_protected_change_before_writing_any_row() {
1612 let service = svc();
1613 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1614 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1615
1616 let res = service
1617 .bulk_partial_update(vec![
1618 (a.id.clone(), patch(&[("name", serde_json::json!("a2"))])),
1619 (b.id.clone(), patch(&[("secret_hash", serde_json::json!("h1"))])),
1620 ])
1621 .await;
1622 assert!(validation_message(res).contains("field_not_writable"));
1623 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1624 }
1625
1626 #[tokio::test]
1627 async fn bulk_partial_update_rejects_an_unknown_field() {
1628 let service = svc();
1629 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1630
1631 let res = service.bulk_partial_update(vec![(a.id.clone(), patch(&[("colour", serde_json::json!("red"))]))]).await;
1632 assert!(validation_message(res).contains("unknown_field"));
1633 }
1634
1635 #[derive(Default)]
1638 struct RecordingPublisher {
1639 kinds: tokio::sync::Mutex<Vec<&'static str>>,
1640 }
1641
1642 #[async_trait]
1643 impl CrudEventPublisher<Widget> for RecordingPublisher {
1644 async fn publish(&self, event: CrudEvent<Widget>) -> Result<(), backbone_messaging::EventError> {
1645 let kind = match event {
1646 CrudEvent::Updated { .. } => "updated",
1647 CrudEvent::Patched { .. } => "patched",
1648 _ => "other",
1649 };
1650 self.kinds.lock().await.push(kind);
1651 Ok(())
1652 }
1653 }
1654
1655 #[tokio::test]
1656 async fn patch_publishes_an_updated_event() {
1657 let publisher = Arc::new(RecordingPublisher::default());
1658 let service = svc().with_event_publisher(publisher.clone());
1659 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1660 publisher.kinds.lock().await.clear();
1661
1662 service.partial_update(&w.id, patch(&[("name", serde_json::json!("b"))])).await.unwrap();
1663 assert_eq!(*publisher.kinds.lock().await, vec!["updated"]);
1664 }
1665
1666 struct RefuseTrash;
1669
1670 #[async_trait]
1671 impl ServiceLifecycle<Widget> for RefuseTrash {
1672 async fn before_restore(&self, _entity: &Widget) -> ServiceResult<()> {
1673 Err(ServiceError::Validation("restore refused".into()))
1674 }
1675 async fn before_hard_delete(&self, _entity: &Widget) -> ServiceResult<()> {
1676 Err(ServiceError::Validation("hard delete refused".into()))
1677 }
1678 async fn before_restore_all(&self) -> ServiceResult<()> {
1679 Err(ServiceError::Validation("restore all refused".into()))
1680 }
1681 async fn before_empty_trash(&self) -> ServiceResult<()> {
1682 Err(ServiceError::Validation("empty trash refused".into()))
1683 }
1684 }
1685
1686 #[tokio::test]
1687 async fn every_trash_operation_consults_the_lifecycle() {
1688 let repo = Arc::new(InMemoryWidgetRepo::new());
1689 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1690 GenericCrudService::with_lifecycle(repo, Arc::new(RefuseTrash));
1691 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1692 service.soft_delete(&w.id).await.unwrap();
1693
1694 assert!(service.restore(&w.id).await.is_err(), "restore");
1695 assert!(service.bulk_restore(vec![w.id.clone()]).await.is_err(), "bulk restore");
1696 assert!(service.restore_all().await.is_err(), "restore all");
1697 assert!(service.hard_delete(&w.id).await.is_err(), "hard delete");
1698 assert!(service.bulk_permanent_delete(vec![w.id.clone()]).await.is_err(), "bulk permanent delete");
1699 assert!(service.empty_trash().await.is_err(), "empty trash");
1700
1701 let still = service.get_deleted_by_id(&w.id).await.unwrap().unwrap();
1703 assert!(still.deleted_at.is_some());
1704 }
1705
1706 #[tokio::test]
1709 async fn a_refused_field_is_a_violation_naming_its_path_and_code() {
1710 let service = svc();
1711 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1712
1713 let err = service
1714 .partial_update(&w.id, patch(&[("secret_hash", serde_json::json!("h1")), ("colour", serde_json::json!("red"))]))
1715 .await
1716 .unwrap_err();
1717 let ServiceError::Violations(v) = &err else { panic!("expected violations, got {err:?}") };
1718 let pairs: Vec<(&str, &str)> = v.iter().map(|x| (x.path.as_str(), x.code.as_str())).collect();
1719 assert_eq!(pairs, vec![("colour", "unknown_field")]);
1720 assert!(err.to_string().starts_with("validation failed: unknown_field: `colour`"), "{err}");
1722
1723 let err = service.partial_update(&w.id, patch(&[("secret_hash", serde_json::json!("h1"))])).await.unwrap_err();
1724 let ServiceError::Violations(v) = &err else { panic!("expected violations, got {err:?}") };
1725 assert_eq!((v[0].path.as_str(), v[0].code.as_str()), ("secret_hash", "field_not_writable"));
1726 }
1727
1728 #[derive(Default)]
1731 struct NameGuard {
1732 seen: std::sync::Mutex<Vec<crate::write_guard::WriteKind>>,
1733 }
1734
1735 #[async_trait]
1736 impl crate::write_guard::WriteGuard<Widget> for NameGuard {
1737 async fn check(&self, ctx: &crate::write_guard::WriteCtx<'_, Widget>) -> crate::write_guard::GuardOutcome {
1738 self.seen.lock().unwrap().push(ctx.kind);
1739 let row = ctx.after.or(ctx.before).expect("a write concerns a row");
1740 let mut out = crate::write_guard::GuardOutcome::default();
1741 match row.name.as_str() {
1742 "forbidden" => out.refuse.push(crate::violation::Violation::new("name", "name_forbidden", "this name is refused")),
1743 "watched" => out.shadow.push(crate::violation::Violation::new("name", "name_watched", "would be refused")),
1744 _ => {}
1745 }
1746 out
1747 }
1748 }
1749
1750 fn guarded(guard: Arc<NameGuard>) -> GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> {
1751 svc().with_guard(guard)
1752 }
1753
1754 #[tokio::test]
1755 async fn a_guard_refusal_stops_create_update_and_patch_with_nothing_written() {
1756 let guard = Arc::new(NameGuard::default());
1757 let service = guarded(guard.clone());
1758
1759 let err = service.create(CreateWidgetDto { name: "forbidden".into() }).await.unwrap_err();
1760 assert!(matches!(&err, ServiceError::Violations(v) if v[0].code == "name_forbidden"), "{err:?}");
1761 assert!(service.list(1, 20, Default::default()).await.unwrap().0.is_empty(), "nothing created");
1762
1763 let w = service.create(CreateWidgetDto { name: "ok".into() }).await.unwrap();
1764 assert!(service.update(&w.id, UpdateWidgetDto { name: "forbidden".into() }).await.is_err());
1765 assert!(service.partial_update(&w.id, patch(&[("name", serde_json::json!("forbidden"))])).await.is_err());
1766 assert!(service
1767 .bulk_partial_update(vec![(w.id.clone(), patch(&[("name", serde_json::json!("forbidden"))]))])
1768 .await
1769 .is_err());
1770 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().name, "ok");
1771 }
1772
1773 #[tokio::test]
1774 async fn a_shadowed_rule_lets_the_write_through() {
1775 let service = guarded(Arc::new(NameGuard::default()));
1776 let w = service.create(CreateWidgetDto { name: "watched".into() }).await.unwrap();
1777 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().name, "watched");
1778 }
1779
1780 #[tokio::test]
1781 async fn the_guard_is_asked_about_every_write_kind() {
1782 use crate::write_guard::WriteKind::*;
1783 let guard = Arc::new(NameGuard::default());
1784 let service = guarded(guard.clone());
1785 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1786 service.update(&w.id, UpdateWidgetDto { name: "b".into() }).await.unwrap();
1787 service.partial_update(&w.id, patch(&[("name", serde_json::json!("c"))])).await.unwrap();
1788 service.soft_delete(&w.id).await.unwrap();
1789 service.restore(&w.id).await.unwrap();
1790 service.soft_delete(&w.id).await.unwrap();
1791 service.hard_delete(&w.id).await.unwrap();
1792 assert_eq!(*guard.seen.lock().unwrap(), vec![Create, Update, Update, Delete, Restore, Delete, HardDelete]);
1793 }
1794
1795 #[test]
1796 fn the_generic_service_hands_its_violations_to_the_http_layer() {
1797 use crate::http::CrudService;
1798 type S = GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo>;
1799 let v = vec![crate::violation::Violation::new("name", "x", "y")];
1800 assert_eq!(<S as CrudService<Widget, CreateWidgetDto, UpdateWidgetDto>>::violations_of(&ServiceError::Violations(v.clone())), Some(v));
1801 assert_eq!(<S as CrudService<Widget, CreateWidgetDto, UpdateWidgetDto>>::violations_of(&ServiceError::NotFound), None);
1802 }
1803}