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 self.consult_as(kind, before, after, true).await
233 }
234
235 async fn consult_quietly(
239 &self,
240 kind: crate::write_guard::WriteKind,
241 before: Option<&E>,
242 after: Option<&E>,
243 ) -> ServiceResult<()> {
244 self.consult_as(kind, before, after, false).await
245 }
246
247 async fn consult_as(
248 &self,
249 kind: crate::write_guard::WriteKind,
250 before: Option<&E>,
251 after: Option<&E>,
252 log_shadow: bool,
253 ) -> ServiceResult<()> {
254 let ctx = crate::write_guard::WriteCtx { kind, before, after };
255 let outcome = self.guard.check(&ctx).await;
256 for v in outcome.shadow.iter().filter(|_| log_shadow) {
257 tracing::warn!(
258 target: "backbone::write_guard",
259 entity = std::any::type_name::<E>(),
260 kind = ?kind,
261 code = %v.code,
262 path = %v.path,
263 "a rule running in shadow would refuse this write: {}",
264 v.message
265 );
266 }
267 if outcome.refuse.is_empty() {
268 Ok(())
269 } else {
270 Err(ServiceError::Violations(outcome.refuse))
271 }
272 }
273
274 pub fn with_lifecycle(repository: Arc<R>, lifecycle: Arc<dyn ServiceLifecycle<E>>) -> Self {
276 Self::new(repository, lifecycle, NoOpCrudEventPublisher::arc())
277 }
278
279 pub fn with_repository(repository: Arc<R>) -> Self {
281 Self::new(
282 repository,
283 Arc::new(NoOpLifecycle::new()),
284 NoOpCrudEventPublisher::arc(),
285 )
286 }
287
288 pub fn with_event_publisher(mut self, publisher: Arc<dyn CrudEventPublisher<E>>) -> Self {
290 self.event_publisher = publisher;
291 self
292 }
293
294 pub fn repository(&self) -> &R {
296 &self.repository
297 }
298
299 pub async fn list(
302 &self,
303 page: u32,
304 limit: u32,
305 filters: HashMap<String, String>,
306 ) -> ServiceResult<(Vec<E>, u64)> {
307 self.repository
308 .list_filtered(page, limit, filters)
309 .await
310 .map_err(ServiceError::Repository)
311 }
312
313 pub async fn list_with_info(
316 &self,
317 page: u32,
318 limit: u32,
319 filters: HashMap<String, String>,
320 ) -> ServiceResult<(Vec<E>, backbone_orm::repository::PaginationInfo)> {
321 self.repository
322 .list_filtered_with_info(page, limit, filters)
323 .await
324 .map_err(ServiceError::Repository)
325 }
326
327 pub async fn aggregate(
329 &self,
330 spec: &backbone_orm::repository::AggregateSpec,
331 filters: HashMap<String, String>,
332 ) -> ServiceResult<backbone_orm::repository::AggregateResult> {
333 self.repository
334 .aggregate_filtered(spec, filters)
335 .await
336 .map_err(ServiceError::Repository)
337 }
338
339 pub async fn create(&self, dto: C) -> ServiceResult<E> {
340 let mut entity = E::from_create_dto(dto)?;
341 self.consult(crate::write_guard::WriteKind::Create, None, Some(&entity)).await?;
342 self.lifecycle.before_create(&mut entity).await?;
343 let saved = self
344 .repository
345 .create(entity)
346 .await
347 .map_err(ServiceError::Repository)?;
348 self.lifecycle.after_create(&saved).await?;
349
350 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
351 let _ = self
352 .event_publisher
353 .publish(CrudEvent::Created { entity: saved.clone(), metadata: meta })
354 .await;
355
356 Ok(saved)
357 }
358
359 pub async fn get_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
360 self.repository
361 .find_by_id(id)
362 .await
363 .map_err(ServiceError::Repository)
364 }
365
366 pub async fn update(&self, id: &str, dto: U) -> ServiceResult<Option<E>> {
367 let existing = self
368 .repository
369 .find_by_id(id)
370 .await
371 .map_err(ServiceError::Repository)?;
372 let Some(before) = existing else {
373 return Ok(None);
374 };
375 let mut updated = before.clone().apply_update(dto)?;
376 refuse_protected_changes::<E>(&as_object(&before)?, &as_object(&updated)?)?;
377 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&updated)).await?;
378 self.lifecycle.before_update(&mut updated).await?;
379 let saved = self
380 .repository
381 .update(updated)
382 .await
383 .map_err(ServiceError::Repository)?;
384 self.lifecycle.after_update(&saved).await?;
385
386 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
387 let _ = self
388 .event_publisher
389 .publish(CrudEvent::Updated { before, after: saved.clone(), metadata: meta })
390 .await;
391
392 Ok(Some(saved))
393 }
394
395 pub async fn soft_delete(&self, id: &str) -> ServiceResult<bool> {
396 let existing = self
397 .repository
398 .find_by_id(id)
399 .await
400 .map_err(ServiceError::Repository)?;
401 let Some(entity) = existing else {
402 return Ok(false);
403 };
404 self.consult(crate::write_guard::WriteKind::Delete, Some(&entity), None).await?;
405 self.lifecycle.before_delete(&entity).await?;
406 let deleted = self
407 .repository
408 .soft_delete(id)
409 .await
410 .map_err(ServiceError::Repository)?;
411 self.lifecycle.after_delete(id).await?;
412
413 if deleted {
414 let meta = EventMetadata::new(id, std::any::type_name::<E>());
415 let _ = self
416 .event_publisher
417 .publish(CrudEvent::SoftDeleted { entity, metadata: meta })
418 .await;
419 }
420
421 Ok(deleted)
422 }
423
424 pub async fn restore(&self, id: &str) -> ServiceResult<Option<E>> {
425 let Some(trashed) = self.get_deleted_by_id(id).await? else {
426 return Ok(None);
427 };
428 self.consult(crate::write_guard::WriteKind::Restore, Some(&trashed), None).await?;
429 self.lifecycle.before_restore(&trashed).await?;
430 let restored = self
431 .repository
432 .restore(id)
433 .await
434 .map_err(ServiceError::Repository)?;
435
436 if let Some(ref entity) = restored {
437 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
438 let _ = self
439 .event_publisher
440 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
441 .await;
442 }
443
444 Ok(restored)
445 }
446
447 pub async fn hard_delete(&self, id: &str) -> ServiceResult<bool> {
448 let Some(entity) = self.get_deleted_by_id(id).await? else {
449 return Ok(false);
450 };
451 self.consult(crate::write_guard::WriteKind::HardDelete, Some(&entity), None).await?;
452 self.lifecycle.before_hard_delete(&entity).await?;
453 let deleted = self
454 .repository
455 .hard_delete(id)
456 .await
457 .map_err(ServiceError::Repository)?;
458
459 if deleted {
460 let meta = EventMetadata::new(id, std::any::type_name::<E>());
461 let _ = self
462 .event_publisher
463 .publish(CrudEvent::HardDeleted {
464 entity_id: id.to_string(),
465 metadata: meta,
466 })
467 .await;
468 }
469
470 Ok(deleted)
471 }
472
473 pub async fn list_deleted(
474 &self,
475 page: u32,
476 limit: u32,
477 ) -> ServiceResult<(Vec<E>, u64)> {
478 self.repository
479 .list_deleted(page, limit)
480 .await
481 .map_err(ServiceError::Repository)
482 }
483
484 pub async fn empty_trash(&self) -> ServiceResult<u64> {
485 self.lifecycle.before_empty_trash().await?;
486 self.repository
487 .empty_trash()
488 .await
489 .map_err(ServiceError::Repository)
490 }
491
492 pub async fn count(&self) -> ServiceResult<u64> {
493 self.repository
494 .count()
495 .await
496 .map_err(ServiceError::Repository)
497 }
498
499 pub async fn count_active(&self) -> ServiceResult<u64> {
500 self.repository
501 .count()
502 .await
503 .map_err(ServiceError::Repository)
504 }
505
506 pub async fn find_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
508 self.get_by_id(id).await
509 }
510
511 pub async fn permanent_delete(&self, id: &str) -> ServiceResult<bool> {
513 self.hard_delete(id).await
514 }
515
516 pub async fn count_deleted(&self) -> ServiceResult<u64> {
518 self.repository
519 .count_deleted()
520 .await
521 .map_err(ServiceError::Repository)
522 }
523
524 pub async fn get_deleted_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
526 self.repository
527 .find_by_id_including_deleted(id)
528 .await
529 .map_err(ServiceError::Repository)
530 }
531
532 pub async fn list_deleted_filtered(
534 &self,
535 page: u32,
536 limit: u32,
537 _filters: HashMap<String, String>,
538 ) -> ServiceResult<(Vec<E>, u64)> {
539 self.list_deleted(page, limit).await
540 }
541
542 pub async fn upsert(&self, dto: C) -> ServiceResult<E> {
544 self.create(dto).await
545 }
546
547 pub async fn partial_update(
552 &self,
553 id: &str,
554 fields: HashMap<String, serde_json::Value>,
555 ) -> ServiceResult<Option<E>>
556 where
557 E: serde::Serialize + serde::de::DeserializeOwned,
558 {
559 let existing = self
560 .repository
561 .find_by_id(id)
562 .await
563 .map_err(ServiceError::Repository)?;
564 let Some(before) = existing else {
565 return Ok(None);
566 };
567 let mut patched = merge_patch(&before, fields)?;
568 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&patched)).await?;
569 self.lifecycle.before_update(&mut patched).await?;
570 let saved = self
571 .repository
572 .update(patched)
573 .await
574 .map_err(ServiceError::Repository)?;
575 self.lifecycle.after_update(&saved).await?;
576
577 let meta = EventMetadata::new(saved.entity_id(), std::any::type_name::<E>());
578 let _ = self
579 .event_publisher
580 .publish(CrudEvent::Updated { before, after: saved.clone(), metadata: meta })
581 .await;
582
583 Ok(Some(saved))
584 }
585
586 pub async fn bulk_create(&self, dtos: Vec<C>) -> ServiceResult<Vec<E>>
587 where
588 E: FromCreateDto<C>,
589 {
590 let mut results = Vec::with_capacity(dtos.len());
591 for dto in dtos {
592 let entity = self.create(dto).await?;
593 results.push(entity);
594 }
595 Ok(results)
596 }
597
598 pub async fn bulk_soft_delete(&self, ids: Vec<String>) -> ServiceResult<u64> {
607 check_batch_size(ids.len())?;
608 let ids = dedup_ids(ids);
609 let mut entities = Vec::with_capacity(ids.len());
614 for id in &ids {
615 match self
616 .repository
617 .find_by_id(id)
618 .await
619 .map_err(ServiceError::Repository)?
620 {
621 Some(e) => entities.push(e),
622 None => return Err(ServiceError::Validation(format!("id '{id}' not found"))),
623 }
624 }
625 for entity in &entities {
626 self.consult(crate::write_guard::WriteKind::Delete, Some(entity), None).await?;
627 self.lifecycle.before_delete(entity).await?;
628 }
629 let affected = self
630 .repository
631 .bulk_soft_delete(&ids)
632 .await
633 .map_err(ServiceError::Repository)?;
634 for (id, entity) in ids.iter().zip(entities.into_iter()) {
635 self.lifecycle.after_delete(id).await?;
636 let meta = EventMetadata::new(id, std::any::type_name::<E>());
637 let _ = self
638 .event_publisher
639 .publish(CrudEvent::SoftDeleted { entity, metadata: meta })
640 .await;
641 }
642 Ok(affected)
643 }
644
645 pub async fn bulk_restore(&self, ids: Vec<String>) -> ServiceResult<Vec<E>> {
647 check_batch_size(ids.len())?;
648 let ids = dedup_ids(ids);
649 for id in &ids {
652 if let Some(trashed) = self.get_deleted_by_id(id).await? {
653 self.consult(crate::write_guard::WriteKind::Restore, Some(&trashed), None).await?;
654 self.lifecycle.before_restore(&trashed).await?;
655 }
656 }
657 let restored = self
658 .repository
659 .bulk_restore(&ids)
660 .await
661 .map_err(ServiceError::Repository)?;
662 for entity in &restored {
663 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
664 let _ = self
665 .event_publisher
666 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
667 .await;
668 }
669 Ok(restored)
670 }
671
672 pub async fn restore_all(&self) -> ServiceResult<u64> {
675 self.lifecycle.before_restore_all().await?;
676 let restored = self
677 .repository
678 .restore_all()
679 .await
680 .map_err(ServiceError::Repository)?;
681 for entity in &restored {
682 let meta = EventMetadata::new(entity.entity_id(), std::any::type_name::<E>());
683 let _ = self
684 .event_publisher
685 .publish(CrudEvent::Restored { entity: entity.clone(), metadata: meta })
686 .await;
687 }
688 Ok(restored.len() as u64)
689 }
690
691 pub async fn bulk_permanent_delete(&self, ids: Vec<String>) -> ServiceResult<u64> {
693 check_batch_size(ids.len())?;
694 let ids = dedup_ids(ids);
695 for id in &ids {
701 if let Some(trashed) = self.get_deleted_by_id(id).await? {
702 self.consult(crate::write_guard::WriteKind::HardDelete, Some(&trashed), None).await?;
703 self.lifecycle.before_hard_delete(&trashed).await?;
704 }
705 }
706 let affected = self
707 .repository
708 .bulk_hard_delete(&ids)
709 .await
710 .map_err(ServiceError::Repository)?;
711 for id in &ids {
712 let meta = EventMetadata::new(id, std::any::type_name::<E>());
713 let _ = self
714 .event_publisher
715 .publish(CrudEvent::HardDeleted {
716 entity_id: id.to_string(),
717 metadata: meta,
718 })
719 .await;
720 }
721 Ok(affected)
722 }
723
724 pub async fn bulk_update(&self, items: Vec<(String, U)>) -> ServiceResult<Vec<E>> {
726 check_batch_size(items.len())?;
727 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
728 return Err(ServiceError::Validation(format!(
729 "duplicate id '{dup}' in bulk update"
730 )));
731 }
732 let mut befores = Vec::with_capacity(items.len());
733 let mut prepared = Vec::with_capacity(items.len());
734 for (id, dto) in items {
735 let Some(before) = self
736 .repository
737 .find_by_id(&id)
738 .await
739 .map_err(ServiceError::Repository)?
740 else {
741 return Err(ServiceError::Validation(format!("id '{id}' not found")));
742 };
743 let mut updated = before.clone().apply_update(dto)?;
744 refuse_protected_changes::<E>(&as_object(&before)?, &as_object(&updated)?)?;
745 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&updated)).await?;
746 self.lifecycle.before_update(&mut updated).await?;
747 befores.push(before);
748 prepared.push(updated);
749 }
750 let saved = self
751 .repository
752 .bulk_update(prepared)
753 .await
754 .map_err(ServiceError::Repository)?;
755 self.publish_bulk_updates(befores, &saved).await?;
756 Ok(saved)
757 }
758
759 async fn publish_bulk_updates(&self, befores: Vec<E>, saved: &[E]) -> ServiceResult<()> {
763 for (before, after) in befores.into_iter().zip(saved.iter()) {
764 self.lifecycle.after_update(after).await?;
765 let meta = EventMetadata::new(after.entity_id(), std::any::type_name::<E>());
766 let _ = self
767 .event_publisher
768 .publish(CrudEvent::Updated {
769 before,
770 after: after.clone(),
771 metadata: meta,
772 })
773 .await;
774 }
775 Ok(())
776 }
777
778 pub async fn bulk_partial_update(
780 &self,
781 items: Vec<(String, HashMap<String, serde_json::Value>)>,
782 ) -> ServiceResult<Vec<E>>
783 where
784 E: serde::Serialize + serde::de::DeserializeOwned,
785 {
786 check_batch_size(items.len())?;
787 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
788 return Err(ServiceError::Validation(format!(
789 "duplicate id '{dup}' in bulk partial update"
790 )));
791 }
792 let mut befores = Vec::with_capacity(items.len());
793 let mut prepared = Vec::with_capacity(items.len());
794 for (id, fields) in items {
795 let Some(before) = self
796 .repository
797 .find_by_id(&id)
798 .await
799 .map_err(ServiceError::Repository)?
800 else {
801 return Err(ServiceError::Validation(format!("id '{id}' not found")));
802 };
803 let mut patched = merge_patch(&before, fields)?;
804 self.consult(crate::write_guard::WriteKind::Update, Some(&before), Some(&patched)).await?;
805 self.lifecycle.before_update(&mut patched).await?;
806 befores.push(before);
807 prepared.push(patched);
808 }
809 let saved = self
810 .repository
811 .bulk_update(prepared)
812 .await
813 .map_err(ServiceError::Repository)?;
814 self.publish_bulk_updates(befores, &saved).await?;
815 Ok(saved)
816 }
817
818 pub async fn preview_bulk_partial_update(
827 &self,
828 items: Vec<(String, HashMap<String, serde_json::Value>)>,
829 ) -> ServiceResult<Vec<crate::http::BulkPreviewRow<E, ServiceError>>>
830 where
831 E: serde::Serialize + serde::de::DeserializeOwned,
832 {
833 check_batch_size(items.len())?;
834 if let Some(dup) = first_duplicate_id(items.iter().map(|(id, _)| id.clone())) {
835 return Err(ServiceError::Validation(format!(
836 "duplicate id '{dup}' in bulk partial update"
837 )));
838 }
839 let mut rows = Vec::with_capacity(items.len());
840 for (id, fields) in items {
841 let outcome = self.preview_one(&id, fields).await;
842 rows.push(crate::http::BulkPreviewRow { id, outcome });
843 }
844 Ok(rows)
845 }
846
847 async fn preview_one(
848 &self,
849 id: &str,
850 fields: HashMap<String, serde_json::Value>,
851 ) -> ServiceResult<(E, E)>
852 where
853 E: serde::Serialize + serde::de::DeserializeOwned,
854 {
855 let Some(before) = self
856 .repository
857 .find_by_id(id)
858 .await
859 .map_err(ServiceError::Repository)?
860 else {
861 return Err(ServiceError::Validation(format!("id '{id}' not found")));
862 };
863 let mut patched = merge_patch(&before, fields)?;
864 self.consult_quietly(crate::write_guard::WriteKind::Update, Some(&before), Some(&patched))
865 .await?;
866 self.lifecycle.before_update(&mut patched).await?;
867 Ok((before, patched))
868 }
869}
870
871pub const ALWAYS_WRITE_PROTECTED: &[&str] = &["id", "metadata"];
878
879fn as_object(entity: &impl serde::Serialize) -> ServiceResult<serde_json::Map<String, serde_json::Value>> {
880 match serde_json::to_value(entity).map_err(|e| ServiceError::Internal(e.to_string()))? {
881 serde_json::Value::Object(map) => Ok(map),
882 _ => Err(ServiceError::Internal("an entity must serialize to a JSON object".into())),
883 }
884}
885
886fn refuse_protected_changes<E: PersistentEntity>(
890 before: &serde_json::Map<String, serde_json::Value>,
891 after: &serde_json::Map<String, serde_json::Value>,
892) -> ServiceResult<()> {
893 let refused: Vec<crate::violation::Violation> = ALWAYS_WRITE_PROTECTED
894 .iter()
895 .chain(E::write_protected_fields())
896 .filter(|key| before.get(**key) != after.get(**key))
897 .map(|key| {
898 crate::violation::Violation::new(
899 *key,
900 "field_not_writable",
901 format!(
902 "field_not_writable: `{key}` cannot be changed by a generic write; \
903 it changes only through the operation that owns it"
904 ),
905 )
906 })
907 .collect();
908 if refused.is_empty() {
909 Ok(())
910 } else {
911 Err(ServiceError::Violations(refused))
912 }
913}
914
915fn merge_patch<E: PersistentEntity>(
923 before: &E,
924 fields: HashMap<String, serde_json::Value>,
925) -> ServiceResult<E> {
926 let before_obj = as_object(before)?;
927 let mut merged = before_obj.clone();
928 let mut given = Vec::new();
929 for (key, value) in fields {
930 if !value.is_null() {
931 given.push(key.clone());
932 }
933 merged.insert(key, value);
934 }
935 let patched: E = serde_json::from_value(serde_json::Value::Object(merged))
936 .map_err(|e| ServiceError::Validation(e.to_string()))?;
937 let after_obj = as_object(&patched)?;
938
939 let mut unknown: Vec<String> = given.into_iter().filter(|k| !after_obj.contains_key(k)).collect();
940 if !unknown.is_empty() {
941 unknown.sort();
942 return Err(ServiceError::Violations(
943 unknown
944 .into_iter()
945 .map(|k| {
946 let message = format!("unknown_field: `{k}` is not a field of this record");
947 crate::violation::Violation::new(k, "unknown_field", message)
948 })
949 .collect(),
950 ));
951 }
952 refuse_protected_changes::<E>(&before_obj, &after_obj)?;
953 Ok(patched)
954}
955
956pub const MAX_BATCH_SIZE: usize = 1000;
960
961fn check_batch_size(count: usize) -> Result<(), ServiceError> {
963 if count > MAX_BATCH_SIZE {
964 return Err(ServiceError::Validation(format!(
965 "Batch too large: {count} items exceeds the maximum of {MAX_BATCH_SIZE}."
966 )));
967 }
968 Ok(())
969}
970
971fn dedup_ids(ids: Vec<String>) -> Vec<String> {
975 let mut seen = std::collections::HashSet::new();
976 ids.into_iter().filter(|id| seen.insert(id.clone())).collect()
977}
978
979fn first_duplicate_id(ids: impl IntoIterator<Item = String>) -> Option<String> {
982 let mut seen = std::collections::HashSet::new();
983 for id in ids {
984 if !seen.insert(id.clone()) {
985 return Some(id);
986 }
987 }
988 None
989}
990
991#[async_trait::async_trait]
1002impl<E, C, U, R> crate::http::CrudService<E, C, U> for GenericCrudService<E, C, U, R>
1003where
1004 E: PersistentEntity
1005 + Clone
1006 + serde::Serialize
1007 + serde::de::DeserializeOwned
1008 + FromCreateDto<C>
1009 + ApplyUpdateDto<U>
1010 + Send
1011 + Sync
1012 + 'static,
1013 C: Send + Sync + 'static,
1014 U: Send + Sync + 'static,
1015 R: CrudRepository<E> + Send + Sync + 'static,
1016{
1017 type Error = ServiceError;
1018
1019 fn violations_of(err: &ServiceError) -> Option<Vec<crate::violation::Violation>> {
1020 match err {
1021 ServiceError::Violations(v) => Some(v.clone()),
1022 _ => None,
1023 }
1024 }
1025
1026 fn entity_name() -> &'static str {
1027 std::any::type_name::<E>()
1028 }
1029
1030 async fn fetch_related_json(&self, table: &str, ids: &[String]) -> Vec<serde_json::Value> {
1031 self.repository().fetch_related_json(table, ids).await
1032 }
1033
1034 async fn list(
1035 &self,
1036 page: u32,
1037 limit: u32,
1038 filters: HashMap<String, String>,
1039 ) -> Result<(Vec<E>, u64), ServiceError> {
1040 self.list(page, limit, filters).await
1041 }
1042
1043 async fn list_with_info(
1044 &self,
1045 page: u32,
1046 limit: u32,
1047 filters: HashMap<String, String>,
1048 ) -> Result<(Vec<E>, backbone_orm::repository::PaginationInfo), ServiceError> {
1049 self.list_with_info(page, limit, filters).await
1050 }
1051
1052 async fn aggregate(
1053 &self,
1054 spec: &backbone_orm::repository::AggregateSpec,
1055 filters: HashMap<String, String>,
1056 ) -> Result<backbone_orm::repository::AggregateResult, ServiceError> {
1057 self.aggregate(spec, filters).await
1058 }
1059
1060 fn table_name(&self) -> Option<&str> {
1061 self.repository.table_name()
1062 }
1063
1064 async fn create(&self, dto: C) -> Result<E, ServiceError> {
1065 self.create(dto).await
1066 }
1067
1068 async fn get_by_id(&self, id: &str) -> Result<Option<E>, ServiceError> {
1069 self.get_by_id(id).await
1070 }
1071
1072 async fn update(&self, id: &str, dto: U) -> Result<Option<E>, ServiceError> {
1073 self.update(id, dto).await
1074 }
1075
1076 async fn partial_update(
1077 &self,
1078 id: &str,
1079 fields: HashMap<String, serde_json::Value>,
1080 ) -> Result<Option<E>, ServiceError> {
1081 self.partial_update(id, fields).await
1082 }
1083
1084 async fn soft_delete(&self, id: &str) -> Result<bool, ServiceError> {
1085 self.soft_delete(id).await
1086 }
1087
1088 async fn bulk_create(&self, items: Vec<C>) -> Result<Vec<E>, ServiceError> {
1089 self.bulk_create(items).await
1090 }
1091
1092 async fn upsert(&self, dto: C) -> Result<E, ServiceError> {
1093 self.upsert(dto).await
1094 }
1095
1096 async fn list_deleted(&self, page: u32, limit: u32) -> Result<(Vec<E>, u64), ServiceError> {
1097 self.list_deleted(page, limit).await
1098 }
1099
1100 async fn restore(&self, id: &str) -> Result<Option<E>, ServiceError> {
1101 self.restore(id).await
1102 }
1103
1104 async fn empty_trash(&self) -> Result<u64, ServiceError> {
1105 self.empty_trash().await
1106 }
1107
1108 async fn get_deleted_by_id(&self, id: &str) -> Result<Option<E>, ServiceError> {
1109 self.get_deleted_by_id(id).await
1110 }
1111
1112 async fn permanent_delete(&self, id: &str) -> Result<bool, ServiceError> {
1113 self.permanent_delete(id).await
1114 }
1115
1116 async fn list_deleted_filtered(
1117 &self,
1118 page: u32,
1119 limit: u32,
1120 filters: HashMap<String, String>,
1121 ) -> Result<(Vec<E>, u64), ServiceError> {
1122 self.list_deleted_filtered(page, limit, filters).await
1123 }
1124
1125 async fn count_active(&self) -> Result<u64, ServiceError> {
1126 self.count_active().await
1127 }
1128
1129 async fn count_deleted(&self) -> Result<u64, ServiceError> {
1130 self.count_deleted().await
1131 }
1132
1133 async fn bulk_soft_delete(&self, ids: Vec<String>) -> Result<u64, ServiceError> {
1134 self.bulk_soft_delete(ids).await
1135 }
1136
1137 async fn bulk_restore(&self, ids: Vec<String>) -> Result<Vec<E>, ServiceError> {
1138 self.bulk_restore(ids).await
1139 }
1140
1141 async fn bulk_permanent_delete(&self, ids: Vec<String>) -> Result<u64, ServiceError> {
1142 self.bulk_permanent_delete(ids).await
1143 }
1144
1145 async fn restore_all(&self) -> Result<u64, ServiceError> {
1146 self.restore_all().await
1147 }
1148
1149 async fn bulk_update(&self, items: Vec<(String, U)>) -> Result<Vec<E>, ServiceError> {
1150 self.bulk_update(items).await
1151 }
1152
1153 async fn bulk_partial_update(
1154 &self,
1155 items: Vec<(String, HashMap<String, serde_json::Value>)>,
1156 ) -> Result<Vec<E>, ServiceError> {
1157 self.bulk_partial_update(items).await
1158 }
1159
1160 async fn preview_bulk_partial_update(
1161 &self,
1162 items: Vec<(String, HashMap<String, serde_json::Value>)>,
1163 ) -> Result<Option<Vec<crate::http::BulkPreviewRow<E, ServiceError>>>, ServiceError> {
1164 self.preview_bulk_partial_update(items).await.map(Some)
1165 }
1166}
1167
1168impl<E, C, U, R> Clone for GenericCrudService<E, C, U, R>
1170where
1171 E: PersistentEntity + Clone,
1172 R: CrudRepository<E>,
1173{
1174 fn clone(&self) -> Self {
1175 Self {
1176 repository: self.repository.clone(),
1177 lifecycle: self.lifecycle.clone(),
1178 event_publisher: self.event_publisher.clone(),
1179 guard: self.guard.clone(),
1180 _phantom: PhantomData,
1181 }
1182 }
1183}
1184
1185#[cfg(test)]
1186mod tests {
1187 use super::*;
1188 use chrono::{DateTime, Utc};
1189 use serde::{Deserialize, Serialize};
1190
1191 #[derive(Debug, Clone, Serialize, Deserialize)]
1194 struct Widget {
1195 id: String,
1196 name: String,
1197 #[serde(skip_serializing_if = "Option::is_none")]
1198 note: Option<String>,
1199 secret_hash: String,
1200 #[serde(skip_serializing_if = "Option::is_none")]
1201 deleted_at: Option<DateTime<Utc>>,
1202 created_at: DateTime<Utc>,
1203 updated_at: DateTime<Utc>,
1204 }
1205
1206 impl crate::persistence::traits::PersistentEntity for Widget {
1207 fn write_protected_fields() -> &'static [&'static str] {
1208 &["secret_hash"]
1209 }
1210 fn entity_id(&self) -> String {
1211 self.id.clone()
1212 }
1213 fn set_entity_id(&mut self, id: String) {
1214 self.id = id;
1215 }
1216 fn created_at(&self) -> Option<DateTime<Utc>> {
1217 Some(self.created_at)
1218 }
1219 fn set_created_at(&mut self, ts: DateTime<Utc>) {
1220 self.created_at = ts;
1221 }
1222 fn updated_at(&self) -> Option<DateTime<Utc>> {
1223 Some(self.updated_at)
1224 }
1225 fn set_updated_at(&mut self, ts: DateTime<Utc>) {
1226 self.updated_at = ts;
1227 }
1228 fn deleted_at(&self) -> Option<DateTime<Utc>> {
1229 self.deleted_at
1230 }
1231 fn set_deleted_at(&mut self, ts: Option<DateTime<Utc>>) {
1232 self.deleted_at = ts;
1233 }
1234 }
1235
1236 struct CreateWidgetDto {
1237 name: String,
1238 }
1239
1240 struct UpdateWidgetDto {
1241 name: String,
1242 }
1243
1244 impl FromCreateDto<CreateWidgetDto> for Widget {
1245 fn from_create_dto(dto: CreateWidgetDto) -> ServiceResult<Self> {
1246 Ok(Widget {
1247 id: uuid::Uuid::new_v4().to_string(),
1248 name: dto.name,
1249 note: None,
1250 secret_hash: "h0".into(),
1251 deleted_at: None,
1252 created_at: Utc::now(),
1253 updated_at: Utc::now(),
1254 })
1255 }
1256 }
1257
1258 impl ApplyUpdateDto<UpdateWidgetDto> for Widget {
1259 fn apply_update(mut self, dto: UpdateWidgetDto) -> ServiceResult<Self> {
1260 self.name = dto.name;
1261 self.updated_at = Utc::now();
1262 Ok(self)
1263 }
1264 }
1265
1266 struct InMemoryWidgetRepo {
1268 data: tokio::sync::Mutex<Vec<Widget>>,
1269 }
1270
1271 impl InMemoryWidgetRepo {
1272 fn new() -> Self {
1273 Self {
1274 data: tokio::sync::Mutex::new(Vec::new()),
1275 }
1276 }
1277 }
1278
1279 #[async_trait]
1280 impl crate::persistence::traits::CrudRepository<Widget> for InMemoryWidgetRepo {
1281 async fn create(&self, entity: Widget) -> Result<Widget, RepositoryError> {
1282 self.data.lock().await.push(entity.clone());
1283 Ok(entity)
1284 }
1285
1286 async fn find_by_id(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1287 Ok(self
1288 .data
1289 .lock()
1290 .await
1291 .iter()
1292 .find(|w| w.id == id && w.deleted_at.is_none())
1293 .cloned())
1294 }
1295
1296 async fn update(&self, entity: Widget) -> Result<Widget, RepositoryError> {
1297 let mut data = self.data.lock().await;
1298 data.retain(|w| w.id != entity.id);
1299 data.push(entity.clone());
1300 Ok(entity)
1301 }
1302
1303 async fn soft_delete(&self, id: &str) -> Result<bool, RepositoryError> {
1304 let mut data = self.data.lock().await;
1305 if let Some(w) = data.iter_mut().find(|w| w.id == id) {
1306 w.deleted_at = Some(Utc::now());
1307 return Ok(true);
1308 }
1309 Ok(false)
1310 }
1311
1312 async fn find_by_id_including_deleted(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1313 Ok(self.data.lock().await.iter().find(|w| w.id == id).cloned())
1314 }
1315
1316 async fn restore(&self, id: &str) -> Result<Option<Widget>, RepositoryError> {
1317 let mut data = self.data.lock().await;
1318 if let Some(w) = data.iter_mut().find(|w| w.id == id) {
1319 w.deleted_at = None;
1320 return Ok(Some(w.clone()));
1321 }
1322 Ok(None)
1323 }
1324
1325 async fn hard_delete(&self, id: &str) -> Result<bool, RepositoryError> {
1326 let mut data = self.data.lock().await;
1327 let before = data.len();
1328 data.retain(|w| w.id != id);
1329 Ok(data.len() < before)
1330 }
1331
1332 async fn list(
1333 &self,
1334 page: u32,
1335 limit: u32,
1336 ) -> Result<(Vec<Widget>, u64), RepositoryError> {
1337 let data = self.data.lock().await;
1338 let active: Vec<_> = data.iter().filter(|w| w.deleted_at.is_none()).cloned().collect();
1339 let total = active.len() as u64;
1340 let offset = ((page.saturating_sub(1)) as usize) * (limit as usize);
1341 let limit = limit as usize;
1342 let page = active.into_iter().skip(offset).take(limit).collect();
1343 Ok((page, total))
1344 }
1345
1346 async fn list_deleted(
1347 &self,
1348 page: u32,
1349 limit: u32,
1350 ) -> Result<(Vec<Widget>, u64), RepositoryError> {
1351 let data = self.data.lock().await;
1352 let deleted: Vec<_> = data.iter().filter(|w| w.deleted_at.is_some()).cloned().collect();
1353 let total = deleted.len() as u64;
1354 let offset = ((page.saturating_sub(1)) as usize) * (limit as usize);
1355 let limit = limit as usize;
1356 let page = deleted.into_iter().skip(offset).take(limit).collect();
1357 Ok((page, total))
1358 }
1359
1360 async fn count(&self) -> Result<u64, RepositoryError> {
1361 let data = self.data.lock().await;
1362 Ok(data.iter().filter(|w| w.deleted_at.is_none()).count() as u64)
1363 }
1364
1365 async fn count_deleted(&self) -> Result<u64, RepositoryError> {
1366 let data = self.data.lock().await;
1367 Ok(data.iter().filter(|w| w.deleted_at.is_some()).count() as u64)
1368 }
1369
1370 async fn bulk_create(&self, entities: Vec<Widget>) -> Result<Vec<Widget>, RepositoryError> {
1371 let mut data = self.data.lock().await;
1372 data.extend(entities.clone());
1373 Ok(entities)
1374 }
1375
1376 async fn empty_trash(&self) -> Result<u64, RepositoryError> {
1377 let mut data = self.data.lock().await;
1378 let before = data.len();
1379 data.retain(|w| w.deleted_at.is_none());
1380 Ok((before - data.len()) as u64)
1381 }
1382 }
1383
1384 #[tokio::test]
1385 async fn create_and_get_roundtrip() {
1386 let repo = Arc::new(InMemoryWidgetRepo::new());
1387 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1388 GenericCrudService::with_repository(repo);
1389
1390 let widget = service
1391 .create(CreateWidgetDto {
1392 name: "sprocket".into(),
1393 })
1394 .await
1395 .unwrap();
1396
1397 let found = service.get_by_id(&widget.id).await.unwrap();
1398 assert!(found.is_some());
1399 assert_eq!(found.unwrap().name, "sprocket");
1400 }
1401
1402 #[tokio::test]
1403 async fn soft_delete_hides_from_list() {
1404 let repo = Arc::new(InMemoryWidgetRepo::new());
1405 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1406 GenericCrudService::with_repository(repo);
1407
1408 let w = service.create(CreateWidgetDto { name: "w".into() }).await.unwrap();
1409 service.soft_delete(&w.id).await.unwrap();
1410
1411 let (items, _) = service.list(1, 20, Default::default()).await.unwrap();
1412 assert!(items.is_empty());
1413
1414 let (deleted, _) = service.list_deleted(1, 20).await.unwrap();
1415 assert_eq!(deleted.len(), 1);
1416 }
1417
1418 #[tokio::test]
1419 async fn update_publishes_event_and_returns_new_entity() {
1420 let repo = Arc::new(InMemoryWidgetRepo::new());
1421 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1422 GenericCrudService::with_repository(repo);
1423
1424 let w = service.create(CreateWidgetDto { name: "old".into() }).await.unwrap();
1425 let updated = service
1426 .update(&w.id, UpdateWidgetDto { name: "new".into() })
1427 .await
1428 .unwrap();
1429 assert_eq!(updated.unwrap().name, "new");
1430 }
1431
1432 fn svc() -> GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> {
1435 GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()))
1436 }
1437
1438 #[tokio::test]
1439 async fn bulk_soft_delete_success() {
1440 let service = svc();
1441 let mut ids = Vec::new();
1442 for n in 0..3 {
1443 ids.push(service.create(CreateWidgetDto { name: format!("w{n}") }).await.unwrap().id);
1444 }
1445 let n = service.bulk_soft_delete(ids).await.unwrap();
1446 assert_eq!(n, 3);
1447 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1448 assert!(active.is_empty());
1449 }
1450
1451 #[tokio::test]
1452 async fn bulk_soft_delete_is_all_or_nothing_on_missing_id() {
1453 let service = svc();
1454 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1455 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1456
1457 let res = service
1459 .bulk_soft_delete(vec![a.id.clone(), b.id.clone(), "does-not-exist".into()])
1460 .await;
1461 assert!(res.is_err());
1462
1463 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1464 assert_eq!(active.len(), 2, "no rows should have been deleted");
1465 }
1466
1467 #[tokio::test]
1468 async fn bulk_restore_and_restore_all() {
1469 let service = svc();
1470 let mut ids = Vec::new();
1471 for n in 0..3 {
1472 ids.push(service.create(CreateWidgetDto { name: format!("w{n}") }).await.unwrap().id);
1473 }
1474 service.bulk_soft_delete(ids.clone()).await.unwrap();
1475
1476 let restored = service.bulk_restore(vec![ids[0].clone(), ids[1].clone()]).await.unwrap();
1478 assert_eq!(restored.len(), 2);
1479
1480 let n = service.restore_all().await.unwrap();
1482 assert_eq!(n, 1);
1483
1484 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1485 assert_eq!(active.len(), 3);
1486 }
1487
1488 #[tokio::test]
1489 async fn bulk_update_is_all_or_nothing_on_missing_id() {
1490 let service = svc();
1491 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1492
1493 let res = service
1494 .bulk_update(vec![
1495 (a.id.clone(), UpdateWidgetDto { name: "A2".into() }),
1496 ("missing".into(), UpdateWidgetDto { name: "X".into() }),
1497 ])
1498 .await;
1499 assert!(res.is_err());
1500
1501 let still = service.get_by_id(&a.id).await.unwrap().unwrap();
1503 assert_eq!(still.name, "a");
1504 }
1505
1506 #[tokio::test]
1507 async fn bulk_partial_update_success() {
1508 let service = svc();
1509 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1510 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1511
1512 let mut patch_a = std::collections::HashMap::new();
1513 patch_a.insert("name".to_string(), serde_json::json!("a2"));
1514 let mut patch_b = std::collections::HashMap::new();
1515 patch_b.insert("name".to_string(), serde_json::json!("b2"));
1516
1517 let saved = service
1518 .bulk_partial_update(vec![(a.id.clone(), patch_a), (b.id.clone(), patch_b)])
1519 .await
1520 .unwrap();
1521 assert_eq!(saved.len(), 2);
1522 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a2");
1523 assert_eq!(service.get_by_id(&b.id).await.unwrap().unwrap().name, "b2");
1524 }
1525
1526 #[tokio::test]
1527 async fn bulk_soft_delete_tolerates_duplicate_ids() {
1528 let service = svc();
1529 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1530 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1531
1532 let n = service
1534 .bulk_soft_delete(vec![a.id.clone(), a.id.clone(), b.id.clone()])
1535 .await
1536 .unwrap();
1537 assert_eq!(n, 2, "two distinct rows soft-deleted");
1538
1539 let (active, _) = service.list(1, 20, Default::default()).await.unwrap();
1540 assert!(active.is_empty());
1541 }
1542
1543 #[tokio::test]
1544 async fn bulk_update_rejects_duplicate_ids() {
1545 let service = svc();
1546 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1547
1548 let res = service
1550 .bulk_update(vec![
1551 (a.id.clone(), UpdateWidgetDto { name: "first".into() }),
1552 (a.id.clone(), UpdateWidgetDto { name: "second".into() }),
1553 ])
1554 .await;
1555 assert!(res.is_err());
1556 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1557 }
1558
1559 #[tokio::test]
1560 async fn bulk_soft_delete_rejects_oversized_batch() {
1561 let service = svc();
1562 let ids: Vec<String> = (0..MAX_BATCH_SIZE + 1).map(|n| format!("id-{n}")).collect();
1563 assert!(service.bulk_soft_delete(ids).await.is_err());
1565 }
1566
1567 struct RekeyWidgetDto {
1572 name: String,
1573 secret_hash: String,
1574 }
1575
1576 impl ApplyUpdateDto<RekeyWidgetDto> for Widget {
1577 fn apply_update(mut self, dto: RekeyWidgetDto) -> ServiceResult<Self> {
1578 self.name = dto.name;
1579 self.secret_hash = dto.secret_hash;
1580 Ok(self)
1581 }
1582 }
1583
1584 type RekeyService = GenericCrudService<Widget, CreateWidgetDto, RekeyWidgetDto, InMemoryWidgetRepo>;
1585
1586 fn patch(pairs: &[(&str, serde_json::Value)]) -> HashMap<String, serde_json::Value> {
1587 pairs.iter().map(|(k, v)| (k.to_string(), v.clone())).collect()
1588 }
1589
1590 fn validation_message(res: ServiceResult<impl std::fmt::Debug>) -> String {
1593 match res {
1594 Err(ServiceError::Validation(m)) => m,
1595 Err(e @ ServiceError::Violations(_)) => e.to_string(),
1596 other => panic!("expected a validation error, got {other:?}"),
1597 }
1598 }
1599
1600 #[tokio::test]
1601 async fn patch_rejects_a_field_the_entity_does_not_have() {
1602 let service = svc();
1603 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1604
1605 let msg = validation_message(service.partial_update(&w.id, patch(&[("nmae", serde_json::json!("b"))])).await);
1606 assert!(msg.contains("unknown_field") && msg.contains("nmae"), "{msg}");
1607 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().name, "a");
1608 }
1609
1610 #[tokio::test]
1611 async fn patch_rejects_changing_the_id() {
1612 let service = svc();
1613 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1614
1615 let msg = validation_message(service.partial_update(&w.id, patch(&[("id", serde_json::json!("other"))])).await);
1616 assert!(msg.contains("field_not_writable") && msg.contains("id"), "{msg}");
1617 assert!(service.get_by_id(&w.id).await.unwrap().is_some());
1618 }
1619
1620 #[tokio::test]
1621 async fn patch_rejects_changing_a_field_the_entity_protects() {
1622 let service = svc();
1623 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1624
1625 let msg = validation_message(
1626 service.partial_update(&w.id, patch(&[("secret_hash", serde_json::json!("h1"))])).await,
1627 );
1628 assert!(msg.contains("field_not_writable") && msg.contains("secret_hash"), "{msg}");
1629 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().secret_hash, "h0");
1630 }
1631
1632 #[tokio::test]
1633 async fn patch_accepts_a_protected_field_echoed_unchanged() {
1634 let service = svc();
1635 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1636
1637 let saved = service
1638 .partial_update(&w.id, patch(&[("name", serde_json::json!("b")), ("secret_hash", serde_json::json!("h0"))]))
1639 .await
1640 .unwrap()
1641 .unwrap();
1642 assert_eq!(saved.name, "b");
1643 }
1644
1645 #[tokio::test]
1646 async fn patch_sets_and_clears_an_optional_field_that_is_not_serialized_when_empty() {
1647 let service = svc();
1648 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1649
1650 let set = service.partial_update(&w.id, patch(&[("note", serde_json::json!("hi"))])).await.unwrap().unwrap();
1651 assert_eq!(set.note.as_deref(), Some("hi"));
1652 let cleared = service.partial_update(&w.id, patch(&[("note", serde_json::Value::Null)])).await.unwrap().unwrap();
1653 assert_eq!(cleared.note, None);
1654 }
1655
1656 #[tokio::test]
1657 async fn put_rejects_changing_a_field_the_entity_protects() {
1658 let service: RekeyService = GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()));
1659 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1660
1661 let msg = validation_message(
1662 service.update(&w.id, RekeyWidgetDto { name: "b".into(), secret_hash: "h1".into() }).await,
1663 );
1664 assert!(msg.contains("field_not_writable") && msg.contains("secret_hash"), "{msg}");
1665
1666 let saved = service
1668 .update(&w.id, RekeyWidgetDto { name: "b".into(), secret_hash: "h0".into() })
1669 .await
1670 .unwrap()
1671 .unwrap();
1672 assert_eq!(saved.name, "b");
1673 }
1674
1675 #[tokio::test]
1676 async fn bulk_update_rejects_a_protected_change_before_writing_any_row() {
1677 let service: RekeyService = GenericCrudService::with_repository(Arc::new(InMemoryWidgetRepo::new()));
1678 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1679 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1680
1681 let res = service
1682 .bulk_update(vec![
1683 (a.id.clone(), RekeyWidgetDto { name: "a2".into(), secret_hash: "h0".into() }),
1684 (b.id.clone(), RekeyWidgetDto { name: "b2".into(), secret_hash: "h1".into() }),
1685 ])
1686 .await;
1687 assert!(validation_message(res).contains("field_not_writable"));
1688 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1689 }
1690
1691 #[tokio::test]
1692 async fn bulk_partial_update_rejects_a_protected_change_before_writing_any_row() {
1693 let service = svc();
1694 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1695 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1696
1697 let res = service
1698 .bulk_partial_update(vec![
1699 (a.id.clone(), patch(&[("name", serde_json::json!("a2"))])),
1700 (b.id.clone(), patch(&[("secret_hash", serde_json::json!("h1"))])),
1701 ])
1702 .await;
1703 assert!(validation_message(res).contains("field_not_writable"));
1704 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1705 }
1706
1707 #[tokio::test]
1708 async fn a_bulk_preview_says_what_each_row_would_do_and_writes_nothing() {
1709 let service = svc();
1710 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1711 let b = service.create(CreateWidgetDto { name: "b".into() }).await.unwrap();
1712
1713 let rows = service
1714 .preview_bulk_partial_update(vec![
1715 (a.id.clone(), patch(&[("name", serde_json::json!("a2"))])),
1716 (b.id.clone(), patch(&[("secret_hash", serde_json::json!("h1"))])),
1717 ("missing".into(), patch(&[("name", serde_json::json!("x"))])),
1718 ])
1719 .await
1720 .unwrap();
1721 assert_eq!(rows.len(), 3);
1722 let (before, after) = rows[0].outcome.as_ref().unwrap();
1723 assert_eq!((before.name.as_str(), after.name.as_str()), ("a", "a2"));
1724 let refused = rows[1].outcome.as_ref().err().unwrap().to_string();
1725 assert!(refused.contains("field_not_writable"), "{refused}");
1726 assert!(rows[2].outcome.is_err(), "a row that is not there is refused, not dropped");
1727 assert_eq!(service.get_by_id(&a.id).await.unwrap().unwrap().name, "a");
1729 }
1730
1731 #[tokio::test]
1732 async fn a_bulk_preview_refuses_a_batch_the_write_would_refuse_whole() {
1733 let service = svc();
1734 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1735 let res = service
1736 .preview_bulk_partial_update(vec![
1737 (a.id.clone(), patch(&[("name", serde_json::json!("a2"))])),
1738 (a.id.clone(), patch(&[("name", serde_json::json!("a3"))])),
1739 ])
1740 .await;
1741 assert!(validation_message(res.map(|r| r.len())).contains("duplicate id"));
1742 }
1743
1744 #[tokio::test]
1745 async fn bulk_partial_update_rejects_an_unknown_field() {
1746 let service = svc();
1747 let a = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1748
1749 let res = service.bulk_partial_update(vec![(a.id.clone(), patch(&[("colour", serde_json::json!("red"))]))]).await;
1750 assert!(validation_message(res).contains("unknown_field"));
1751 }
1752
1753 #[derive(Default)]
1756 struct RecordingPublisher {
1757 kinds: tokio::sync::Mutex<Vec<&'static str>>,
1758 }
1759
1760 #[async_trait]
1761 impl CrudEventPublisher<Widget> for RecordingPublisher {
1762 async fn publish(&self, event: CrudEvent<Widget>) -> Result<(), backbone_messaging::EventError> {
1763 let kind = match event {
1764 CrudEvent::Updated { .. } => "updated",
1765 CrudEvent::Patched { .. } => "patched",
1766 _ => "other",
1767 };
1768 self.kinds.lock().await.push(kind);
1769 Ok(())
1770 }
1771 }
1772
1773 #[tokio::test]
1774 async fn patch_publishes_an_updated_event() {
1775 let publisher = Arc::new(RecordingPublisher::default());
1776 let service = svc().with_event_publisher(publisher.clone());
1777 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1778 publisher.kinds.lock().await.clear();
1779
1780 service.partial_update(&w.id, patch(&[("name", serde_json::json!("b"))])).await.unwrap();
1781 assert_eq!(*publisher.kinds.lock().await, vec!["updated"]);
1782 }
1783
1784 struct RefuseTrash;
1787
1788 #[async_trait]
1789 impl ServiceLifecycle<Widget> for RefuseTrash {
1790 async fn before_restore(&self, _entity: &Widget) -> ServiceResult<()> {
1791 Err(ServiceError::Validation("restore refused".into()))
1792 }
1793 async fn before_hard_delete(&self, _entity: &Widget) -> ServiceResult<()> {
1794 Err(ServiceError::Validation("hard delete refused".into()))
1795 }
1796 async fn before_restore_all(&self) -> ServiceResult<()> {
1797 Err(ServiceError::Validation("restore all refused".into()))
1798 }
1799 async fn before_empty_trash(&self) -> ServiceResult<()> {
1800 Err(ServiceError::Validation("empty trash refused".into()))
1801 }
1802 }
1803
1804 #[tokio::test]
1805 async fn every_trash_operation_consults_the_lifecycle() {
1806 let repo = Arc::new(InMemoryWidgetRepo::new());
1807 let service: GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> =
1808 GenericCrudService::with_lifecycle(repo, Arc::new(RefuseTrash));
1809 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1810 service.soft_delete(&w.id).await.unwrap();
1811
1812 assert!(service.restore(&w.id).await.is_err(), "restore");
1813 assert!(service.bulk_restore(vec![w.id.clone()]).await.is_err(), "bulk restore");
1814 assert!(service.restore_all().await.is_err(), "restore all");
1815 assert!(service.hard_delete(&w.id).await.is_err(), "hard delete");
1816 assert!(service.bulk_permanent_delete(vec![w.id.clone()]).await.is_err(), "bulk permanent delete");
1817 assert!(service.empty_trash().await.is_err(), "empty trash");
1818
1819 let still = service.get_deleted_by_id(&w.id).await.unwrap().unwrap();
1821 assert!(still.deleted_at.is_some());
1822 }
1823
1824 #[tokio::test]
1827 async fn a_refused_field_is_a_violation_naming_its_path_and_code() {
1828 let service = svc();
1829 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1830
1831 let err = service
1832 .partial_update(&w.id, patch(&[("secret_hash", serde_json::json!("h1")), ("colour", serde_json::json!("red"))]))
1833 .await
1834 .unwrap_err();
1835 let ServiceError::Violations(v) = &err else { panic!("expected violations, got {err:?}") };
1836 let pairs: Vec<(&str, &str)> = v.iter().map(|x| (x.path.as_str(), x.code.as_str())).collect();
1837 assert_eq!(pairs, vec![("colour", "unknown_field")]);
1838 assert!(err.to_string().starts_with("validation failed: unknown_field: `colour`"), "{err}");
1840
1841 let err = service.partial_update(&w.id, patch(&[("secret_hash", serde_json::json!("h1"))])).await.unwrap_err();
1842 let ServiceError::Violations(v) = &err else { panic!("expected violations, got {err:?}") };
1843 assert_eq!((v[0].path.as_str(), v[0].code.as_str()), ("secret_hash", "field_not_writable"));
1844 }
1845
1846 #[derive(Default)]
1849 struct NameGuard {
1850 seen: std::sync::Mutex<Vec<crate::write_guard::WriteKind>>,
1851 }
1852
1853 #[async_trait]
1854 impl crate::write_guard::WriteGuard<Widget> for NameGuard {
1855 async fn check(&self, ctx: &crate::write_guard::WriteCtx<'_, Widget>) -> crate::write_guard::GuardOutcome {
1856 self.seen.lock().unwrap().push(ctx.kind);
1857 let row = ctx.after.or(ctx.before).expect("a write concerns a row");
1858 let mut out = crate::write_guard::GuardOutcome::default();
1859 match row.name.as_str() {
1860 "forbidden" => out.refuse.push(crate::violation::Violation::new("name", "name_forbidden", "this name is refused")),
1861 "watched" => out.shadow.push(crate::violation::Violation::new("name", "name_watched", "would be refused")),
1862 _ => {}
1863 }
1864 out
1865 }
1866 }
1867
1868 fn guarded(guard: Arc<NameGuard>) -> GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo> {
1869 svc().with_guard(guard)
1870 }
1871
1872 #[tokio::test]
1873 async fn a_guard_refusal_stops_create_update_and_patch_with_nothing_written() {
1874 let guard = Arc::new(NameGuard::default());
1875 let service = guarded(guard.clone());
1876
1877 let err = service.create(CreateWidgetDto { name: "forbidden".into() }).await.unwrap_err();
1878 assert!(matches!(&err, ServiceError::Violations(v) if v[0].code == "name_forbidden"), "{err:?}");
1879 assert!(service.list(1, 20, Default::default()).await.unwrap().0.is_empty(), "nothing created");
1880
1881 let w = service.create(CreateWidgetDto { name: "ok".into() }).await.unwrap();
1882 assert!(service.update(&w.id, UpdateWidgetDto { name: "forbidden".into() }).await.is_err());
1883 assert!(service.partial_update(&w.id, patch(&[("name", serde_json::json!("forbidden"))])).await.is_err());
1884 assert!(service
1885 .bulk_partial_update(vec![(w.id.clone(), patch(&[("name", serde_json::json!("forbidden"))]))])
1886 .await
1887 .is_err());
1888 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().name, "ok");
1889 }
1890
1891 #[tokio::test]
1892 async fn a_shadowed_rule_lets_the_write_through() {
1893 let service = guarded(Arc::new(NameGuard::default()));
1894 let w = service.create(CreateWidgetDto { name: "watched".into() }).await.unwrap();
1895 assert_eq!(service.get_by_id(&w.id).await.unwrap().unwrap().name, "watched");
1896 }
1897
1898 #[tokio::test]
1899 async fn the_guard_is_asked_about_every_write_kind() {
1900 use crate::write_guard::WriteKind::*;
1901 let guard = Arc::new(NameGuard::default());
1902 let service = guarded(guard.clone());
1903 let w = service.create(CreateWidgetDto { name: "a".into() }).await.unwrap();
1904 service.update(&w.id, UpdateWidgetDto { name: "b".into() }).await.unwrap();
1905 service.partial_update(&w.id, patch(&[("name", serde_json::json!("c"))])).await.unwrap();
1906 service.soft_delete(&w.id).await.unwrap();
1907 service.restore(&w.id).await.unwrap();
1908 service.soft_delete(&w.id).await.unwrap();
1909 service.hard_delete(&w.id).await.unwrap();
1910 assert_eq!(*guard.seen.lock().unwrap(), vec![Create, Update, Update, Delete, Restore, Delete, HardDelete]);
1911 }
1912
1913 #[test]
1914 fn the_generic_service_hands_its_violations_to_the_http_layer() {
1915 use crate::http::CrudService;
1916 type S = GenericCrudService<Widget, CreateWidgetDto, UpdateWidgetDto, InMemoryWidgetRepo>;
1917 let v = vec![crate::violation::Violation::new("name", "x", "y")];
1918 assert_eq!(<S as CrudService<Widget, CreateWidgetDto, UpdateWidgetDto>>::violations_of(&ServiceError::Violations(v.clone())), Some(v));
1919 assert_eq!(<S as CrudService<Widget, CreateWidgetDto, UpdateWidgetDto>>::violations_of(&ServiceError::NotFound), None);
1920 }
1921}