1use async_trait::async_trait;
21use std::collections::HashMap;
22use std::marker::PhantomData;
23use std::sync::Arc;
24
25use crate::service::{ServiceError, ServiceResult};
26
27#[derive(Debug, Clone)]
34pub struct GenericCreateCommand<E, DTO> {
35 pub payload: DTO,
36 pub correlation_id: Option<String>,
37 _phantom: PhantomData<E>,
38}
39
40impl<E, DTO> GenericCreateCommand<E, DTO> {
41 pub fn new(payload: DTO) -> Self {
42 Self {
43 payload,
44 correlation_id: None,
45 _phantom: PhantomData,
46 }
47 }
48
49 pub fn with_correlation(mut self, id: impl Into<String>) -> Self {
50 self.correlation_id = Some(id.into());
51 self
52 }
53}
54
55impl<E: Send + Sync, DTO: Send + Sync> crate::command::Command
56 for GenericCreateCommand<E, DTO>
57{
58 type Result = E;
59}
60
61#[derive(Debug, Clone)]
63pub struct GenericUpdateCommand<E, DTO> {
64 pub id: String,
65 pub payload: DTO,
66 pub correlation_id: Option<String>,
67 _phantom: PhantomData<E>,
68}
69
70impl<E, DTO> GenericUpdateCommand<E, DTO> {
71 pub fn new(id: impl Into<String>, payload: DTO) -> Self {
72 Self {
73 id: id.into(),
74 payload,
75 correlation_id: None,
76 _phantom: PhantomData,
77 }
78 }
79}
80
81impl<E: Send + Sync, DTO: Send + Sync> crate::command::Command
82 for GenericUpdateCommand<E, DTO>
83{
84 type Result = Option<E>;
85}
86
87#[derive(Debug, Clone)]
89pub struct GenericDeleteCommand<E> {
90 pub id: String,
91 _phantom: PhantomData<E>,
92}
93
94impl<E> GenericDeleteCommand<E> {
95 pub fn new(id: impl Into<String>) -> Self {
96 Self {
97 id: id.into(),
98 _phantom: PhantomData,
99 }
100 }
101}
102
103impl<E: Send + Sync> crate::command::Command for GenericDeleteCommand<E> {
104 type Result = bool;
105}
106
107#[derive(Debug, Clone)]
109pub struct GenericRestoreCommand<E> {
110 pub id: String,
111 _phantom: PhantomData<E>,
112}
113
114impl<E> GenericRestoreCommand<E> {
115 pub fn new(id: impl Into<String>) -> Self {
116 Self {
117 id: id.into(),
118 _phantom: PhantomData,
119 }
120 }
121}
122
123impl<E: Send + Sync> crate::command::Command for GenericRestoreCommand<E> {
124 type Result = Option<E>;
125}
126
127#[derive(Debug, Clone)]
131pub struct GenericGetQuery<E> {
132 pub id: String,
133 _phantom: PhantomData<E>,
134}
135
136impl<E> GenericGetQuery<E> {
137 pub fn new(id: impl Into<String>) -> Self {
138 Self {
139 id: id.into(),
140 _phantom: PhantomData,
141 }
142 }
143}
144
145impl<E: Send + Sync> crate::query::Query for GenericGetQuery<E> {
146 type Result = Option<E>;
147}
148
149#[derive(Debug, Clone)]
153pub struct GenericListQuery<E, F = HashMap<String, String>> {
154 pub page: u32,
155 pub limit: u32,
156 pub filters: F,
157 _phantom: PhantomData<E>,
158}
159
160impl<E, F: Default> GenericListQuery<E, F> {
161 pub fn new(page: u32, limit: u32) -> Self {
162 Self {
163 page,
164 limit,
165 filters: F::default(),
166 _phantom: PhantomData,
167 }
168 }
169
170 pub fn with_filters(mut self, filters: F) -> Self {
171 self.filters = filters;
172 self
173 }
174}
175
176impl<E: Send + Sync, F: Send + Sync> crate::query::Query for GenericListQuery<E, F> {
177 type Result = (Vec<E>, u64);
178}
179
180#[derive(Debug, Clone)]
182pub struct GenericListDeletedQuery<E> {
183 pub page: u32,
184 pub limit: u32,
185 _phantom: PhantomData<E>,
186}
187
188impl<E> GenericListDeletedQuery<E> {
189 pub fn new(page: u32, limit: u32) -> Self {
190 Self {
191 page,
192 limit,
193 _phantom: PhantomData,
194 }
195 }
196}
197
198impl<E: Send + Sync> crate::query::Query for GenericListDeletedQuery<E> {
199 type Result = (Vec<E>, u64);
200}
201
202pub struct GenericCommandHandler<E, C, U, S>
211where
212 E: Clone + Send + Sync + 'static,
213 C: Clone + Send + Sync + 'static,
214 U: Clone + Send + Sync + 'static,
215 S: Send + Sync + 'static,
216{
217 service: Arc<S>,
218 _phantom: PhantomData<(E, C, U)>,
219}
220
221impl<E, C, U, S> GenericCommandHandler<E, C, U, S>
222where
223 E: Clone + Send + Sync + 'static,
224 C: Clone + Send + Sync + 'static,
225 U: Clone + Send + Sync + 'static,
226 S: Send + Sync + 'static,
227{
228 pub fn new(service: Arc<S>) -> Self {
229 Self {
230 service,
231 _phantom: PhantomData,
232 }
233 }
234}
235
236#[async_trait]
238pub trait CqrsService<E, C, U>: Send + Sync + 'static
239where
240 E: Clone + Send + Sync + 'static,
241 C: Send + Sync + 'static,
242 U: Send + Sync + 'static,
243{
244 async fn create(&self, dto: C) -> ServiceResult<E>;
245 async fn update(&self, id: &str, dto: U) -> ServiceResult<Option<E>>;
246 async fn soft_delete(&self, id: &str) -> ServiceResult<bool>;
247 async fn restore(&self, id: &str) -> ServiceResult<Option<E>>;
248 async fn get_by_id(&self, id: &str) -> ServiceResult<Option<E>>;
249 async fn list(&self, page: u32, limit: u32, filters: HashMap<String, String>) -> ServiceResult<(Vec<E>, u64)>;
250 async fn list_deleted(&self, page: u32, limit: u32) -> ServiceResult<(Vec<E>, u64)>;
251}
252
253#[async_trait]
255impl<E, C, U, S> crate::command::CommandHandler<GenericCreateCommand<E, C>>
256 for GenericCommandHandler<E, C, U, S>
257where
258 E: Clone + Send + Sync + 'static,
259 C: Clone + Send + Sync + 'static,
260 U: Clone + Send + Sync + 'static,
261 S: CqrsService<E, C, U>,
262{
263 type Error = ServiceError;
264
265 async fn handle(&self, command: GenericCreateCommand<E, C>) -> Result<E, Self::Error> {
266 self.service.create(command.payload).await
267 }
268}
269
270#[async_trait]
272impl<E, C, U, S> crate::command::CommandHandler<GenericUpdateCommand<E, U>>
273 for GenericCommandHandler<E, C, U, S>
274where
275 E: Clone + Send + Sync + 'static,
276 C: Clone + Send + Sync + 'static,
277 U: Clone + Send + Sync + 'static,
278 S: CqrsService<E, C, U>,
279{
280 type Error = ServiceError;
281
282 async fn handle(&self, command: GenericUpdateCommand<E, U>) -> Result<Option<E>, Self::Error> {
283 self.service.update(&command.id, command.payload).await
284 }
285}
286
287#[async_trait]
289impl<E, C, U, S> crate::command::CommandHandler<GenericDeleteCommand<E>>
290 for GenericCommandHandler<E, C, U, S>
291where
292 E: Clone + Send + Sync + 'static,
293 C: Clone + Send + Sync + 'static,
294 U: Clone + Send + Sync + 'static,
295 S: CqrsService<E, C, U>,
296{
297 type Error = ServiceError;
298
299 async fn handle(&self, command: GenericDeleteCommand<E>) -> Result<bool, Self::Error> {
300 self.service.soft_delete(&command.id).await
301 }
302}
303
304pub struct GenericQueryHandler<E, F, S>
308where
309 E: Clone + Send + Sync + 'static,
310 F: Clone + Send + Sync + Default + 'static,
311 S: Send + Sync + 'static,
312{
313 service: Arc<S>,
314 _phantom: PhantomData<(E, F)>,
315}
316
317impl<E, F, S> GenericQueryHandler<E, F, S>
318where
319 E: Clone + Send + Sync + 'static,
320 F: Clone + Send + Sync + Default + 'static,
321 S: Send + Sync + 'static,
322{
323 pub fn new(service: Arc<S>) -> Self {
324 Self {
325 service,
326 _phantom: PhantomData,
327 }
328 }
329}
330
331#[async_trait]
334pub trait CqrsReadService<E>: Send + Sync + 'static
335where
336 E: Clone + Send + Sync + 'static,
337{
338 async fn get_by_id(&self, id: &str) -> ServiceResult<Option<E>>;
339 async fn list(&self, page: u32, limit: u32, filters: HashMap<String, String>) -> ServiceResult<(Vec<E>, u64)>;
340 async fn list_deleted(&self, page: u32, limit: u32) -> ServiceResult<(Vec<E>, u64)>;
341}
342
343#[async_trait]
345impl<E, F, S> crate::query::QueryHandler<GenericGetQuery<E>>
346 for GenericQueryHandler<E, F, S>
347where
348 E: Clone + Send + Sync + 'static,
349 F: Clone + Send + Sync + Default + 'static,
350 S: CqrsReadService<E>,
351{
352 type Error = ServiceError;
353
354 async fn handle(&self, query: GenericGetQuery<E>) -> Result<Option<E>, Self::Error> {
355 self.service.get_by_id(&query.id).await
356 }
357}
358
359#[cfg(test)]
360mod tests {
361 use super::*;
362
363 #[test]
364 fn create_command_carries_payload() {
365 let cmd: GenericCreateCommand<String, i32> =
366 GenericCreateCommand::new(42).with_correlation("corr-1");
367 assert_eq!(cmd.payload, 42);
368 assert_eq!(cmd.correlation_id.as_deref(), Some("corr-1"));
369 }
370
371 #[test]
372 fn delete_command_carries_id() {
373 let cmd: GenericDeleteCommand<String> = GenericDeleteCommand::new("entity-1");
374 assert_eq!(cmd.id, "entity-1");
375 }
376
377 #[test]
378 fn list_query_default_filters() {
379 let q: GenericListQuery<String> = GenericListQuery::new(1, 20);
380 assert_eq!(q.page, 1);
381 assert_eq!(q.limit, 20);
382 assert!(q.filters.is_empty());
383 }
384}