Skip to main content

backbone_core/persistence/
adapter.rs

1//! CrudService Adapter
2//!
3//! Bridges the `CrudRepository` trait (from persistence) to the `CrudService` trait
4//! (from HTTP layer), enabling automatic HTTP endpoint generation for any repository.
5//!
6//! # Usage
7//!
8//! ```ignore
9//! use backbone_core::persistence::{InMemoryRepository, CrudServiceAdapter};
10//! use backbone_core::{BackboneCrudHandler};
11//!
12//! // Create repository
13//! let repo = InMemoryRepository::<User>::new();
14//!
15//! // Wrap with CrudServiceAdapter
16//! let service = CrudServiceAdapter::new(
17//!     Arc::new(repo),
18//!     |dto: CreateUserDto| User { /* convert DTO to entity */ },
19//!     |entity: &mut User, dto: UpdateUserDto| { /* apply update */ },
20//! );
21//!
22//! // Generate all 12 HTTP endpoints
23//! let routes = BackboneCrudHandler::<_, User, CreateUserDto, UpdateUserDto, UserResponse>
24//!     ::routes(Arc::new(service), "/api/v1/users");
25//! ```
26
27use async_trait::async_trait;
28use std::collections::HashMap;
29use std::sync::Arc;
30
31use super::traits::{CrudRepository, PersistentEntity, RepositoryError};
32use crate::http::CrudService;
33
34/// Type alias for create mapper function
35type CreateMapperFn<C, E> = Box<dyn Fn(C) -> E + Send + Sync>;
36
37/// Type alias for update mapper function
38type UpdateMapperFn<E, U> = Box<dyn Fn(&mut E, U) + Send + Sync>;
39
40/// Adapter that wraps a `CrudRepository` to implement `CrudService`
41///
42/// This enables using any repository implementation with the `BackboneCrudHandler`
43/// to automatically generate all 12 HTTP endpoints.
44///
45/// # Type Parameters
46///
47/// - `R`: Repository implementing `CrudRepository<E>`
48/// - `E`: Entity type (must implement `PersistentEntity`)
49/// - `C`: Create DTO type
50/// - `U`: Update DTO type
51pub struct CrudServiceAdapter<R, E, C, U>
52where
53    R: CrudRepository<E> + Send + Sync,
54    E: PersistentEntity,
55    C: Send + Sync,
56    U: Send + Sync,
57{
58    repository: Arc<R>,
59    /// Function to convert CreateDto to Entity
60    create_mapper: CreateMapperFn<C, E>,
61    /// Function to apply UpdateDto to Entity
62    update_mapper: UpdateMapperFn<E, U>,
63    /// Entity name for error messages (reserved for future use)
64    #[allow(dead_code)]
65    entity_name: &'static str,
66}
67
68impl<R, E, C, U> CrudServiceAdapter<R, E, C, U>
69where
70    R: CrudRepository<E> + Send + Sync,
71    E: PersistentEntity,
72    C: Send + Sync,
73    U: Send + Sync,
74{
75    /// Create a new adapter with custom mappers
76    pub fn new<F1, F2>(
77        repository: Arc<R>,
78        entity_name: &'static str,
79        create_mapper: F1,
80        update_mapper: F2,
81    ) -> Self
82    where
83        F1: Fn(C) -> E + Send + Sync + 'static,
84        F2: Fn(&mut E, U) + Send + Sync + 'static,
85    {
86        Self {
87            repository,
88            create_mapper: Box::new(create_mapper),
89            update_mapper: Box::new(update_mapper),
90            entity_name,
91        }
92    }
93
94    /// Get a reference to the underlying repository
95    pub fn repository(&self) -> &R {
96        &self.repository
97    }
98}
99
100/// Error type for CrudServiceAdapter
101#[derive(Debug, thiserror::Error)]
102pub enum AdapterError {
103    #[error("Repository error: {0}")]
104    Repository(#[from] RepositoryError),
105
106    #[error("Validation error: {0}")]
107    Validation(String),
108
109    #[error("Not found: {0}")]
110    NotFound(String),
111}
112
113impl From<AdapterError> for String {
114    fn from(e: AdapterError) -> Self {
115        e.to_string()
116    }
117}
118
119#[async_trait]
120impl<R, E, C, U> CrudService<E, C, U> for CrudServiceAdapter<R, E, C, U>
121where
122    R: CrudRepository<E> + Send + Sync + 'static,
123    E: PersistentEntity + Clone + 'static,
124    C: Send + Sync + 'static,
125    U: Send + Sync + 'static,
126{
127    type Error = AdapterError;
128
129    fn entity_name() -> &'static str {
130        // Note: This is a static method, so we can't access self.entity_name
131        // For dynamic entity names, consider using a different error approach
132        "Entity"
133    }
134
135    async fn list(
136        &self,
137        page: u32,
138        limit: u32,
139        _filters: HashMap<String, String>,
140    ) -> Result<(Vec<E>, u64), Self::Error> {
141        // Note: Filter support requires the repository to implement SearchableRepository.
142        // For repositories that only implement CrudRepository, filters are ignored.
143        // Use SearchableCrudServiceAdapter for full filter support.
144        Ok(self.repository.list(page, limit).await?)
145    }
146
147    async fn aggregate(
148        &self,
149        spec: &backbone_orm::repository::AggregateSpec,
150        filters: HashMap<String, String>,
151    ) -> Result<backbone_orm::repository::AggregateResult, Self::Error> {
152        Ok(self.repository.aggregate_filtered(spec, filters).await?)
153    }
154
155    fn table_name(&self) -> Option<&str> {
156        self.repository.table_name()
157    }
158
159    async fn create(&self, dto: C) -> Result<E, Self::Error> {
160        let entity = (self.create_mapper)(dto);
161        Ok(self.repository.create(entity).await?)
162    }
163
164    async fn get_by_id(&self, id: &str) -> Result<Option<E>, Self::Error> {
165        Ok(self.repository.find_by_id(id).await?)
166    }
167
168    async fn update(&self, id: &str, dto: U) -> Result<Option<E>, Self::Error> {
169        let entity = self.repository.find_by_id(id).await?;
170
171        let Some(mut entity) = entity else {
172            return Ok(None);
173        };
174
175        (self.update_mapper)(&mut entity, dto);
176        Ok(Some(self.repository.update(entity).await?))
177    }
178
179    async fn partial_update(
180        &self,
181        id: &str,
182        _fields: HashMap<String, serde_json::Value>,
183    ) -> Result<Option<E>, Self::Error> {
184        // Partial update requires entity to implement PartialUpdatable
185        // For now, just return the entity as-is (no-op)
186        // Modules can override this behavior
187        Ok(self.repository.find_by_id(id).await?)
188    }
189
190    async fn soft_delete(&self, id: &str) -> Result<bool, Self::Error> {
191        Ok(self.repository.soft_delete(id).await?)
192    }
193
194    async fn bulk_create(&self, items: Vec<C>) -> Result<Vec<E>, Self::Error> {
195        let entities: Vec<E> = items.into_iter().map(&*self.create_mapper).collect();
196        Ok(self.repository.bulk_create(entities).await?)
197    }
198
199    async fn upsert(&self, dto: C) -> Result<E, Self::Error> {
200        let entity = (self.create_mapper)(dto);
201        let id = entity.entity_id();
202
203        // Check if exists
204        if let Some(existing) = self.repository.find_by_id(&id).await? {
205            // Update existing
206            Ok(self.repository.update(existing).await?)
207        } else {
208            // Create new
209            Ok(self.repository.create(entity).await?)
210        }
211    }
212
213    async fn list_deleted(&self, page: u32, limit: u32) -> Result<(Vec<E>, u64), Self::Error> {
214        Ok(self.repository.list_deleted(page, limit).await?)
215    }
216
217    async fn restore(&self, id: &str) -> Result<Option<E>, Self::Error> {
218        Ok(self.repository.restore(id).await?)
219    }
220
221    async fn empty_trash(&self) -> Result<u64, Self::Error> {
222        Ok(self.repository.empty_trash().await?)
223    }
224
225    async fn get_deleted_by_id(&self, id: &str) -> Result<Option<E>, Self::Error> {
226        Ok(self.repository.find_by_id_including_deleted(id).await?)
227    }
228
229    async fn permanent_delete(&self, id: &str) -> Result<bool, Self::Error> {
230        Ok(self.repository.hard_delete(id).await?)
231    }
232
233    async fn list_deleted_filtered(&self, page: u32, limit: u32, filters: std::collections::HashMap<String, String>) -> Result<(Vec<E>, u64), Self::Error> {
234        // Default implementation: ignore filters and use regular list_deleted
235        let _ = filters;
236        Ok(self.repository.list_deleted(page, limit).await?)
237    }
238
239    async fn count_active(&self) -> Result<u64, Self::Error> {
240        Ok(self.repository.count().await?)
241    }
242
243    async fn count_deleted(&self) -> Result<u64, Self::Error> {
244        Ok(self.repository.count_deleted().await?)
245    }
246}
247
248// ============================================================
249// Simple Adapter for Identity Mapping
250// ============================================================
251
252/// Simple adapter for entities where Create/Update DTOs are the same as Entity
253///
254/// Use this when your DTOs are identical to your entity type.
255pub struct SimpleCrudServiceAdapter<R, E>
256where
257    R: CrudRepository<E> + Send + Sync,
258    E: PersistentEntity + Clone,
259{
260    repository: Arc<R>,
261    /// Entity name for error messages (reserved for future use)
262    #[allow(dead_code)]
263    entity_name: &'static str,
264    _phantom: std::marker::PhantomData<E>,
265}
266
267impl<R, E> SimpleCrudServiceAdapter<R, E>
268where
269    R: CrudRepository<E> + Send + Sync,
270    E: PersistentEntity + Clone,
271{
272    pub fn new(repository: Arc<R>, entity_name: &'static str) -> Self {
273        Self {
274            repository,
275            entity_name,
276            _phantom: std::marker::PhantomData,
277        }
278    }
279}
280
281#[async_trait]
282impl<R, E> CrudService<E, E, E> for SimpleCrudServiceAdapter<R, E>
283where
284    R: CrudRepository<E> + Send + Sync + 'static,
285    E: PersistentEntity + Clone + 'static,
286{
287    type Error = AdapterError;
288
289    fn entity_name() -> &'static str {
290        "Entity"
291    }
292
293    async fn list(
294        &self,
295        page: u32,
296        limit: u32,
297        _filters: HashMap<String, String>,
298    ) -> Result<(Vec<E>, u64), Self::Error> {
299        Ok(self.repository.list(page, limit).await?)
300    }
301
302    async fn aggregate(
303        &self,
304        spec: &backbone_orm::repository::AggregateSpec,
305        filters: HashMap<String, String>,
306    ) -> Result<backbone_orm::repository::AggregateResult, Self::Error> {
307        Ok(self.repository.aggregate_filtered(spec, filters).await?)
308    }
309
310    fn table_name(&self) -> Option<&str> {
311        self.repository.table_name()
312    }
313
314    async fn create(&self, entity: E) -> Result<E, Self::Error> {
315        Ok(self.repository.create(entity).await?)
316    }
317
318    async fn get_by_id(&self, id: &str) -> Result<Option<E>, Self::Error> {
319        Ok(self.repository.find_by_id(id).await?)
320    }
321
322    async fn update(&self, id: &str, entity: E) -> Result<Option<E>, Self::Error> {
323        if self.repository.find_by_id(id).await?.is_none() {
324            return Ok(None);
325        }
326        Ok(Some(self.repository.update(entity).await?))
327    }
328
329    async fn partial_update(
330        &self,
331        id: &str,
332        _fields: HashMap<String, serde_json::Value>,
333    ) -> Result<Option<E>, Self::Error> {
334        Ok(self.repository.find_by_id(id).await?)
335    }
336
337    async fn soft_delete(&self, id: &str) -> Result<bool, Self::Error> {
338        Ok(self.repository.soft_delete(id).await?)
339    }
340
341    async fn bulk_create(&self, items: Vec<E>) -> Result<Vec<E>, Self::Error> {
342        Ok(self.repository.bulk_create(items).await?)
343    }
344
345    async fn upsert(&self, entity: E) -> Result<E, Self::Error> {
346        let id = entity.entity_id();
347        if self.repository.find_by_id(&id).await?.is_some() {
348            Ok(self.repository.update(entity).await?)
349        } else {
350            Ok(self.repository.create(entity).await?)
351        }
352    }
353
354    async fn list_deleted(&self, page: u32, limit: u32) -> Result<(Vec<E>, u64), Self::Error> {
355        Ok(self.repository.list_deleted(page, limit).await?)
356    }
357
358    async fn restore(&self, id: &str) -> Result<Option<E>, Self::Error> {
359        Ok(self.repository.restore(id).await?)
360    }
361
362    async fn empty_trash(&self) -> Result<u64, Self::Error> {
363        Ok(self.repository.empty_trash().await?)
364    }
365
366    async fn get_deleted_by_id(&self, id: &str) -> Result<Option<E>, Self::Error> {
367        Ok(self.repository.find_by_id_including_deleted(id).await?)
368    }
369
370    async fn permanent_delete(&self, id: &str) -> Result<bool, Self::Error> {
371        Ok(self.repository.hard_delete(id).await?)
372    }
373
374    async fn list_deleted_filtered(&self, page: u32, limit: u32, filters: std::collections::HashMap<String, String>) -> Result<(Vec<E>, u64), Self::Error> {
375        // Default implementation: ignore filters and use regular list_deleted
376        let _ = filters;
377        Ok(self.repository.list_deleted(page, limit).await?)
378    }
379
380    async fn count_active(&self) -> Result<u64, Self::Error> {
381        Ok(self.repository.count().await?)
382    }
383
384    async fn count_deleted(&self) -> Result<u64, Self::Error> {
385        Ok(self.repository.count_deleted().await?)
386    }
387}
388
389// ============================================================
390// Searchable Adapter with Filter Support
391// ============================================================
392
393use super::traits::SearchableRepository;
394
395/// Adapter for repositories that implement `SearchableRepository`
396///
397/// This adapter provides full filter support in the `list` method by leveraging
398/// the `SearchableRepository::search` method.
399///
400/// # Example
401///
402/// ```ignore
403/// use backbone_core::persistence::{InMemoryRepository, SearchableCrudServiceAdapter};
404///
405/// let repo = Arc::new(InMemoryRepository::<User>::new());
406/// let service = SearchableCrudServiceAdapter::new(
407///     repo,
408///     "User",
409///     |dto: CreateUserDto| User::from(dto),
410///     |entity: &mut User, dto: UpdateUserDto| entity.apply_update(dto),
411/// );
412/// ```
413pub struct SearchableCrudServiceAdapter<R, E, C, U>
414where
415    R: SearchableRepository<E> + Send + Sync,
416    E: PersistentEntity,
417    C: Send + Sync,
418    U: Send + Sync,
419{
420    repository: Arc<R>,
421    create_mapper: CreateMapperFn<C, E>,
422    update_mapper: UpdateMapperFn<E, U>,
423    #[allow(dead_code)]
424    entity_name: &'static str,
425}
426
427impl<R, E, C, U> SearchableCrudServiceAdapter<R, E, C, U>
428where
429    R: SearchableRepository<E> + Send + Sync,
430    E: PersistentEntity,
431    C: Send + Sync,
432    U: Send + Sync,
433{
434    /// Create a new searchable adapter with custom mappers
435    pub fn new<F1, F2>(
436        repository: Arc<R>,
437        entity_name: &'static str,
438        create_mapper: F1,
439        update_mapper: F2,
440    ) -> Self
441    where
442        F1: Fn(C) -> E + Send + Sync + 'static,
443        F2: Fn(&mut E, U) + Send + Sync + 'static,
444    {
445        Self {
446            repository,
447            create_mapper: Box::new(create_mapper),
448            update_mapper: Box::new(update_mapper),
449            entity_name,
450        }
451    }
452
453    /// Get a reference to the underlying repository
454    pub fn repository(&self) -> &R {
455        &self.repository
456    }
457}
458
459#[async_trait]
460impl<R, E, C, U> CrudService<E, C, U> for SearchableCrudServiceAdapter<R, E, C, U>
461where
462    R: SearchableRepository<E> + Send + Sync + 'static,
463    E: PersistentEntity + Clone + 'static,
464    C: Send + Sync + 'static,
465    U: Send + Sync + 'static,
466{
467    type Error = AdapterError;
468
469    fn entity_name() -> &'static str {
470        "Entity"
471    }
472
473    async fn list(
474        &self,
475        page: u32,
476        limit: u32,
477        filters: HashMap<String, String>,
478    ) -> Result<(Vec<E>, u64), Self::Error> {
479        // Use SearchableRepository::search for full filter support
480        if filters.is_empty() {
481            Ok(self.repository.list(page, limit).await?)
482        } else {
483            Ok(self.repository.search(filters, page, limit).await?)
484        }
485    }
486
487    async fn aggregate(
488        &self,
489        spec: &backbone_orm::repository::AggregateSpec,
490        filters: HashMap<String, String>,
491    ) -> Result<backbone_orm::repository::AggregateResult, Self::Error> {
492        Ok(self.repository.aggregate_filtered(spec, filters).await?)
493    }
494
495    fn table_name(&self) -> Option<&str> {
496        self.repository.table_name()
497    }
498
499    async fn create(&self, dto: C) -> Result<E, Self::Error> {
500        let entity = (self.create_mapper)(dto);
501        Ok(self.repository.create(entity).await?)
502    }
503
504    async fn get_by_id(&self, id: &str) -> Result<Option<E>, Self::Error> {
505        Ok(self.repository.find_by_id(id).await?)
506    }
507
508    async fn update(&self, id: &str, dto: U) -> Result<Option<E>, Self::Error> {
509        let entity = self.repository.find_by_id(id).await?;
510
511        let Some(mut entity) = entity else {
512            return Ok(None);
513        };
514
515        (self.update_mapper)(&mut entity, dto);
516        Ok(Some(self.repository.update(entity).await?))
517    }
518
519    async fn partial_update(
520        &self,
521        id: &str,
522        _fields: HashMap<String, serde_json::Value>,
523    ) -> Result<Option<E>, Self::Error> {
524        Ok(self.repository.find_by_id(id).await?)
525    }
526
527    async fn soft_delete(&self, id: &str) -> Result<bool, Self::Error> {
528        Ok(self.repository.soft_delete(id).await?)
529    }
530
531    async fn bulk_create(&self, items: Vec<C>) -> Result<Vec<E>, Self::Error> {
532        let entities: Vec<E> = items.into_iter().map(&*self.create_mapper).collect();
533        Ok(self.repository.bulk_create(entities).await?)
534    }
535
536    async fn upsert(&self, dto: C) -> Result<E, Self::Error> {
537        let entity = (self.create_mapper)(dto);
538        let id = entity.entity_id();
539
540        if let Some(existing) = self.repository.find_by_id(&id).await? {
541            Ok(self.repository.update(existing).await?)
542        } else {
543            Ok(self.repository.create(entity).await?)
544        }
545    }
546
547    async fn list_deleted(&self, page: u32, limit: u32) -> Result<(Vec<E>, u64), Self::Error> {
548        Ok(self.repository.list_deleted(page, limit).await?)
549    }
550
551    async fn restore(&self, id: &str) -> Result<Option<E>, Self::Error> {
552        Ok(self.repository.restore(id).await?)
553    }
554
555    async fn empty_trash(&self) -> Result<u64, Self::Error> {
556        Ok(self.repository.empty_trash().await?)
557    }
558
559    async fn get_deleted_by_id(&self, id: &str) -> Result<Option<E>, Self::Error> {
560        Ok(self.repository.find_by_id_including_deleted(id).await?)
561    }
562
563    async fn permanent_delete(&self, id: &str) -> Result<bool, Self::Error> {
564        Ok(self.repository.hard_delete(id).await?)
565    }
566
567    async fn list_deleted_filtered(&self, page: u32, limit: u32, filters: std::collections::HashMap<String, String>) -> Result<(Vec<E>, u64), Self::Error> {
568        // Default implementation: ignore filters and use regular list_deleted
569        let _ = filters;
570        Ok(self.repository.list_deleted(page, limit).await?)
571    }
572
573    async fn count_active(&self) -> Result<u64, Self::Error> {
574        Ok(self.repository.count().await?)
575    }
576
577    async fn count_deleted(&self) -> Result<u64, Self::Error> {
578        Ok(self.repository.count_deleted().await?)
579    }
580}