1use async_trait::async_trait;
28use std::collections::HashMap;
29use std::sync::Arc;
30
31use super::traits::{CrudRepository, PersistentEntity, RepositoryError};
32use crate::http::CrudService;
33
34type CreateMapperFn<C, E> = Box<dyn Fn(C) -> E + Send + Sync>;
36
37type UpdateMapperFn<E, U> = Box<dyn Fn(&mut E, U) + Send + Sync>;
39
40pub 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 create_mapper: CreateMapperFn<C, E>,
61 update_mapper: UpdateMapperFn<E, U>,
63 #[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 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 pub fn repository(&self) -> &R {
96 &self.repository
97 }
98}
99
100#[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 "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 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 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 if let Some(existing) = self.repository.find_by_id(&id).await? {
205 Ok(self.repository.update(existing).await?)
207 } else {
208 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 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
248pub struct SimpleCrudServiceAdapter<R, E>
256where
257 R: CrudRepository<E> + Send + Sync,
258 E: PersistentEntity + Clone,
259{
260 repository: Arc<R>,
261 #[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 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
389use super::traits::SearchableRepository;
394
395pub 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 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 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 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 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}