1use crate::{
4 application::{ApplicationError, ApplicationResult, handlers::QueryHandlerGat, queries::*},
5 domain::{
6 SessionState,
7 entities::Stream,
8 ports::{
9 FrameStoreGat, SessionPagination, SessionQueryCriteria, SortOrder as RepoSortOrder,
10 StreamRepositoryGat, StreamStoreGat,
11 },
12 },
13};
14use std::{marker::PhantomData, sync::Arc, time::Instant};
15
16const MAX_FRAMES_PAGE_SIZE: usize = crate::domain::config::MAX_PAGINATION_LIMIT;
20
21#[derive(Debug)]
23pub struct SessionQueryHandler<R>
24where
25 R: StreamRepositoryGat + 'static,
26{
27 repository: Arc<R>,
28}
29
30impl<R> SessionQueryHandler<R>
31where
32 R: StreamRepositoryGat + 'static,
33{
34 pub fn new(repository: Arc<R>) -> Self {
36 Self { repository }
37 }
38}
39
40impl<R> QueryHandlerGat<GetSessionQuery> for SessionQueryHandler<R>
41where
42 R: StreamRepositoryGat + Send + Sync,
43{
44 type Response = SessionResponse;
45
46 type HandleFuture<'a>
47 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
48 where
49 Self: 'a;
50
51 fn handle(&self, query: GetSessionQuery) -> Self::HandleFuture<'_> {
52 async move {
53 let session = self
54 .repository
55 .find_session(query.session_id.into())
56 .await
57 .map_err(ApplicationError::Domain)?
58 .ok_or_else(|| {
59 ApplicationError::NotFound(format!("Session {} not found", query.session_id))
60 })?;
61
62 Ok(SessionResponse { session })
63 }
64 }
65}
66
67impl<R> QueryHandlerGat<GetActiveSessionsQuery> for SessionQueryHandler<R>
68where
69 R: StreamRepositoryGat + Send + Sync,
70{
71 type Response = SessionsResponse;
72
73 type HandleFuture<'a>
74 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
75 where
76 Self: 'a;
77
78 fn handle(&self, query: GetActiveSessionsQuery) -> Self::HandleFuture<'_> {
79 async move {
80 const MAX_PAGE_SIZE: usize = 100;
81 let limit = query.limit.unwrap_or(MAX_PAGE_SIZE).min(MAX_PAGE_SIZE);
82 let offset = query.offset.unwrap_or(0);
83
84 let pagination = SessionPagination {
85 offset,
86 limit,
87 sort_by: None,
88 sort_order: RepoSortOrder::Ascending,
89 };
90
91 let result = self
92 .repository
93 .find_sessions_by_criteria(SessionQueryCriteria::default(), pagination)
94 .await
95 .map_err(ApplicationError::Domain)?;
96
97 Ok(SessionsResponse {
98 sessions: result.sessions,
99 total_count: result.total_count,
100 has_more: result.has_more,
101 })
102 }
103 }
104}
105
106impl<R> QueryHandlerGat<GetSessionHealthQuery> for SessionQueryHandler<R>
107where
108 R: StreamRepositoryGat + Send + Sync,
109{
110 type Response = HealthResponse;
111
112 type HandleFuture<'a>
113 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
114 where
115 Self: 'a;
116
117 fn handle(&self, query: GetSessionHealthQuery) -> Self::HandleFuture<'_> {
118 async move {
119 let session = self
120 .repository
121 .find_session(query.session_id.into())
122 .await
123 .map_err(ApplicationError::Domain)?
124 .ok_or_else(|| {
125 ApplicationError::NotFound(format!("Session {} not found", query.session_id))
126 })?;
127
128 let health = session.health_check();
129
130 Ok(HealthResponse { health })
131 }
132 }
133}
134
135impl<R> QueryHandlerGat<SearchSessionsQuery> for SessionQueryHandler<R>
136where
137 R: StreamRepositoryGat + Send + Sync,
138{
139 type Response = SessionsResponse;
140
141 type HandleFuture<'a>
142 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
143 where
144 Self: 'a;
145
146 fn handle(&self, query: SearchSessionsQuery) -> Self::HandleFuture<'_> {
147 async move {
148 const MAX_PAGE_SIZE: usize = 100;
149 let limit = query.limit.unwrap_or(MAX_PAGE_SIZE).clamp(1, MAX_PAGE_SIZE);
150 let offset = query
151 .offset
152 .unwrap_or(0)
153 .min(crate::domain::config::MAX_PAGINATION_OFFSET);
154
155 let (states, exclude_expired) = match query.filters.state {
161 Some(state) => (Some(vec![state.as_str().to_string()]), false),
162 None => (Some(vec![SessionState::Active.as_str().to_string()]), true),
163 };
164
165 let criteria = SessionQueryCriteria {
166 states,
167 exclude_expired,
168 created_after: query.filters.created_after,
169 created_before: query.filters.created_before,
170 client_info_pattern: query.filters.client_info,
171 has_active_streams: query.filters.has_active_streams,
172 ..SessionQueryCriteria::default()
173 };
174
175 let sort_by = query.sort_by;
176 let sort_order = match query.sort_order {
177 Some(SortOrder::Descending) => RepoSortOrder::Descending,
178 _ => RepoSortOrder::Ascending,
179 };
180
181 let pagination = SessionPagination {
182 offset,
183 limit,
184 sort_by,
185 sort_order,
186 };
187
188 let result = self
189 .repository
190 .find_sessions_by_criteria(criteria, pagination)
191 .await
192 .map_err(ApplicationError::Domain)?;
193
194 Ok(SessionsResponse {
195 sessions: result.sessions,
196 total_count: result.total_count,
197 has_more: result.has_more,
198 })
199 }
200 }
201}
202
203impl<R> QueryHandlerGat<GetSessionStatsQuery> for SessionQueryHandler<R>
204where
205 R: StreamRepositoryGat + Send + Sync,
206{
207 type Response = SessionStatsResponse;
208
209 type HandleFuture<'a>
210 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
211 where
212 Self: 'a;
213
214 fn handle(&self, query: GetSessionStatsQuery) -> Self::HandleFuture<'_> {
215 async move {
216 let session = self
217 .repository
218 .find_session(query.session_id.into())
219 .await
220 .map_err(ApplicationError::Domain)?
221 .ok_or_else(|| {
222 ApplicationError::NotFound(format!("Session {} not found", query.session_id))
223 })?;
224
225 let streams = session.streams();
226 let active_stream_count = streams.values().filter(|s| s.is_active()).count();
227
228 Ok(SessionStatsResponse {
229 session_id: session.id().into(),
230 stats: session.stats().clone(),
231 stream_count: streams.len(),
232 active_stream_count,
233 created_at: session.created_at(),
234 updated_at: session.updated_at(),
235 duration_ms: session.duration().map(|d| d.num_milliseconds()),
236 })
237 }
238 }
239}
240
241#[derive(Debug)]
243pub struct StreamQueryHandler<R, S, F>
244where
245 R: StreamRepositoryGat + 'static,
246 S: StreamStoreGat + 'static,
247 F: FrameStoreGat + 'static,
248{
249 session_repository: Arc<R>,
250 frame_store: Arc<F>,
251 _phantom: PhantomData<S>,
252}
253
254impl<R, S, F> StreamQueryHandler<R, S, F>
255where
256 R: StreamRepositoryGat + 'static,
257 S: StreamStoreGat + 'static,
258 F: FrameStoreGat + 'static,
259{
260 pub fn new(session_repository: Arc<R>, _stream_store: Arc<S>, frame_store: Arc<F>) -> Self {
262 Self {
263 session_repository,
264 frame_store,
265 _phantom: PhantomData,
266 }
267 }
268}
269
270impl<R, S, F> QueryHandlerGat<GetStreamQuery> for StreamQueryHandler<R, S, F>
271where
272 R: StreamRepositoryGat + Send + Sync,
273 S: StreamStoreGat + Send + Sync,
274 F: FrameStoreGat + Send + Sync,
275{
276 type Response = StreamResponse;
277
278 type HandleFuture<'a>
279 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
280 where
281 Self: 'a;
282
283 fn handle(&self, query: GetStreamQuery) -> Self::HandleFuture<'_> {
284 async move {
285 let session = self
286 .session_repository
287 .find_session(query.session_id.into())
288 .await
289 .map_err(ApplicationError::Domain)?
290 .ok_or_else(|| {
291 ApplicationError::NotFound(format!("Session {} not found", query.session_id))
292 })?;
293
294 let stream = session
295 .stream(query.stream_id.into())
296 .ok_or_else(|| {
297 ApplicationError::NotFound(format!("Stream {} not found", query.stream_id))
298 })?
299 .clone();
300
301 Ok(StreamResponse { stream })
302 }
303 }
304}
305
306impl<R, S, F> QueryHandlerGat<GetStreamsForSessionQuery> for StreamQueryHandler<R, S, F>
307where
308 R: StreamRepositoryGat + Send + Sync,
309 S: StreamStoreGat + Send + Sync,
310 F: FrameStoreGat + Send + Sync,
311{
312 type Response = StreamsResponse;
313
314 type HandleFuture<'a>
315 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
316 where
317 Self: 'a;
318
319 fn handle(&self, query: GetStreamsForSessionQuery) -> Self::HandleFuture<'_> {
320 async move {
321 let session = self
322 .session_repository
323 .find_session(query.session_id.into())
324 .await
325 .map_err(ApplicationError::Domain)?
326 .ok_or_else(|| {
327 ApplicationError::NotFound(format!("Session {} not found", query.session_id))
328 })?;
329
330 let streams: Vec<Stream> = session
331 .streams()
332 .values()
333 .filter(|stream| query.include_inactive || stream.is_active())
334 .cloned()
335 .collect();
336
337 Ok(StreamsResponse { streams })
338 }
339 }
340}
341
342impl<R, S, F> QueryHandlerGat<GetStreamFramesQuery> for StreamQueryHandler<R, S, F>
343where
344 R: StreamRepositoryGat + Send + Sync,
345 S: StreamStoreGat + Send + Sync,
346 F: FrameStoreGat + Send + Sync,
347{
348 type Response = FramesResponse;
349
350 type HandleFuture<'a>
351 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
352 where
353 Self: 'a;
354
355 fn handle(&self, query: GetStreamFramesQuery) -> Self::HandleFuture<'_> {
356 async move {
357 let session = self
358 .session_repository
359 .find_session(query.session_id.into())
360 .await
361 .map_err(ApplicationError::Domain)?
362 .ok_or_else(|| {
363 ApplicationError::NotFound(format!("Session {} not found", query.session_id))
364 })?;
365
366 session.stream(query.stream_id.into()).ok_or_else(|| {
368 ApplicationError::NotFound(format!("Stream {} not found", query.stream_id))
369 })?;
370
371 let limit = Some(
375 query
376 .limit
377 .map(|l| l.min(MAX_FRAMES_PAGE_SIZE))
378 .unwrap_or(MAX_FRAMES_PAGE_SIZE),
379 );
380 let priority_filter = query
381 .priority_filter
382 .map(crate::domain::value_objects::Priority::try_from)
383 .transpose()
384 .map_err(ApplicationError::Domain)?;
385
386 let page = self
387 .frame_store
388 .get_frames(
389 query.stream_id.into(),
390 query.since_sequence,
391 priority_filter,
392 limit,
393 )
394 .await
395 .map_err(ApplicationError::Domain)?;
396
397 Ok(FramesResponse {
398 frames: page.frames,
399 total_count: page.total_matching,
400 })
401 }
402 }
403}
404
405#[derive(Debug)]
407pub struct SystemQueryHandler<R>
408where
409 R: StreamRepositoryGat + 'static,
410{
411 repository: Arc<R>,
412 started_at: Instant,
413}
414
415impl<R> SystemQueryHandler<R>
416where
417 R: StreamRepositoryGat + 'static,
418{
419 pub fn new(repository: Arc<R>) -> Self {
421 Self {
422 repository,
423 started_at: Instant::now(),
424 }
425 }
426
427 pub fn with_start_time(repository: Arc<R>, started_at: Instant) -> Self {
438 Self {
439 repository,
440 started_at,
441 }
442 }
443}
444
445impl<R> QueryHandlerGat<GetSystemStatsQuery> for SystemQueryHandler<R>
446where
447 R: StreamRepositoryGat + Send + Sync,
448{
449 type Response = SystemStatsResponse;
450
451 type HandleFuture<'a>
452 = impl std::future::Future<Output = ApplicationResult<Self::Response>> + Send + 'a
453 where
454 Self: 'a;
455
456 fn handle(&self, _query: GetSystemStatsQuery) -> Self::HandleFuture<'_> {
457 async move {
458 let sessions = self
459 .repository
460 .find_active_sessions()
461 .await
462 .map_err(ApplicationError::Domain)?;
463
464 let total_sessions = sessions.len() as u64;
465 let active_sessions = sessions.iter().filter(|s| s.is_active()).count() as u64;
466
467 let mut total_streams = 0u64;
468 let mut active_streams = 0u64;
469 let mut total_frames = 0u64;
470 let mut total_bytes = 0u64;
471 let mut total_duration_ms = 0f64;
472 let mut completed_sessions = 0u64;
473
474 for session in &sessions {
475 let stats = session.stats();
476 total_streams += stats.total_streams;
477 active_streams += stats.active_streams;
478 total_frames += stats.total_frames;
479 total_bytes += stats.total_bytes;
480
481 if let Some(duration) = session.duration() {
482 total_duration_ms += duration.num_milliseconds() as f64;
483 completed_sessions += 1;
484 }
485 }
486
487 let average_session_duration_seconds = if completed_sessions > 0 {
488 total_duration_ms / completed_sessions as f64 / 1000.0
489 } else {
490 0.0
491 };
492
493 let uptime_seconds = self.started_at.elapsed().as_secs().max(1);
495 let frames_per_second = total_frames as f64 / uptime_seconds as f64;
496 let bytes_per_second = total_bytes as f64 / uptime_seconds as f64;
497
498 Ok(SystemStatsResponse {
499 total_sessions,
500 active_sessions,
501 total_streams,
502 active_streams,
503 total_frames,
504 total_bytes,
505 average_session_duration_seconds,
506 frames_per_second,
507 bytes_per_second,
508 uptime_seconds,
509 })
510 }
511 }
512}
513
514#[cfg(test)]
515mod tests {
516 use super::*;
517 use crate::domain::{
518 aggregates::{StreamSession, stream_session::SessionConfig},
519 ports::{PriorityDistribution, StreamFilter, StreamStatistics, StreamStatus, TimeProvider},
520 value_objects::{SessionId, StreamId},
521 };
522 use crate::test_support::MockRepository;
523 use chrono::Utc;
524
525 struct MockStreamStore;
526
527 impl StreamStoreGat for MockStreamStore {
528 type StoreStreamFuture<'a>
529 = impl std::future::Future<Output = crate::domain::DomainResult<()>> + Send + 'a
530 where
531 Self: 'a;
532
533 type GetStreamFuture<'a>
534 = impl std::future::Future<
535 Output = crate::domain::DomainResult<Option<crate::domain::entities::Stream>>,
536 > + Send
537 + 'a
538 where
539 Self: 'a;
540
541 type DeleteStreamFuture<'a>
542 = impl std::future::Future<Output = crate::domain::DomainResult<()>> + Send + 'a
543 where
544 Self: 'a;
545
546 type ListStreamsForSessionFuture<'a>
547 = impl std::future::Future<
548 Output = crate::domain::DomainResult<Vec<crate::domain::entities::Stream>>,
549 > + Send
550 + 'a
551 where
552 Self: 'a;
553
554 type FindStreamsBySessionFuture<'a>
555 = impl std::future::Future<
556 Output = crate::domain::DomainResult<Vec<crate::domain::entities::Stream>>,
557 > + Send
558 + 'a
559 where
560 Self: 'a;
561
562 type UpdateStreamStatusFuture<'a>
563 = impl std::future::Future<Output = crate::domain::DomainResult<()>> + Send + 'a
564 where
565 Self: 'a;
566
567 type GetStreamStatisticsFuture<'a>
568 = impl std::future::Future<Output = crate::domain::DomainResult<StreamStatistics>>
569 + Send
570 + 'a
571 where
572 Self: 'a;
573
574 fn store_stream(
575 &self,
576 _stream: crate::domain::entities::Stream,
577 ) -> Self::StoreStreamFuture<'_> {
578 async move { Ok(()) }
579 }
580
581 fn get_stream(&self, _stream_id: StreamId) -> Self::GetStreamFuture<'_> {
582 async move { Ok(None) }
583 }
584
585 fn delete_stream(&self, _stream_id: StreamId) -> Self::DeleteStreamFuture<'_> {
586 async move { Ok(()) }
587 }
588
589 fn list_streams_for_session(
590 &self,
591 _session_id: SessionId,
592 ) -> Self::ListStreamsForSessionFuture<'_> {
593 async move { Ok(vec![]) }
594 }
595
596 fn find_streams_by_session(
597 &self,
598 _session_id: SessionId,
599 _filter: StreamFilter,
600 ) -> Self::FindStreamsBySessionFuture<'_> {
601 async move { Ok(vec![]) }
602 }
603
604 fn update_stream_status(
605 &self,
606 _stream_id: StreamId,
607 _status: StreamStatus,
608 ) -> Self::UpdateStreamStatusFuture<'_> {
609 async move { Ok(()) }
610 }
611
612 fn get_stream_statistics(
613 &self,
614 _stream_id: StreamId,
615 ) -> Self::GetStreamStatisticsFuture<'_> {
616 async move {
617 Ok(StreamStatistics {
618 total_frames: 0,
619 total_bytes: 0,
620 priority_distribution: PriorityDistribution::default(),
621 avg_frame_size: 0.0,
622 creation_time: Utc::now(),
623 completion_time: None,
624 processing_duration: None,
625 })
626 }
627 }
628 }
629
630 #[tokio::test]
631 async fn test_get_session_query() {
632 let repository = Arc::new(MockRepository::new());
633 let handler = SessionQueryHandler::new(repository.clone());
634
635 let mut session = StreamSession::new(SessionConfig::default());
637 let _ = session.activate();
638 let session_id = session.id();
639 repository.add_session(session);
640
641 let query = GetSessionQuery {
643 session_id: session_id.into(),
644 };
645 let result = handler.handle(query).await;
646
647 assert!(result.is_ok());
648 let response = result.unwrap();
649 assert_eq!(response.session.id(), session_id);
650 }
651
652 #[tokio::test]
653 async fn test_get_session_not_found() {
654 let repository = Arc::new(MockRepository::new());
655 let handler = SessionQueryHandler::new(repository);
656
657 let query = GetSessionQuery {
658 session_id: SessionId::new().into(),
659 };
660 let result = handler.handle(query).await;
661
662 assert!(result.is_err());
663 match result.err().unwrap() {
664 ApplicationError::NotFound(_) => {}
665 _ => panic!("Expected NotFound error"),
666 }
667 }
668
669 #[tokio::test]
670 async fn test_get_active_sessions_query() {
671 let repository = Arc::new(MockRepository::new());
672 let handler = SessionQueryHandler::new(repository.clone());
673
674 for i in 0..5 {
676 let mut session = StreamSession::new(SessionConfig::default());
677 if i < 3 {
678 let _ = session.activate();
679 }
680 repository.add_session(session);
681 }
682
683 let query = GetActiveSessionsQuery {
685 offset: None,
686 limit: None,
687 };
688 let result = handler.handle(query).await;
689
690 assert!(result.is_ok());
691 let response = result.unwrap();
692 assert_eq!(response.sessions.len(), 5);
693 assert_eq!(response.total_count, 5);
694 }
695
696 #[tokio::test]
697 async fn test_get_active_sessions_with_pagination() {
698 let repository = Arc::new(MockRepository::new());
699 let handler = SessionQueryHandler::new(repository.clone());
700
701 for _ in 0..10 {
703 let mut session = StreamSession::new(SessionConfig::default());
704 let _ = session.activate();
705 repository.add_session(session);
706 }
707
708 let query = GetActiveSessionsQuery {
710 offset: Some(3),
711 limit: Some(4),
712 };
713 let result = handler.handle(query).await;
714
715 assert!(result.is_ok());
716 let response = result.unwrap();
717 assert_eq!(response.sessions.len(), 4);
718 assert_eq!(response.total_count, 10);
719 assert!(response.has_more);
720 }
721
722 #[tokio::test]
723 async fn test_get_active_sessions_last_page_has_more_false() {
724 let repository = Arc::new(MockRepository::new());
725 let handler = SessionQueryHandler::new(repository.clone());
726
727 for _ in 0..5 {
728 let mut session = StreamSession::new(SessionConfig::default());
729 let _ = session.activate();
730 repository.add_session(session);
731 }
732
733 let query = GetActiveSessionsQuery {
735 offset: Some(3),
736 limit: Some(4),
737 };
738 let response = handler.handle(query).await.unwrap();
739 assert_eq!(response.sessions.len(), 2);
740 assert!(!response.has_more);
741 }
742
743 #[tokio::test]
744 async fn test_get_active_sessions_page_cap() {
745 let repository = Arc::new(MockRepository::new());
746 let handler = SessionQueryHandler::new(repository.clone());
747
748 for _ in 0..110 {
749 let mut session = StreamSession::new(SessionConfig::default());
750 let _ = session.activate();
751 repository.add_session(session);
752 }
753
754 let query = GetActiveSessionsQuery {
756 offset: Some(0),
757 limit: Some(200),
758 };
759 let response = handler.handle(query).await.unwrap();
760 assert!(response.sessions.len() <= 100);
761 assert!(response.has_more);
762 }
763
764 #[tokio::test]
765 async fn test_get_session_health_query() {
766 let repository = Arc::new(MockRepository::new());
767 let handler = SessionQueryHandler::new(repository.clone());
768
769 let mut session = StreamSession::new(SessionConfig::default());
771 let _ = session.activate();
772 let session_id = session.id();
773 repository.add_session(session);
774
775 let query = GetSessionHealthQuery {
777 session_id: session_id.into(),
778 };
779 let result = handler.handle(query).await;
780
781 assert!(result.is_ok());
782 let response = result.unwrap();
783 assert!(response.health.is_healthy);
784 }
785
786 #[tokio::test]
787 async fn test_session_handler_creation() {
788 let repository = Arc::new(MockRepository::new());
789 let handler = SessionQueryHandler::new(repository.clone());
790
791 assert!(std::ptr::eq(
793 handler.repository.as_ref(),
794 repository.as_ref()
795 ));
796 }
797
798 #[tokio::test]
799 async fn test_stream_handler_creation() {
800 let session_repository = Arc::new(MockRepository::new());
801 let stream_store = Arc::new(MockStreamStore);
802 let frame_store = Arc::new(crate::infrastructure::adapters::InMemoryFrameStore::new());
803 let handler =
804 StreamQueryHandler::new(session_repository.clone(), stream_store, frame_store);
805
806 assert!(std::ptr::eq(
808 handler.session_repository.as_ref(),
809 session_repository.as_ref()
810 ));
811 }
812
813 #[tokio::test]
814 async fn test_system_handler_creation() {
815 let repository = Arc::new(MockRepository::new());
816 let handler = SystemQueryHandler::new(repository.clone());
817
818 assert!(std::ptr::eq(
820 handler.repository.as_ref(),
821 repository.as_ref()
822 ));
823 }
824
825 #[tokio::test]
826 async fn test_system_handler_real_uptime() {
827 use std::time::{Duration, Instant};
828
829 let repository = Arc::new(MockRepository::new());
830 let started_at = Instant::now() - Duration::from_secs(10);
832 let handler = SystemQueryHandler::with_start_time(repository, started_at);
833
834 let query = GetSystemStatsQuery {
835 include_historical: false,
836 };
837 let result = QueryHandlerGat::handle(&handler, query).await.unwrap();
838
839 assert!(
840 result.uptime_seconds >= 10,
841 "uptime_seconds should be at least 10, got {}",
842 result.uptime_seconds
843 );
844 }
845
846 #[tokio::test]
847 async fn test_get_stream_frames_session_not_found() {
848 use crate::domain::value_objects::{SessionId, StreamId};
849
850 let session_repository = Arc::new(MockRepository::new());
851 let stream_store = Arc::new(MockStreamStore);
852 let frame_store = Arc::new(crate::infrastructure::adapters::InMemoryFrameStore::new());
853 let handler = StreamQueryHandler::new(session_repository, stream_store, frame_store);
854
855 let query = GetStreamFramesQuery {
856 session_id: SessionId::new().into(),
857 stream_id: StreamId::new().into(),
858 since_sequence: None,
859 priority_filter: None,
860 limit: None,
861 };
862
863 let result: ApplicationResult<FramesResponse> =
864 QueryHandlerGat::handle(&handler, query).await;
865 assert!(matches!(result, Err(ApplicationError::NotFound(_))));
866 }
867
868 #[tokio::test]
869 async fn test_get_stream_frames_stream_not_found() {
870 use crate::domain::value_objects::StreamId;
871
872 let session_repository = Arc::new(MockRepository::new());
873 let mut session = StreamSession::new(SessionConfig::default());
874 let _ = session.activate();
875 let session_id = session.id();
876 session_repository.add_session(session);
877
878 let stream_store = Arc::new(MockStreamStore);
879 let frame_store = Arc::new(crate::infrastructure::adapters::InMemoryFrameStore::new());
880 let handler = StreamQueryHandler::new(session_repository, stream_store, frame_store);
881
882 let query = GetStreamFramesQuery {
883 session_id: session_id.into(),
884 stream_id: StreamId::new().into(),
885 since_sequence: None,
886 priority_filter: None,
887 limit: None,
888 };
889
890 let result: ApplicationResult<FramesResponse> =
891 QueryHandlerGat::handle(&handler, query).await;
892 assert!(matches!(result, Err(ApplicationError::NotFound(_))));
893 }
894
895 #[tokio::test]
896 async fn test_get_stream_frames_returns_empty() {
897 use crate::domain::value_objects::JsonData;
898
899 let session_repository = Arc::new(MockRepository::new());
900 let mut session = StreamSession::new(SessionConfig::default());
901 let _ = session.activate();
902 let session_id = session.id();
903 let stream_id = session
904 .create_stream(JsonData::from(serde_json::json!({"k": "v"})))
905 .unwrap();
906 session_repository.add_session(session);
907
908 let stream_store = Arc::new(MockStreamStore);
909 let frame_store = Arc::new(crate::infrastructure::adapters::InMemoryFrameStore::new());
910 let handler = StreamQueryHandler::new(session_repository, stream_store, frame_store);
911
912 let query = GetStreamFramesQuery {
913 session_id: session_id.into(),
914 stream_id: stream_id.into(),
915 since_sequence: None,
916 priority_filter: None,
917 limit: None,
918 };
919
920 let result = QueryHandlerGat::handle(&handler, query).await.unwrap();
921 assert_eq!(result.frames.len(), 0);
922 assert_eq!(result.total_count, 0);
923 }
924
925 #[tokio::test]
926 async fn test_get_stream_frames_returns_persisted_frames() {
927 use crate::domain::{
928 entities::frame::FramePatch,
929 ports::FrameStoreGat,
930 value_objects::{JsonData, JsonPath, Priority},
931 };
932 use crate::infrastructure::adapters::InMemoryFrameStore;
933
934 let session_repository = Arc::new(MockRepository::new());
935 let mut session = StreamSession::new(SessionConfig::default());
936 let _ = session.activate();
937 let session_id = session.id();
938 let stream_id = session
939 .create_stream(JsonData::from(serde_json::json!({"k": "v"})))
940 .unwrap();
941 session_repository.add_session(session);
942
943 let stream_store = Arc::new(MockStreamStore);
944 let frame_store = Arc::new(InMemoryFrameStore::new());
945
946 let frames = (1..=3)
948 .map(|seq| {
949 let patch = FramePatch::set(
950 JsonPath::new(format!("$.field_{seq}")).unwrap(),
951 JsonData::Integer(seq as i64),
952 );
953 crate::domain::entities::Frame::patch(stream_id, seq, Priority::HIGH, vec![patch])
954 .unwrap()
955 })
956 .collect();
957 frame_store.append_frames(stream_id, frames).await.unwrap();
958
959 let handler = StreamQueryHandler::new(session_repository, stream_store, frame_store);
960
961 let query = GetStreamFramesQuery {
962 session_id: session_id.into(),
963 stream_id: stream_id.into(),
964 since_sequence: None,
965 priority_filter: None,
966 limit: None,
967 };
968
969 let result = QueryHandlerGat::handle(&handler, query).await.unwrap();
970 assert_eq!(result.frames.len(), 3);
971 assert_eq!(result.total_count, 3);
972 assert_eq!(
973 result
974 .frames
975 .iter()
976 .map(crate::domain::entities::Frame::sequence)
977 .collect::<Vec<_>>(),
978 vec![1, 2, 3]
979 );
980 }
981
982 #[tokio::test]
983 async fn test_get_stream_frames_applies_filters() {
984 use crate::domain::{
985 entities::frame::FramePatch,
986 ports::FrameStoreGat,
987 value_objects::{JsonData, JsonPath, Priority},
988 };
989 use crate::infrastructure::adapters::InMemoryFrameStore;
990
991 let session_repository = Arc::new(MockRepository::new());
992 let mut session = StreamSession::new(SessionConfig::default());
993 let _ = session.activate();
994 let session_id = session.id();
995 let stream_id = session
996 .create_stream(JsonData::from(serde_json::json!({"k": "v"})))
997 .unwrap();
998 session_repository.add_session(session);
999
1000 let stream_store = Arc::new(MockStreamStore);
1001 let frame_store = Arc::new(InMemoryFrameStore::new());
1002
1003 let frames: Vec<_> = [
1005 (1u64, Priority::LOW),
1006 (2, Priority::HIGH),
1007 (3, Priority::CRITICAL),
1008 ]
1009 .into_iter()
1010 .map(|(seq, prio)| {
1011 let patch = FramePatch::set(
1012 JsonPath::new(format!("$.field_{seq}")).unwrap(),
1013 JsonData::Integer(seq as i64),
1014 );
1015 crate::domain::entities::Frame::patch(stream_id, seq, prio, vec![patch]).unwrap()
1016 })
1017 .collect();
1018 frame_store.append_frames(stream_id, frames).await.unwrap();
1019
1020 let handler = StreamQueryHandler::new(session_repository, stream_store, frame_store);
1021
1022 let query = GetStreamFramesQuery {
1024 session_id: session_id.into(),
1025 stream_id: stream_id.into(),
1026 since_sequence: Some(1),
1027 priority_filter: Some(Priority::HIGH.into()),
1028 limit: Some(10),
1029 };
1030
1031 let result = QueryHandlerGat::handle(&handler, query).await.unwrap();
1032 assert_eq!(result.frames.len(), 2);
1033 assert_eq!(result.total_count, 2);
1034 assert_eq!(
1035 result
1036 .frames
1037 .iter()
1038 .map(crate::domain::entities::Frame::sequence)
1039 .collect::<Vec<_>>(),
1040 vec![2, 3]
1041 );
1042 }
1043
1044 #[tokio::test]
1045 async fn test_get_stream_frames_caps_limit_to_max() {
1046 use crate::domain::{
1047 entities::frame::FramePatch,
1048 ports::FrameStoreGat,
1049 value_objects::{JsonData, JsonPath, Priority},
1050 };
1051 use crate::infrastructure::adapters::InMemoryFrameStore;
1052
1053 let session_repository = Arc::new(MockRepository::new());
1054 let mut session = StreamSession::new(SessionConfig::default());
1055 let _ = session.activate();
1056 let session_id = session.id();
1057 let stream_id = session
1058 .create_stream(JsonData::from(serde_json::json!({"k": "v"})))
1059 .unwrap();
1060 session_repository.add_session(session);
1061
1062 let frame_store = Arc::new(InMemoryFrameStore::new());
1063 let frames: Vec<_> = (1..=5)
1064 .map(|seq| {
1065 let patch = FramePatch::set(
1066 JsonPath::new(format!("$.field_{seq}")).unwrap(),
1067 JsonData::Integer(seq as i64),
1068 );
1069 crate::domain::entities::Frame::patch(stream_id, seq, Priority::HIGH, vec![patch])
1070 .unwrap()
1071 })
1072 .collect();
1073 frame_store.append_frames(stream_id, frames).await.unwrap();
1074
1075 let handler =
1076 StreamQueryHandler::new(session_repository, Arc::new(MockStreamStore), frame_store);
1077
1078 let query = GetStreamFramesQuery {
1080 session_id: session_id.into(),
1081 stream_id: stream_id.into(),
1082 since_sequence: None,
1083 priority_filter: None,
1084 limit: Some(999_999),
1085 };
1086
1087 let result = QueryHandlerGat::handle(&handler, query).await.unwrap();
1088 assert_eq!(result.frames.len(), 5);
1089 assert_eq!(result.total_count, 5);
1090 }
1091
1092 #[tokio::test]
1093 async fn test_get_session_stats_not_found() {
1094 use crate::domain::value_objects::SessionId;
1095
1096 let repository = Arc::new(MockRepository::new());
1097 let handler = SessionQueryHandler::new(repository);
1098
1099 let query = GetSessionStatsQuery {
1100 session_id: SessionId::new().into(),
1101 };
1102
1103 let result: ApplicationResult<SessionStatsResponse> =
1104 QueryHandlerGat::handle(&handler, query).await;
1105 assert!(matches!(result, Err(ApplicationError::NotFound(_))));
1106 }
1107
1108 #[tokio::test]
1109 async fn test_get_session_stats_returns_metadata() {
1110 use crate::domain::value_objects::JsonData;
1111
1112 let repository = Arc::new(MockRepository::new());
1113 let mut session = StreamSession::new(SessionConfig::default());
1114 let _ = session.activate();
1115 let session_id = session.id();
1116 let created_at = session.created_at();
1117 let _ = session.create_stream(JsonData::from(serde_json::json!({"a": 1})));
1119 let _ = session.create_stream(JsonData::from(serde_json::json!({"b": 2})));
1120 repository.add_session(session);
1121
1122 let handler = SessionQueryHandler::new(repository);
1123
1124 let query = GetSessionStatsQuery {
1125 session_id: session_id.into(),
1126 };
1127
1128 let result = QueryHandlerGat::handle(&handler, query).await.unwrap();
1129 assert_eq!(result.stream_count, 2);
1130 assert_eq!(result.created_at, created_at);
1131 }
1132
1133 #[tokio::test]
1136 async fn test_get_active_sessions_empty() {
1137 let repository = Arc::new(MockRepository::new());
1138 let handler = SessionQueryHandler::new(repository);
1139
1140 let query = GetActiveSessionsQuery {
1141 limit: Some(10),
1142 offset: Some(0),
1143 };
1144
1145 let result = QueryHandlerGat::handle(&handler, query).await;
1146 assert!(result.is_ok());
1147 let response = result.unwrap();
1148 assert_eq!(response.sessions.len(), 0);
1149 assert_eq!(response.total_count, 0);
1150 }
1151
1152 #[tokio::test]
1153 async fn test_get_active_sessions_with_limit() {
1154 let repository = Arc::new(MockRepository::new());
1155 let handler = SessionQueryHandler::new(repository);
1156
1157 let query = GetActiveSessionsQuery {
1158 limit: Some(5),
1159 offset: None,
1160 };
1161
1162 let result = QueryHandlerGat::handle(&handler, query).await;
1163 assert!(result.is_ok());
1164 }
1165
1166 #[tokio::test]
1167 async fn test_get_active_sessions_with_offset() {
1168 let repository = Arc::new(MockRepository::new());
1169 let handler = SessionQueryHandler::new(repository);
1170
1171 let query = GetActiveSessionsQuery {
1172 limit: None,
1173 offset: Some(10),
1174 };
1175
1176 let result = QueryHandlerGat::handle(&handler, query).await;
1177 assert!(result.is_ok());
1178 }
1179
1180 #[tokio::test]
1181 async fn test_get_active_sessions_offset_beyond_count() {
1182 let repository = Arc::new(MockRepository::new());
1183 let handler = SessionQueryHandler::new(repository);
1184
1185 let query = GetActiveSessionsQuery {
1186 limit: Some(10),
1187 offset: Some(1000),
1188 };
1189
1190 let result = QueryHandlerGat::handle(&handler, query).await;
1191 assert!(result.is_ok());
1192 let response = result.unwrap();
1193 assert_eq!(response.sessions.len(), 0);
1194 }
1195
1196 #[tokio::test]
1197 async fn test_get_stream_not_found() {
1198 use crate::domain::value_objects::{SessionId, StreamId};
1199
1200 let session_repository = Arc::new(MockRepository::new());
1201 let stream_store = Arc::new(MockStreamStore);
1202 let frame_store = Arc::new(crate::infrastructure::adapters::InMemoryFrameStore::new());
1203 let handler = StreamQueryHandler::new(session_repository, stream_store, frame_store);
1204
1205 let query = GetStreamQuery {
1206 session_id: SessionId::new().into(),
1207 stream_id: StreamId::new().into(),
1208 };
1209
1210 let result: ApplicationResult<StreamResponse> =
1211 QueryHandlerGat::handle(&handler, query).await;
1212 assert!(result.is_err());
1213 }
1214
1215 #[tokio::test]
1216 async fn test_session_query_not_found() {
1217 use crate::domain::value_objects::SessionId;
1218
1219 let repository = Arc::new(MockRepository::new());
1220 let handler = SessionQueryHandler::new(repository);
1221
1222 let query = GetSessionQuery {
1223 session_id: SessionId::new().into(),
1224 };
1225
1226 let result = QueryHandlerGat::handle(&handler, query).await;
1227 assert!(result.is_err());
1228 }
1229
1230 fn make_session(client_info: Option<&str>) -> StreamSession {
1232 let mut session = StreamSession::new(SessionConfig::default());
1233 let _ = session.activate();
1234 if let Some(info) = client_info {
1235 session.set_client_info(info.to_owned(), None, None);
1236 }
1237 session
1238 }
1239
1240 #[tokio::test]
1241 async fn test_client_info_filter_matching_passes() {
1242 let repository = Arc::new(MockRepository::new());
1243 repository.add_session(make_session(Some("Mozilla/5.0 (compatible; TestBot/1.0)")));
1244 let handler = SessionQueryHandler::new(repository);
1245
1246 let query = SearchSessionsQuery {
1248 filters: SessionFilters {
1249 client_info: Some("testbot".to_owned()),
1250 ..Default::default()
1251 },
1252 sort_by: None,
1253 sort_order: None,
1254 limit: None,
1255 offset: None,
1256 };
1257
1258 let result = QueryHandlerGat::handle(&handler, query).await;
1259 assert!(result.is_ok());
1260 let response = result.unwrap();
1261 assert_eq!(response.sessions.len(), 1);
1262 assert!(!response.has_more);
1263 }
1264
1265 #[tokio::test]
1266 async fn test_client_info_filter_non_matching_rejected() {
1267 let repository = Arc::new(MockRepository::new());
1268 repository.add_session(make_session(Some("Mozilla/5.0 (compatible; TestBot/1.0)")));
1269 let handler = SessionQueryHandler::new(repository);
1270
1271 let query = SearchSessionsQuery {
1272 filters: SessionFilters {
1273 client_info: Some("OtherClient".to_owned()),
1274 ..Default::default()
1275 },
1276 sort_by: None,
1277 sort_order: None,
1278 limit: None,
1279 offset: None,
1280 };
1281
1282 let result = QueryHandlerGat::handle(&handler, query).await;
1283 assert!(result.is_ok());
1284 let response = result.unwrap();
1285 assert_eq!(response.sessions.len(), 0);
1286 assert!(!response.has_more);
1287 }
1288
1289 #[tokio::test]
1290 async fn test_client_info_filter_no_info_rejected() {
1291 let repository = Arc::new(MockRepository::new());
1292 repository.add_session(make_session(None));
1293 let handler = SessionQueryHandler::new(repository);
1294
1295 let query = SearchSessionsQuery {
1296 filters: SessionFilters {
1297 client_info: Some("TestBot".to_owned()),
1298 ..Default::default()
1299 },
1300 sort_by: None,
1301 sort_order: None,
1302 limit: None,
1303 offset: None,
1304 };
1305
1306 let result = QueryHandlerGat::handle(&handler, query).await;
1307 assert!(result.is_ok());
1308 let response = result.unwrap();
1309 assert_eq!(response.sessions.len(), 0);
1310 assert!(!response.has_more);
1311 }
1312
1313 #[tokio::test]
1314 async fn test_client_info_filter_none_passes_all() {
1315 let repository = Arc::new(MockRepository::new());
1316 repository.add_session(make_session(Some("SomeAgent/2.0")));
1317 repository.add_session(make_session(None));
1318 let handler = SessionQueryHandler::new(repository);
1319
1320 let query = SearchSessionsQuery {
1321 filters: SessionFilters::default(),
1322 sort_by: None,
1323 sort_order: None,
1324 limit: None,
1325 offset: None,
1326 };
1327
1328 let result = QueryHandlerGat::handle(&handler, query).await;
1329 assert!(result.is_ok());
1330 let response = result.unwrap();
1331 assert_eq!(response.sessions.len(), 2);
1332 assert!(!response.has_more);
1333 }
1334
1335 #[tokio::test]
1336 async fn test_client_info_filter_case_insensitive() {
1337 let repository = Arc::new(MockRepository::new());
1338 repository.add_session(make_session(Some("Mozilla/5.0 (compatible; TESTBOT/2.0)")));
1339 let handler = SessionQueryHandler::new(repository);
1340
1341 let query = SearchSessionsQuery {
1343 filters: SessionFilters {
1344 client_info: Some("testbot".to_owned()),
1345 ..Default::default()
1346 },
1347 sort_by: None,
1348 sort_order: None,
1349 limit: None,
1350 offset: None,
1351 };
1352
1353 let result = QueryHandlerGat::handle(&handler, query).await;
1354 assert!(result.is_ok());
1355 let response = result.unwrap();
1356 assert_eq!(response.sessions.len(), 1);
1357 assert!(!response.has_more);
1358 }
1359
1360 #[tokio::test]
1361 async fn test_search_sessions_with_pagination_has_more_true() {
1362 let repository = Arc::new(MockRepository::new());
1363 for _ in 0..10 {
1364 repository.add_session(make_session(None));
1365 }
1366 let handler = SessionQueryHandler::new(repository);
1367
1368 let query = SearchSessionsQuery {
1369 filters: SessionFilters::default(),
1370 sort_by: None,
1371 sort_order: None,
1372 limit: Some(4),
1373 offset: Some(3),
1374 };
1375
1376 let result = QueryHandlerGat::handle(&handler, query).await;
1377 assert!(result.is_ok());
1378 let response = result.unwrap();
1379 assert_eq!(response.sessions.len(), 4);
1380 assert_eq!(response.total_count, 10);
1381 assert!(response.has_more);
1382 }
1383
1384 #[tokio::test]
1385 async fn test_search_sessions_last_page_has_more_false() {
1386 let repository = Arc::new(MockRepository::new());
1387 for _ in 0..7 {
1388 repository.add_session(make_session(None));
1389 }
1390 let handler = SessionQueryHandler::new(repository);
1391
1392 let query = SearchSessionsQuery {
1394 filters: SessionFilters::default(),
1395 sort_by: None,
1396 sort_order: None,
1397 limit: Some(4),
1398 offset: Some(3),
1399 };
1400
1401 let result = QueryHandlerGat::handle(&handler, query).await;
1402 assert!(result.is_ok());
1403 let response = result.unwrap();
1404 assert_eq!(response.sessions.len(), 4);
1405 assert_eq!(response.total_count, 7);
1406 assert!(!response.has_more);
1407 }
1408
1409 #[tokio::test]
1410 async fn test_search_sessions_filter_and_pagination_combined() {
1411 let repository = Arc::new(MockRepository::new());
1412 for _ in 0..8 {
1414 repository.add_session(make_session(None));
1415 }
1416 for _ in 0..3 {
1418 repository.add_session(StreamSession::new(SessionConfig::default()));
1419 }
1420 let handler = SessionQueryHandler::new(repository);
1421
1422 let query = SearchSessionsQuery {
1423 filters: SessionFilters {
1424 state: Some(SessionState::Active),
1425 ..Default::default()
1426 },
1427 sort_by: None,
1428 sort_order: None,
1429 limit: Some(3),
1430 offset: Some(2),
1431 };
1432
1433 let result = QueryHandlerGat::handle(&handler, query).await;
1434 assert!(result.is_ok());
1435 let response = result.unwrap();
1436 assert_eq!(response.sessions.len(), 3);
1437 assert_eq!(response.total_count, 8);
1438 assert!(response.has_more);
1439 }
1440
1441 #[tokio::test]
1442 async fn test_search_sessions_max_page_size_clamped() {
1443 let repository = Arc::new(MockRepository::new());
1444 for _ in 0..150 {
1445 repository.add_session(make_session(None));
1446 }
1447 let handler = SessionQueryHandler::new(repository);
1448
1449 let query = SearchSessionsQuery {
1451 filters: SessionFilters::default(),
1452 sort_by: None,
1453 sort_order: None,
1454 limit: Some(500),
1455 offset: None,
1456 };
1457
1458 let result = QueryHandlerGat::handle(&handler, query).await;
1459 assert!(result.is_ok());
1460 let response = result.unwrap();
1461 assert_eq!(response.sessions.len(), 100);
1462 assert_eq!(response.total_count, 150);
1463 assert!(response.has_more);
1464 }
1465
1466 #[tokio::test]
1467 async fn test_search_sessions_sort_by_stream_count_descending() {
1468 use crate::domain::value_objects::JsonData;
1469
1470 let repository = Arc::new(MockRepository::new());
1471
1472 let mut few = make_session(None);
1473 few.create_stream(JsonData::from(serde_json::json!({"k": "v"})))
1474 .unwrap();
1475 let few_id = few.id();
1476
1477 let mut many = make_session(None);
1478 for _ in 0..3 {
1479 many.create_stream(JsonData::from(serde_json::json!({"k": "v"})))
1480 .unwrap();
1481 }
1482 let many_id = many.id();
1483
1484 let none = make_session(None);
1485 let none_id = none.id();
1486
1487 repository.add_session(few);
1488 repository.add_session(many);
1489 repository.add_session(none);
1490 let handler = SessionQueryHandler::new(repository);
1491
1492 let query = SearchSessionsQuery {
1493 filters: SessionFilters::default(),
1494 sort_by: Some(SessionSortField::StreamCount),
1495 sort_order: Some(SortOrder::Descending),
1496 limit: None,
1497 offset: None,
1498 };
1499
1500 let result = QueryHandlerGat::handle(&handler, query).await;
1501 assert!(result.is_ok());
1502 let response = result.unwrap();
1503 let ids: Vec<_> = response.sessions.iter().map(|s| s.id()).collect();
1504 assert_eq!(ids, vec![many_id, few_id, none_id]);
1505 }
1506
1507 #[tokio::test]
1508 async fn test_search_sessions_sort_by_total_bytes_descending() {
1509 use crate::domain::value_objects::{JsonData, Priority};
1510
1511 let repository = Arc::new(MockRepository::new());
1512
1513 let mut heavy = make_session(None);
1515 let heavy_stream = heavy
1516 .create_stream(JsonData::String(
1517 "hello world payload, quite a few bytes here".to_owned(),
1518 ))
1519 .unwrap();
1520 heavy.start_stream(heavy_stream).unwrap();
1521 heavy
1522 .create_stream_patch_frames(heavy_stream, Priority::LOW, 100)
1523 .unwrap();
1524 let heavy_id = heavy.id();
1525
1526 let empty = make_session(None);
1528 let empty_id = empty.id();
1529
1530 assert!(heavy.stats().total_bytes > empty.stats().total_bytes);
1531
1532 repository.add_session(heavy);
1533 repository.add_session(empty);
1534 let handler = SessionQueryHandler::new(repository);
1535
1536 let query = SearchSessionsQuery {
1537 filters: SessionFilters::default(),
1538 sort_by: Some(SessionSortField::TotalBytes),
1539 sort_order: Some(SortOrder::Descending),
1540 limit: None,
1541 offset: None,
1542 };
1543
1544 let result = QueryHandlerGat::handle(&handler, query).await;
1545 assert!(result.is_ok());
1546 let response = result.unwrap();
1547 let ids: Vec<_> = response.sessions.iter().map(|s| s.id()).collect();
1548 assert_eq!(ids, vec![heavy_id, empty_id]);
1549 }
1550
1551 struct FixedTimeProvider {
1556 counter: std::sync::atomic::AtomicI64,
1557 }
1558
1559 impl FixedTimeProvider {
1560 fn new() -> Self {
1561 Self {
1562 counter: std::sync::atomic::AtomicI64::new(0),
1563 }
1564 }
1565 }
1566
1567 impl TimeProvider for FixedTimeProvider {
1568 fn now(&self) -> chrono::DateTime<Utc> {
1569 let offset = self
1570 .counter
1571 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1572 Utc::now() + chrono::Duration::seconds(offset)
1573 }
1574 }
1575
1576 #[tokio::test]
1577 async fn test_search_sessions_sort_by_created_at_descending() {
1578 let clock = std::sync::Arc::new(FixedTimeProvider::new());
1579 let repository = Arc::new(MockRepository::new());
1580
1581 let mut first = StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1582 first.activate().unwrap();
1583 let first_id = first.id();
1584
1585 let mut second = StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1586 second.activate().unwrap();
1587 let second_id = second.id();
1588
1589 repository.add_session(first);
1590 repository.add_session(second);
1591 let handler = SessionQueryHandler::new(repository);
1592
1593 let query = SearchSessionsQuery {
1594 filters: SessionFilters::default(),
1595 sort_by: Some(SessionSortField::CreatedAt),
1596 sort_order: Some(SortOrder::Descending),
1597 limit: None,
1598 offset: None,
1599 };
1600
1601 let result = QueryHandlerGat::handle(&handler, query).await;
1602 assert!(result.is_ok());
1603 let response = result.unwrap();
1604 let ids: Vec<_> = response.sessions.iter().map(|s| s.id()).collect();
1605 assert_eq!(ids, vec![second_id, first_id]);
1606 }
1607
1608 #[tokio::test]
1609 async fn test_search_sessions_sort_by_updated_at_descending() {
1610 use crate::domain::value_objects::JsonData;
1611
1612 let clock = std::sync::Arc::new(FixedTimeProvider::new());
1613 let repository = Arc::new(MockRepository::new());
1614
1615 let mut untouched =
1618 StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1619 untouched.activate().unwrap();
1620 let untouched_id = untouched.id();
1621
1622 let mut touched =
1623 StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1624 touched.activate().unwrap();
1625 touched
1626 .create_stream(JsonData::from(serde_json::json!({"k": "v"})))
1627 .unwrap();
1628 let touched_id = touched.id();
1629
1630 assert!(touched.updated_at() > untouched.updated_at());
1631
1632 repository.add_session(untouched);
1633 repository.add_session(touched);
1634 let handler = SessionQueryHandler::new(repository);
1635
1636 let query = SearchSessionsQuery {
1637 filters: SessionFilters::default(),
1638 sort_by: Some(SessionSortField::UpdatedAt),
1639 sort_order: Some(SortOrder::Descending),
1640 limit: None,
1641 offset: None,
1642 };
1643
1644 let result = QueryHandlerGat::handle(&handler, query).await;
1645 assert!(result.is_ok());
1646 let response = result.unwrap();
1647 let ids: Vec<_> = response.sessions.iter().map(|s| s.id()).collect();
1648 assert_eq!(ids, vec![touched_id, untouched_id]);
1649 }
1650
1651 #[tokio::test]
1652 async fn test_search_sessions_created_after_before_filter() {
1653 let clock = std::sync::Arc::new(FixedTimeProvider::new());
1654 let repository = Arc::new(MockRepository::new());
1655
1656 let mut early = StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1658 early.activate().unwrap();
1659
1660 let mut middle = StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1661 middle.activate().unwrap();
1662 let middle_id = middle.id();
1663 let middle_created_at = middle.created_at();
1664
1665 let mut late = StreamSession::with_time_provider(SessionConfig::default(), clock.clone());
1666 late.activate().unwrap();
1667
1668 repository.add_session(early);
1669 repository.add_session(middle);
1670 repository.add_session(late);
1671 let handler = SessionQueryHandler::new(repository);
1672
1673 let query = SearchSessionsQuery {
1674 filters: SessionFilters {
1675 created_after: Some(middle_created_at - chrono::Duration::milliseconds(1)),
1676 created_before: Some(middle_created_at + chrono::Duration::milliseconds(1)),
1677 ..Default::default()
1678 },
1679 sort_by: None,
1680 sort_order: None,
1681 limit: None,
1682 offset: None,
1683 };
1684
1685 let result = QueryHandlerGat::handle(&handler, query).await;
1686 assert!(result.is_ok());
1687 let response = result.unwrap();
1688 assert_eq!(response.sessions.len(), 1);
1689 assert_eq!(response.sessions[0].id(), middle_id);
1690 }
1691
1692 #[tokio::test]
1693 async fn test_search_sessions_has_active_streams_filter() {
1694 use crate::domain::value_objects::JsonData;
1695
1696 let repository = Arc::new(MockRepository::new());
1697
1698 let mut with_active = make_session(None);
1699 let stream_id = with_active
1700 .create_stream(JsonData::from(serde_json::json!({"k": "v"})))
1701 .unwrap();
1702 with_active.start_stream(stream_id).unwrap();
1703 let with_active_id = with_active.id();
1704
1705 let without_active = make_session(None);
1706 let without_active_id = without_active.id();
1707
1708 repository.add_session(with_active);
1709 repository.add_session(without_active);
1710 let handler = SessionQueryHandler::new(repository);
1711
1712 let query_active_only = SearchSessionsQuery {
1713 filters: SessionFilters {
1714 has_active_streams: Some(true),
1715 ..Default::default()
1716 },
1717 sort_by: None,
1718 sort_order: None,
1719 limit: None,
1720 offset: None,
1721 };
1722 let result = QueryHandlerGat::handle(&handler, query_active_only).await;
1723 assert!(result.is_ok());
1724 let response = result.unwrap();
1725 assert_eq!(response.sessions.len(), 1);
1726 assert_eq!(response.sessions[0].id(), with_active_id);
1727
1728 let query_inactive_only = SearchSessionsQuery {
1729 filters: SessionFilters {
1730 has_active_streams: Some(false),
1731 ..Default::default()
1732 },
1733 sort_by: None,
1734 sort_order: None,
1735 limit: None,
1736 offset: None,
1737 };
1738 let result = QueryHandlerGat::handle(&handler, query_inactive_only).await;
1739 assert!(result.is_ok());
1740 let response = result.unwrap();
1741 assert_eq!(response.sessions.len(), 1);
1742 assert_eq!(response.sessions[0].id(), without_active_id);
1743 }
1744
1745 #[tokio::test]
1755 async fn test_search_sessions_state_filter_matches_only_specified_variant() {
1756 let repository = Arc::new(MockRepository::new());
1757 let mut session = make_session(None);
1758 session.close().unwrap(); repository.add_session(session);
1760 let handler = SessionQueryHandler::new(repository);
1761
1762 let non_matching_query = SearchSessionsQuery {
1764 filters: SessionFilters {
1765 state: Some(SessionState::Failed),
1766 ..Default::default()
1767 },
1768 sort_by: None,
1769 sort_order: None,
1770 limit: None,
1771 offset: None,
1772 };
1773 let result = QueryHandlerGat::handle(&handler, non_matching_query).await;
1774 assert!(result.is_ok());
1775 assert_eq!(result.unwrap().sessions.len(), 0);
1776
1777 let exact_query = SearchSessionsQuery {
1779 filters: SessionFilters {
1780 state: Some(SessionState::Completed),
1781 ..Default::default()
1782 },
1783 sort_by: None,
1784 sort_order: None,
1785 limit: None,
1786 offset: None,
1787 };
1788 let result = QueryHandlerGat::handle(&handler, exact_query).await;
1789 assert!(result.is_ok());
1790 assert_eq!(result.unwrap().sessions.len(), 1);
1791 }
1792
1793 #[tokio::test]
1801 async fn test_search_sessions_default_scope_excludes_non_active_and_expired() {
1802 use crate::infrastructure::GatInMemoryStreamRepository;
1803
1804 let repository = Arc::new(GatInMemoryStreamRepository::new());
1805
1806 let mut active = StreamSession::new(SessionConfig::default());
1807 active.activate().unwrap();
1808 let active_id = active.id();
1809 repository.save_session(active).await.unwrap();
1810
1811 let mut completed = StreamSession::new(SessionConfig::default());
1812 completed.activate().unwrap();
1813 completed.close().unwrap();
1814 repository.save_session(completed).await.unwrap();
1815
1816 let mut expired = StreamSession::new(SessionConfig {
1817 session_timeout_seconds: 0,
1818 ..SessionConfig::default()
1819 });
1820 expired.activate().unwrap();
1821 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1823 repository.save_session(expired).await.unwrap();
1824
1825 let handler = SessionQueryHandler::new(repository);
1826
1827 let query = SearchSessionsQuery {
1828 filters: SessionFilters::default(),
1829 sort_by: None,
1830 sort_order: None,
1831 limit: None,
1832 offset: None,
1833 };
1834
1835 let result = QueryHandlerGat::handle(&handler, query).await;
1836 assert!(result.is_ok());
1837 let response = result.unwrap();
1838 assert_eq!(response.sessions.len(), 1);
1839 assert_eq!(response.sessions[0].id(), active_id);
1840 }
1841
1842 #[tokio::test]
1848 async fn test_search_sessions_limit_zero_and_excessive_offset_do_not_error() {
1849 use crate::infrastructure::GatInMemoryStreamRepository;
1850
1851 let repository = Arc::new(GatInMemoryStreamRepository::new());
1852 let mut session = StreamSession::new(SessionConfig::default());
1853 session.activate().unwrap();
1854 repository.save_session(session).await.unwrap();
1855 let handler = SessionQueryHandler::new(repository);
1856
1857 let query = SearchSessionsQuery {
1858 filters: SessionFilters::default(),
1859 sort_by: None,
1860 sort_order: None,
1861 limit: Some(0),
1862 offset: Some(usize::MAX),
1863 };
1864
1865 let result = QueryHandlerGat::handle(&handler, query).await;
1866 assert!(
1867 result.is_ok(),
1868 "expected clamped pagination, got {result:?}"
1869 );
1870 let response = result.unwrap();
1871 assert_eq!(response.total_count, 1);
1872 assert!(response.sessions.is_empty());
1873 }
1874}