Skip to main content

backbone_core/
service.rs

1//! Generic CRUD service — one implementation used by all entities.
2//!
3//! Generated code emits only a type alias:
4//!
5//! ```rust,ignore
6//! // Generated (was ~250 lines, now 1 line):
7//! pub type StoredFileService = GenericCrudService<
8//!     StoredFile,
9//!     CreateStoredFileDto,
10//!     UpdateStoredFileDto,
11//!     PostgresStoredFileRepository,
12//! >;
13//! ```
14//!
15//! Custom services wrap the alias with a decorator:
16//!
17//! ```rust,ignore
18//! pub struct StoredFileServiceCustom {
19//!     inner: Arc<StoredFileService>,
20//!     notifications: Arc<NotificationService>,
21//! }
22//!
23//! impl StoredFileServiceCustom {
24//!     pub async fn archive(&self, id: &str, reason: String) -> ServiceResult<StoredFile> {
25//!         let entity = self.inner.get_by_id(id).await?.ok_or(ServiceError::NotFound)?;
26//!         // custom business logic
27//!     }
28//! }
29//! ```
30
31use 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// ─── Service error ───────────────────────────────────────────────────────────
41
42/// Standard service-layer error type.
43#[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    /// The write broke named rules: one violation per rule, each with the
55    /// field it concerns and a stable code. Displays as the same
56    /// `validation failed: …` sentence a plain validation error does.
57    #[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
69// ─── DTO conversion traits ────────────────────────────────────────────────────
70
71/// Convert a Create DTO into a new entity instance.
72///
73/// Implemented once per entity by generated code, may be overridden in custom
74/// decorators via `UseCaseHooks`.
75pub trait FromCreateDto<DTO>: Sized {
76    fn from_create_dto(dto: DTO) -> ServiceResult<Self>;
77}
78
79/// Apply an Update DTO's fields to an existing entity (full update).
80pub trait ApplyUpdateDto<DTO> {
81    fn apply_update(self, dto: DTO) -> ServiceResult<Self>
82    where
83        Self: Sized;
84}
85
86// ─── ServiceLifecycle hooks ──────────────────────────────────────────────────
87
88/// Optional lifecycle callbacks injected into `GenericCrudService`.
89///
90/// Override individual methods — all default to no-op.
91#[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    /// Called before a soft-deleted entity is restored, singly or in a batch.
124    async fn before_restore(&self, entity: &E) -> ServiceResult<()> {
125        let _ = entity;
126        Ok(())
127    }
128
129    /// Called before a soft-deleted entity is removed for good, singly or in
130    /// a batch.
131    async fn before_hard_delete(&self, entity: &E) -> ServiceResult<()> {
132        let _ = entity;
133        Ok(())
134    }
135
136    /// Called once before the whole trash is restored. The rows are not
137    /// loaded for this; refuse here to keep the bulk path closed.
138    async fn before_restore_all(&self) -> ServiceResult<()> {
139        Ok(())
140    }
141
142    /// Called once before the whole trash is removed for good. The rows are
143    /// not loaded for this; refuse here to keep the bulk path closed.
144    async fn before_empty_trash(&self) -> ServiceResult<()> {
145        Ok(())
146    }
147}
148
149/// No-op lifecycle — the default used by generated type aliases.
150pub 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
171// ─── GenericCrudService ──────────────────────────────────────────────────────
172
173/// A fully generic CRUD service.
174///
175/// Type parameters:
176/// - `E`  — entity type (must implement `PersistentEntity`)
177/// - `C`  — create DTO (entity implements `FromCreateDto<C>`)
178/// - `U`  — update DTO (entity implements `ApplyUpdateDto<U>`)
179/// - `R`  — repository (implements `CrudRepository<E>`)
180///
181/// The service always holds an `Arc<dyn CrudEventPublisher<E>>`.  When no
182/// publisher is needed, `NoOpCrudEventPublisher` is used (zero overhead).
183/// Event-publish errors are fire-and-forget — they never fail the operation.
184pub 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    /// Full constructor — supply all three dependencies.
204    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    /// Check every write against `guard`: its refusals refuse the write, its
219    /// shadowed violations are logged and the write proceeds.
220    pub fn with_guard(mut self, guard: Arc<dyn crate::write_guard::WriteGuard<E>>) -> Self {
221        self.guard = guard;
222        self
223    }
224
225    /// Ask the guard about one write.
226    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    /// Create with a custom lifecycle and the no-op event publisher.
253    pub fn with_lifecycle(repository: Arc<R>, lifecycle: Arc<dyn ServiceLifecycle<E>>) -> Self {
254        Self::new(repository, lifecycle, NoOpCrudEventPublisher::arc())
255    }
256
257    /// Create with the no-op lifecycle and no-op event publisher (generated default).
258    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    /// Replace the event publisher on an existing service (builder pattern).
267    pub fn with_event_publisher(mut self, publisher: Arc<dyn CrudEventPublisher<E>>) -> Self {
268        self.event_publisher = publisher;
269        self
270    }
271
272    /// Access the underlying repository for custom query extensions.
273    pub fn repository(&self) -> &R {
274        &self.repository
275    }
276
277    // ─── Individual operations ────────────────────────────────────────────
278
279    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    /// `list`, carrying the pagination info (cursor positions included) so
292    /// the HTTP layer can surface keyset paging on every generic route.
293    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    /// Group and reduce the rows `list` would return, under the same filters.
306    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    /// Alias for `get_by_id` — used by generated handlers.
485    pub async fn find_by_id(&self, id: &str) -> ServiceResult<Option<E>> {
486        self.get_by_id(id).await
487    }
488
489    /// Alias for `hard_delete` — used by generated handlers.
490    pub async fn permanent_delete(&self, id: &str) -> ServiceResult<bool> {
491        self.hard_delete(id).await
492    }
493
494    /// Count soft-deleted entities.
495    pub async fn count_deleted(&self) -> ServiceResult<u64> {
496        self.repository
497            .count_deleted()
498            .await
499            .map_err(ServiceError::Repository)
500    }
501
502    /// Retrieve a soft-deleted entity by ID (includes deleted records).
503    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    /// List deleted entities with filter support (filters currently ignored, delegates to list_deleted).
511    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    /// Upsert: create the entity if it doesn't exist (no ID-based lookup — always creates).
521    pub async fn upsert(&self, dto: C) -> ServiceResult<E> {
522        self.create(dto).await
523    }
524
525    /// Apply a partial update from a JSON map of field values.
526    ///
527    /// Fields not present in the map are left unchanged. A key the entity does
528    /// not have, or a new value for a protected field, refuses the whole write.
529    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    // ── Atomic batch operations ───────────────────────────────────────────────
577    //
578    // All-or-nothing: ids are validated up-front (so a missing id fails the whole
579    // request before any write), then the repository applies the change inside a
580    // single transaction. Lifecycle hooks and CRUD events fire per affected
581    // entity, mirroring the single-row methods. Events are best-effort.
582
583    /// Soft-delete many entities by id (atomic).
584    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        // Load the entities up-front: they are the payloads for the lifecycle
588        // `before_delete` hook and the `SoftDeleted` event below. The repository's
589        // atomic affected-row check (not this loop) is what guarantees the write
590        // is all-or-nothing.
591        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    /// Restore many soft-deleted entities by id (atomic).
624    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        // Ask the lifecycle about every row it can see; an id that is not in
628        // the trash is left to the repository's atomic check below.
629        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    /// Restore every soft-deleted entity (atomic). Returns the number restored.
651    /// Emits a `Restored` event per entity, mirroring [`bulk_restore`](Self::bulk_restore).
652    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    /// Permanently delete many soft-deleted entities by id (atomic).
670    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        // The lifecycle is asked about every row it can see. Whether every id
674        // is in the trash stays with `bulk_hard_delete`, which deletes only
675        // trashed rows and rolls the whole batch back unless every id matched,
676        // so a row that moves between this loop and the write cannot slip
677        // through.
678        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    /// Full-update many entities (atomic). Each item is `(id, UpdateDto)`.
703    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    /// Fire `after_update` and a `Updated` event for each `(before, after)` pair.
738    /// Shared by [`bulk_update`](Self::bulk_update) and
739    /// [`bulk_partial_update`](Self::bulk_partial_update); events are best-effort.
740    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    /// Partial-update many entities (atomic). Each item is `(id, field_map)`.
757    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
797// ─── Generic write protection ────────────────────────────────────────────────
798
799/// Serialized fields no generic write may change on any entity: the row's
800/// identity, and the audit block that records who created, changed and
801/// deleted it. Entities add their own through
802/// [`PersistentEntity::write_protected_fields`].
803pub 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
812/// Refuse a generic write that gives a protected field a new value. Carrying
813/// the stored value unchanged is allowed, so a form that sends the whole
814/// record is not refused for fields it did not touch.
815fn 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
841/// Merge a partial update onto an entity.
842///
843/// The merge goes through the entity's JSON form, so a key the entity has no
844/// field for would otherwise be dropped without a word and the edit lost. It
845/// is refused instead: a key given a value that does not survive the round
846/// trip is not a field of this record. A key set to `null` is exempt, since
847/// an empty optional field may be left out of the JSON form entirely.
848fn 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
882/// Maximum number of ids / items a single batch operation may contain. Enforced
883/// here at the service layer (not only in the HTTP handlers) so every caller —
884/// gRPC, background jobs, internal code — gets the same bound.
885pub const MAX_BATCH_SIZE: usize = 1000;
886
887/// Reject an over-sized batch before any work is done.
888fn 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
897/// Drop duplicate ids while preserving first-seen order. Id-list batch ops issue
898/// a SQL `IN (...)` whose distinct-row count would never match a list that
899/// repeats an id, so a stray duplicate would otherwise fail the whole batch.
900fn 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
905/// Return the first id that appears more than once, if any. Used to reject
906/// ambiguous bulk updates (the same id mapped to two different payloads).
907fn 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// ─── Blanket CrudService impl ────────────────────────────────────────────────
918
919/// Blanket impl so that any `GenericCrudService<E, C, U, R>` automatically
920/// satisfies `CrudService<E, C, U>` without generated adapter structs.
921///
922/// Generated code emits only a type alias:
923///   `pub type UserService = GenericCrudService<User, CreateUserDto, UpdateUserDto, UserRepository>;`
924///
925/// That alias now directly implements `CrudService` and can be passed to
926/// `BackboneCrudHandler` without any additional boilerplate.
927#[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
1087// Make the service cloneable (needed when Arc-ing it and sharing across handlers).
1088impl<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    // ── Minimal in-memory repository for testing ──────────────────────────
1111
1112    #[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    // Minimal in-memory CrudRepository
1186    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    // ── Batch operations ──────────────────────────────────────────────────────
1352
1353    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        // One id is missing → whole batch must be rejected with nothing deleted.
1377        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        // Restore two explicitly.
1396        let restored = service.bulk_restore(vec![ids[0].clone(), ids[1].clone()]).await.unwrap();
1397        assert_eq!(restored.len(), 2);
1398
1399        // Restore the remaining one via restore_all.
1400        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        // The valid row must be untouched because the batch aborted pre-write.
1421        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        // A repeated id must not fail the batch: ids are de-duplicated first.
1452        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        // The same id mapped to two payloads is ambiguous → rejected, nothing written.
1468        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        // Enforced at the service layer, independent of any HTTP handler.
1483        assert!(service.bulk_soft_delete(ids).await.is_err());
1484    }
1485
1486    // ── Generic writes may not change protected or unknown fields ─────────
1487
1488    /// A full update that carries the protected field, the way a form that
1489    /// replaces the whole record would.
1490    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    /// The refusal's message, whether the service refused with named
1510    /// violations or with a plain validation sentence.
1511    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        // The same form echoing the stored value is an ordinary edit.
1586        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    // ── Events and trash lifecycle ────────────────────────────────────────
1636
1637    #[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    /// A lifecycle that refuses every trash operation, so a test can tell
1667    /// whether the operation consulted it.
1668    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        // Nothing moved: the row is still in the trash.
1702        let still = service.get_deleted_by_id(&w.id).await.unwrap().unwrap();
1703        assert!(still.deleted_at.is_some());
1704    }
1705
1706    // ── Named violations and the write guard ──────────────────────────────
1707
1708    #[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        // The sentence clients already read is unchanged.
1721        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    /// Refuses any write whose resulting widget is named `forbidden`, and
1729    /// shadow-reports any named `watched`. Records every kind it was asked about.
1730    #[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}