allsource_core/application/use_cases/
query_events.rs1use crate::{
2 application::dto::{EventDto, QueryEventsRequest, QueryEventsResponse},
3 domain::repositories::EventRepository,
4 error::Result,
5};
6use std::sync::Arc;
7
8pub 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 let tenant_id = request.tenant_id.unwrap_or_else(|| "default".to_string());
30
31 let mut events = if let Some(entity_id) = request.entity_id {
33 if let Some(as_of) = request.as_of {
35 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 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 self.repository
52 .find_by_time_range(&tenant_id, since, until)
53 .await?
54 } else {
55 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 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 let total_count = events.len();
73
74 if let Some(limit) = request.limit {
76 events.truncate(limit);
77 }
78
79 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 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}