Skip to main content

allsource_core/application/use_cases/
query_events.rs

1use crate::{
2    application::dto::{EventDto, QueryEventsRequest, QueryEventsResponse},
3    domain::repositories::EventRepository,
4    error::Result,
5};
6use std::sync::Arc;
7
8/// Use Case: Query Events
9///
10/// This use case handles querying events from the event store with various filters.
11///
12/// Responsibilities:
13/// - Validate query parameters
14/// - Determine query strategy (by entity, by type, by time range, etc.)
15/// - Execute query via repository
16/// - Transform domain events to DTOs
17/// - Apply limits and pagination
18pub struct QueryEventsUseCase {
19    repository: Arc<dyn EventRepository>,
20}
21
22impl QueryEventsUseCase {
23    pub fn new(repository: Arc<dyn EventRepository>) -> Self {
24        Self { repository }
25    }
26
27    pub async fn execute(&self, request: QueryEventsRequest) -> Result<QueryEventsResponse> {
28        // Determine tenant_id (default to "default" if not provided)
29        let tenant_id = request.tenant_id.unwrap_or_else(|| "default".to_string());
30
31        // Determine query strategy based on filters
32        let mut events = if let Some(entity_id) = request.entity_id {
33            // Query by entity (most specific)
34            if let Some(as_of) = request.as_of {
35                // Time-travel query
36                self.repository
37                    .find_by_entity_as_of(&entity_id, &tenant_id, as_of)
38                    .await?
39            } else {
40                self.repository
41                    .find_by_entity(&entity_id, &tenant_id)
42                    .await?
43            }
44        } else if let Some(event_type) = request.event_type {
45            // Query by type
46            self.repository
47                .find_by_type(&event_type, &tenant_id)
48                .await?
49        } else if let (Some(since), Some(until)) = (request.since, request.until) {
50            // Query by time range
51            self.repository
52                .find_by_time_range(&tenant_id, since, until)
53                .await?
54        } else {
55            // No specific filter - this could be expensive!
56            // In production, you might want to require at least one filter
57            return Err(crate::error::Error::InvalidInput(
58                "Query requires at least one filter (entity_id, event_type, or time range)"
59                    .to_string(),
60            ));
61        };
62
63        // Apply time filters if provided (for non-time-range queries)
64        if let Some(since) = request.since {
65            events.retain(|e| e.occurred_after(since));
66        }
67        if let Some(until) = request.until {
68            events.retain(|e| e.occurred_before(until));
69        }
70
71        // Capture total before applying limit
72        let total_count = events.len();
73
74        // Apply limit
75        if let Some(limit) = request.limit {
76            events.truncate(limit);
77        }
78
79        // Convert to DTOs
80        let event_dtos: Vec<EventDto> = events.iter().map(EventDto::from).collect();
81        let count = event_dtos.len();
82        let has_more = count < total_count;
83
84        Ok(QueryEventsResponse {
85            events: event_dtos,
86            count,
87            total_count,
88            has_more,
89            entity_version: None,
90            archive_integrity: None,
91        })
92    }
93}
94
95#[cfg(test)]
96mod tests {
97    use super::*;
98    use crate::domain::entities::Event;
99    use async_trait::async_trait;
100    use chrono::Utc;
101    use serde_json::json;
102    use uuid::Uuid;
103
104    // Mock repository for testing
105    struct MockEventRepository {
106        events: Vec<Event>,
107    }
108
109    impl MockEventRepository {
110        fn with_events(events: Vec<Event>) -> Self {
111            Self { events }
112        }
113    }
114
115    #[async_trait]
116    impl EventRepository for MockEventRepository {
117        async fn save(&self, _event: &Event) -> Result<()> {
118            Ok(())
119        }
120
121        async fn save_batch(&self, _events: &[Event]) -> Result<()> {
122            Ok(())
123        }
124
125        async fn find_by_id(&self, id: Uuid) -> Result<Option<Event>> {
126            Ok(self.events.iter().find(|e| e.id() == id).cloned())
127        }
128
129        async fn find_by_entity(&self, entity_id: &str, tenant_id: &str) -> Result<Vec<Event>> {
130            Ok(self
131                .events
132                .iter()
133                .filter(|e| e.entity_id_str() == entity_id && e.tenant_id_str() == tenant_id)
134                .cloned()
135                .collect())
136        }
137
138        async fn find_by_type(&self, event_type: &str, tenant_id: &str) -> Result<Vec<Event>> {
139            Ok(self
140                .events
141                .iter()
142                .filter(|e| e.event_type_str() == event_type && e.tenant_id_str() == tenant_id)
143                .cloned()
144                .collect())
145        }
146
147        async fn find_by_time_range(
148            &self,
149            tenant_id: &str,
150            start: chrono::DateTime<Utc>,
151            end: chrono::DateTime<Utc>,
152        ) -> Result<Vec<Event>> {
153            Ok(self
154                .events
155                .iter()
156                .filter(|e| e.tenant_id_str() == tenant_id && e.occurred_between(start, end))
157                .cloned()
158                .collect())
159        }
160
161        async fn find_by_entity_as_of(
162            &self,
163            entity_id: &str,
164            tenant_id: &str,
165            as_of: chrono::DateTime<Utc>,
166        ) -> Result<Vec<Event>> {
167            Ok(self
168                .events
169                .iter()
170                .filter(|e| {
171                    e.entity_id_str() == entity_id
172                        && e.tenant_id_str() == tenant_id
173                        && e.occurred_before(as_of)
174                })
175                .cloned()
176                .collect())
177        }
178
179        async fn count(&self, tenant_id: &str) -> Result<usize> {
180            Ok(self
181                .events
182                .iter()
183                .filter(|e| e.tenant_id_str() == tenant_id)
184                .count())
185        }
186
187        async fn health_check(&self) -> Result<()> {
188            Ok(())
189        }
190    }
191
192    fn create_test_events() -> Vec<Event> {
193        vec![
194            Event::from_strings(
195                "user.created".to_string(),
196                "user-1".to_string(),
197                "tenant-1".to_string(),
198                json!({"name": "Alice"}),
199                None,
200            )
201            .unwrap(),
202            Event::from_strings(
203                "user.created".to_string(),
204                "user-2".to_string(),
205                "tenant-1".to_string(),
206                json!({"name": "Bob"}),
207                None,
208            )
209            .unwrap(),
210            Event::from_strings(
211                "order.placed".to_string(),
212                "order-1".to_string(),
213                "tenant-1".to_string(),
214                json!({"amount": 100}),
215                None,
216            )
217            .unwrap(),
218        ]
219    }
220
221    #[tokio::test]
222    async fn test_query_by_entity() {
223        let events = create_test_events();
224        let entity_id = events[0].entity_id().to_string();
225        let repo = Arc::new(MockEventRepository::with_events(events));
226        let use_case = QueryEventsUseCase::new(repo);
227
228        let request = QueryEventsRequest {
229            entity_id: Some(entity_id),
230            event_type: None,
231            tenant_id: Some("tenant-1".to_string()),
232            as_of: None,
233            since: None,
234            until: None,
235            limit: None,
236            event_type_prefix: None,
237            exclude_event_type_prefix: None,
238            payload_filter: None,
239        };
240
241        let response = use_case.execute(request).await;
242        assert!(response.is_ok());
243
244        let response = response.unwrap();
245        assert_eq!(response.count, 1);
246    }
247
248    #[tokio::test]
249    async fn test_query_by_type() {
250        let events = create_test_events();
251        let repo = Arc::new(MockEventRepository::with_events(events));
252        let use_case = QueryEventsUseCase::new(repo);
253
254        let request = QueryEventsRequest {
255            entity_id: None,
256            event_type: Some("user.created".to_string()),
257            tenant_id: Some("tenant-1".to_string()),
258            as_of: None,
259            since: None,
260            until: None,
261            limit: None,
262            event_type_prefix: None,
263            exclude_event_type_prefix: None,
264            payload_filter: None,
265        };
266
267        let response = use_case.execute(request).await;
268        assert!(response.is_ok());
269
270        let response = response.unwrap();
271        assert_eq!(response.count, 2);
272    }
273
274    #[tokio::test]
275    async fn test_query_with_limit() {
276        let events = create_test_events();
277        let repo = Arc::new(MockEventRepository::with_events(events));
278        let use_case = QueryEventsUseCase::new(repo);
279
280        let request = QueryEventsRequest {
281            entity_id: None,
282            event_type: Some("user.created".to_string()),
283            tenant_id: Some("tenant-1".to_string()),
284            as_of: None,
285            since: None,
286            until: None,
287            limit: Some(1),
288            event_type_prefix: None,
289            exclude_event_type_prefix: None,
290            payload_filter: None,
291        };
292
293        let response = use_case.execute(request).await;
294        assert!(response.is_ok());
295
296        let response = response.unwrap();
297        assert_eq!(response.count, 1);
298    }
299
300    #[tokio::test]
301    async fn test_query_requires_filter() {
302        let events = create_test_events();
303        let repo = Arc::new(MockEventRepository::with_events(events));
304        let use_case = QueryEventsUseCase::new(repo);
305
306        let request = QueryEventsRequest {
307            entity_id: None,
308            event_type: None,
309            tenant_id: Some("tenant-1".to_string()),
310            as_of: None,
311            since: None,
312            until: None,
313            limit: None,
314            event_type_prefix: None,
315            exclude_event_type_prefix: None,
316            payload_filter: None,
317        };
318
319        let response = use_case.execute(request).await;
320        assert!(response.is_err());
321    }
322}