Skip to main content

pjson_rs/application/handlers/
query_handlers.rs

1//! Query handlers for read operations
2
3use 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
16/// Hard cap on frames returned by `GET /streams/{id}/frames`, matching the
17/// shared pagination ceiling. Defends the HTTP layer against a client that
18/// passes an oversized `limit`.
19const MAX_FRAMES_PAGE_SIZE: usize = crate::domain::config::MAX_PAGINATION_LIMIT;
20
21/// Handler for session-related queries
22#[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    /// Construct a handler that reads sessions from `repository`.
35    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            // With no explicit state filter, default to the same active-only
156            // scope the legacy `find_active_sessions()` provided (state ==
157            // Active and not expired), so the endpoint's default result set
158            // is unchanged. An explicit `filters.state` opts out of that
159            // default and matches only the requested state(s), expired or not.
160            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/// Handler for stream-related queries
242#[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    /// Construct a system-level query handler from the supplied dependencies.
261    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            // Validate stream exists within the session.
367            session.stream(query.stream_id.into()).ok_or_else(|| {
368                ApplicationError::NotFound(format!("Stream {} not found", query.stream_id))
369            })?;
370
371            // Cap requested page size to MAX_FRAMES_PAGE_SIZE so a client cannot
372            // ask for an unbounded response. Honor None as "give me everything
373            // up to the cap".
374            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/// Handler for system statistics
406#[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    /// Create a new handler, recording `Instant::now()` as the startup time.
420    pub fn new(repository: Arc<R>) -> Self {
421        Self {
422            repository,
423            started_at: Instant::now(),
424        }
425    }
426
427    /// Create a handler with an explicit startup instant.
428    ///
429    /// Useful when multiple handlers share a single process-start time.
430    ///
431    /// # Examples
432    ///
433    /// ```ignore
434    /// let started_at = std::time::Instant::now();
435    /// let handler = SystemQueryHandler::with_start_time(repo, started_at);
436    /// ```
437    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            // Floor to 1 to avoid divide-by-zero when the query runs immediately on startup.
494            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        // Create and add a session
636        let mut session = StreamSession::new(SessionConfig::default());
637        let _ = session.activate();
638        let session_id = session.id();
639        repository.add_session(session);
640
641        // Query the session
642        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        // Add multiple sessions
675        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        // Query active sessions
684        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        // Add 10 sessions
702        for _ in 0..10 {
703            let mut session = StreamSession::new(SessionConfig::default());
704            let _ = session.activate();
705            repository.add_session(session);
706        }
707
708        // Query with pagination
709        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        // offset=3, limit=4 → only 2 remain → last page
734        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        // limit=200 must be capped to 100
755        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        // Create and add a session
770        let mut session = StreamSession::new(SessionConfig::default());
771        let _ = session.activate();
772        let session_id = session.id();
773        repository.add_session(session);
774
775        // Query session health
776        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        // Test that handlers can be created successfully
792        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        // Test that handlers can be created successfully
807        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        // Test that handlers can be created successfully
819        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        // Simulate a handler that started 10 seconds ago.
831        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        // Simulate the command handler having appended frames.
947        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        // Three frames with mixed priorities and sequences.
1004        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        // since_sequence=1 + priority>=HIGH should leave seq 2 and 3.
1023        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        // limit=999_999 must be capped to MAX_FRAMES_PAGE_SIZE without erroring.
1079        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        // Add two streams so we can assert stream_count.
1118        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    // ===== Additional Query Handler Tests for CQ-003 (Coverage Improvement) =====
1134
1135    #[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    // Helper to build an active session with optional client_info
1231    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        // Use mixed-case filter to verify case-insensitive matching.
1247        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        // Filter uses lowercase while session value is uppercase — must still match.
1342        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        // offset=3, limit=4 → page of 4 exactly reaches total_count=7 → full last page
1393        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        // 8 sessions matching the "active" state filter.
1413        for _ in 0..8 {
1414            repository.add_session(make_session(None));
1415        }
1416        // 3 sessions in a non-matching state (never activated).
1417        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        // Requested limit exceeds MAX_PAGE_SIZE (100); the handler must clamp it.
1450        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        // Session with a real byte-carrying patch frame batch.
1514        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        // Session with no streams, so total_bytes stays zero.
1527        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    /// Deterministic clock letting tests control `created_at`/`updated_at`
1552    /// ordering without relying on wall-clock sleeps. Each call to `now()`
1553    /// advances the clock so successive sessions/mutations get distinct,
1554    /// strictly increasing timestamps.
1555    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        // Both sessions created back-to-back; only `touched`'s later mutation
1616        // advances its `updated_at` past `untouched`'s.
1617        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        // Three sessions created at strictly increasing timestamps (t=0,1,2).
1657        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    /// Locks in the exact-match state filter semantics that replaced the old
1746    /// in-process substring match (see `matches_filters`, removed in #391).
1747    /// `SessionFilters::state` is now typed as `SessionState` (#414), so a
1748    /// partial word like "complet" is rejected at compile time rather than
1749    /// silently substring-matching a "Completed" session; only a fellow
1750    /// non-matching variant can be exercised as the negative case here.
1751    /// The HTTP-boundary rejection of an actual malformed/partial `state`
1752    /// string (e.g. `?state=complet`) is covered separately by
1753    /// `axum_adapter::tests::search_sessions_route_rejects_unknown_state`.
1754    #[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(); // Active -> Completed
1759        repository.add_session(session);
1760        let handler = SessionQueryHandler::new(repository);
1761
1762        // Non-matching variant: must not match the Completed session.
1763        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        // Matching variant: matches exactly.
1778        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    /// Regression test for the search endpoint's default scope: with no
1794    /// explicit `filters.state`, results must stay limited to active,
1795    /// non-expired sessions, matching the legacy `find_active_sessions()`
1796    /// contract (`state == Active && !is_expired()`). Uses the real
1797    /// `GatInMemoryStreamRepository` (not `MockRepository`) so the criteria
1798    /// actually flow through `matches_criteria`, the same path `GET
1799    /// /pjs/sessions/search` exercises in production.
1800    #[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        // Give the real clock a moment to move past `expires_at`.
1822        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    /// Regression test: `limit=0` and an offset beyond `MAX_PAGINATION_OFFSET`
1843    /// must not surface `SessionPagination::validate()`'s `DomainError::InvalidInput`
1844    /// (which the HTTP layer maps to a 500). The handler must clamp both
1845    /// before they reach the repository. Uses the real
1846    /// `GatInMemoryStreamRepository` so `SessionPagination::validate()` actually runs.
1847    #[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}