Skip to main content

kmp_application/queries/
rehydrate_session.rs

1use std::sync::Arc;
2use std::time::SystemTime;
3
4use kmp_domain::{
5    BundleMetadata, GraphNeighborhoodReader, KmpBundle, NeighborhoodRequest, NodeDetailReader,
6    SnapshotSaveOptions, SnapshotStore,
7};
8
9use crate::ApplicationError;
10use crate::queries::{
11    ContextRenderOptions, NodeCentricProjectionReader, QueryApplicationService,
12    QueryTimingBreakdown,
13    context_render_options::EndpointHint,
14    render_graph_bundle::{RenderedContext, render_graph_bundle_with_options},
15};
16
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct RehydrateSessionQuery {
19    pub root_node_id: String,
20    pub roles: Vec<String>,
21    pub persist_snapshot: bool,
22    pub timeline_window: u32,
23    pub snapshot_ttl_seconds: u64,
24}
25
26#[derive(Debug, Clone, PartialEq)]
27pub struct RehydrateSessionResult {
28    pub root_node_id: String,
29    pub bundles: Vec<KmpBundle>,
30    /// Per-role rendered context with quality metrics, tiers, and truncation.
31    pub rendered_contexts: Vec<RenderedContext>,
32    pub timeline_events: u32,
33    pub version: BundleMetadata,
34    pub snapshot_persisted: bool,
35    pub snapshot_id: Option<String>,
36    pub generated_at: SystemTime,
37    pub timing: Option<QueryTimingBreakdown>,
38}
39
40#[derive(Debug)]
41pub struct RehydrateSessionUseCase<G, D, S> {
42    graph_reader: G,
43    detail_reader: D,
44    snapshot_store: S,
45    generator_version: &'static str,
46}
47
48impl<G, D, S> RehydrateSessionUseCase<G, D, S>
49where
50    G: GraphNeighborhoodReader + Send + Sync,
51    D: NodeDetailReader + Send + Sync,
52    S: SnapshotStore + Send + Sync,
53{
54    pub fn new(
55        graph_reader: G,
56        detail_reader: D,
57        snapshot_store: S,
58        generator_version: &'static str,
59    ) -> Self {
60        Self {
61            graph_reader,
62            detail_reader,
63            snapshot_store,
64            generator_version,
65        }
66    }
67
68    pub async fn execute(
69        &self,
70        root_node_id: &str,
71        role: &str,
72        persist_snapshot: bool,
73        snapshot_options: SnapshotSaveOptions,
74    ) -> Result<(KmpBundle, QueryTimingBreakdown), ApplicationError> {
75        self.execute_with_depth(
76            root_node_id,
77            role,
78            crate::queries::DEFAULT_NATIVE_GRAPH_TRAVERSAL_DEPTH,
79            persist_snapshot,
80            snapshot_options,
81        )
82        .await
83    }
84
85    pub async fn execute_with_depth(
86        &self,
87        root_node_id: &str,
88        role: &str,
89        depth: u32,
90        persist_snapshot: bool,
91        snapshot_options: SnapshotSaveOptions,
92    ) -> Result<(KmpBundle, QueryTimingBreakdown), ApplicationError> {
93        self.execute_for(
94            &NeighborhoodRequest::new(root_node_id, depth),
95            role,
96            persist_snapshot,
97            snapshot_options,
98        )
99        .await
100    }
101
102    /// Rehydrates for a request that already names the axes the caller
103    /// resolved, so the store can narrow instead of the caller discarding.
104    pub async fn execute_for(
105        &self,
106        request: &NeighborhoodRequest,
107        role: &str,
108        persist_snapshot: bool,
109        snapshot_options: SnapshotSaveOptions,
110    ) -> Result<(KmpBundle, QueryTimingBreakdown), ApplicationError> {
111        let root_node_id = request.root_node_id();
112        let bundle_reader =
113            NodeCentricProjectionReader::new(&self.graph_reader, &self.detail_reader);
114        let (bundle, timing) = match bundle_reader
115            .load_bundle_for(request, role, self.generator_version)
116            .await?
117        {
118            (Some(bundle), timing) => (bundle, timing),
119            (None, _) => {
120                return Err(ApplicationError::NotFound(format!(
121                    "node '{}' not found",
122                    root_node_id
123                )));
124            }
125        };
126
127        if persist_snapshot {
128            self.snapshot_store
129                .save_bundle_with_options(&bundle, snapshot_options)
130                .await?;
131        }
132        Ok((bundle, timing))
133    }
134}
135
136impl<G, D, S> QueryApplicationService<G, D, S>
137where
138    G: GraphNeighborhoodReader + Send + Sync,
139    D: NodeDetailReader + Send + Sync,
140    S: SnapshotStore + Send + Sync,
141{
142    pub async fn rehydrate_session(
143        &self,
144        query: RehydrateSessionQuery,
145    ) -> Result<RehydrateSessionResult, ApplicationError> {
146        if query.roles.is_empty() {
147            return Err(ApplicationError::Validation(
148                "roles cannot be empty".to_string(),
149            ));
150        }
151
152        let bundle_reader = NodeCentricProjectionReader::new(
153            self.graph_reader.as_ref(),
154            self.detail_reader.as_ref(),
155        );
156        let snapshot_options = SnapshotSaveOptions::new(Some(query.snapshot_ttl_seconds));
157
158        let (bundles, timing) = bundle_reader
159            .load_bundles_for_roles(
160                &query.root_node_id,
161                &query.roles,
162                self.generator_version,
163                crate::queries::DEFAULT_NATIVE_GRAPH_TRAVERSAL_DEPTH,
164            )
165            .await?;
166
167        let bundles = match bundles {
168            Some(bundles) => bundles,
169            None => {
170                return Err(ApplicationError::NotFound(format!(
171                    "node '{}' not found",
172                    query.root_node_id
173                )));
174            }
175        };
176
177        let session_options = ContextRenderOptions {
178            endpoint_hint: EndpointHint::SessionSnapshot,
179            ..Default::default()
180        };
181        let rendered_contexts: Vec<RenderedContext> = bundles
182            .iter()
183            .map(|b| render_graph_bundle_with_options(b, &session_options))
184            .collect();
185
186        if query.persist_snapshot {
187            for bundle in &bundles {
188                self.snapshot_store
189                    .save_bundle_with_options(bundle, snapshot_options)
190                    .await?;
191            }
192        }
193
194        let snapshot_id = if query.persist_snapshot {
195            Some(format!(
196                "snapshot:{}:{}",
197                query.root_node_id,
198                query.roles.join(",")
199            ))
200        } else {
201            None
202        };
203
204        Ok(RehydrateSessionResult {
205            root_node_id: query.root_node_id,
206            bundles,
207            rendered_contexts,
208            timeline_events: query.timeline_window,
209            version: BundleMetadata::initial(self.generator_version),
210            snapshot_persisted: query.persist_snapshot,
211            snapshot_id,
212            generated_at: SystemTime::now(),
213            timing: Some(timing),
214        })
215    }
216
217    pub async fn warmup_bundle(&self) -> Result<KmpBundle, ApplicationError> {
218        let (bundle, _timing) = RehydrateSessionUseCase::new(
219            Arc::clone(&self.graph_reader),
220            Arc::clone(&self.detail_reader),
221            Arc::clone(&self.snapshot_store),
222            self.generator_version,
223        )
224        .execute(
225            "bootstrap-node",
226            "system",
227            false,
228            SnapshotSaveOptions::default(),
229        )
230        .await?;
231        Ok(bundle)
232    }
233}
234
235#[cfg(test)]
236mod tests {
237    use std::collections::BTreeMap;
238    use std::sync::Arc;
239
240    use tokio::sync::Mutex;
241
242    use kmp_domain::{
243        ContextPathNeighborhood, NodeDetailProjection, NodeNeighborhood, NodeProjection, PortError,
244        SnapshotSaveOptions,
245    };
246
247    use super::{QueryApplicationService, RehydrateSessionQuery};
248
249    struct SeededGraphReader;
250
251    impl kmp_domain::GraphNeighborhoodReader for SeededGraphReader {
252        async fn load_neighborhood(
253            &self,
254            root_node_id: &str,
255            _depth: u32,
256        ) -> Result<Option<NodeNeighborhood>, PortError> {
257            Ok(Some(NodeNeighborhood {
258                root: NodeProjection {
259                    node_id: root_node_id.to_string(),
260                    node_kind: "story".to_string(),
261                    title: "Root".to_string(),
262                    summary: "Root summary".to_string(),
263                    status: "ACTIVE".to_string(),
264                    labels: vec!["Story".to_string()],
265                    properties: BTreeMap::new(),
266                    provenance: None,
267                },
268                neighbors: Vec::new(),
269                relations: Vec::new(),
270            }))
271        }
272
273        async fn load_context_path(
274            &self,
275            _root_node_id: &str,
276            _target_node_id: &str,
277            _subtree_depth: u32,
278        ) -> Result<Option<ContextPathNeighborhood>, PortError> {
279            Ok(None)
280        }
281    }
282
283    struct SeededDetailReader;
284
285    impl kmp_domain::NodeDetailReader for SeededDetailReader {
286        async fn load_node_detail(
287            &self,
288            node_id: &str,
289        ) -> Result<Option<NodeDetailProjection>, PortError> {
290            Ok(Some(NodeDetailProjection {
291                node_id: node_id.to_string(),
292                detail: "Expanded detail".to_string(),
293                content_hash: "hash-1".to_string(),
294                revision: 2,
295            }))
296        }
297
298        async fn load_node_details_batch(
299            &self,
300            node_ids: Vec<String>,
301        ) -> Result<Vec<Option<NodeDetailProjection>>, PortError> {
302            let mut results = Vec::with_capacity(node_ids.len());
303            for node_id in &node_ids {
304                results.push(self.load_node_detail(node_id).await?);
305            }
306            Ok(results)
307        }
308    }
309
310    #[derive(Debug, Default)]
311    struct RecordingSnapshotStore {
312        options: Mutex<Vec<SnapshotSaveOptions>>,
313    }
314
315    impl kmp_domain::SnapshotStore for RecordingSnapshotStore {
316        async fn save_bundle_with_options(
317            &self,
318            _bundle: &kmp_domain::KmpBundle,
319            options: SnapshotSaveOptions,
320        ) -> Result<(), PortError> {
321            self.options.lock().await.push(options);
322            Ok(())
323        }
324    }
325
326    #[tokio::test]
327    async fn rehydrate_session_propagates_snapshot_ttl_to_store() {
328        let snapshot_store = Arc::new(RecordingSnapshotStore::default());
329        let service = QueryApplicationService::new(
330            Arc::new(SeededGraphReader),
331            Arc::new(SeededDetailReader),
332            Arc::clone(&snapshot_store),
333            "0.1.0",
334        );
335
336        let result = service
337            .rehydrate_session(RehydrateSessionQuery {
338                root_node_id: "story-123".to_string(),
339                roles: vec!["developer".to_string()],
340                persist_snapshot: true,
341                timeline_window: 50,
342                snapshot_ttl_seconds: 1800,
343            })
344            .await
345            .expect("rehydration should succeed");
346
347        assert!(result.snapshot_persisted);
348        assert_eq!(
349            snapshot_store.options.lock().await.as_slice(),
350            &[SnapshotSaveOptions::new(Some(1800))]
351        );
352    }
353
354    struct EmptyGraphReader;
355
356    impl kmp_domain::GraphNeighborhoodReader for EmptyGraphReader {
357        async fn load_neighborhood(
358            &self,
359            _root_node_id: &str,
360            _depth: u32,
361        ) -> Result<Option<NodeNeighborhood>, PortError> {
362            Ok(None)
363        }
364
365        async fn load_context_path(
366            &self,
367            _root_node_id: &str,
368            _target_node_id: &str,
369            _subtree_depth: u32,
370        ) -> Result<Option<ContextPathNeighborhood>, PortError> {
371            Ok(None)
372        }
373    }
374
375    #[tokio::test]
376    async fn rehydrate_session_returns_not_found_when_node_does_not_exist() {
377        let service = QueryApplicationService::new(
378            Arc::new(EmptyGraphReader),
379            Arc::new(SeededDetailReader),
380            Arc::new(RecordingSnapshotStore::default()),
381            "0.1.0",
382        );
383
384        let result = service
385            .rehydrate_session(RehydrateSessionQuery {
386                root_node_id: "nonexistent-node".to_string(),
387                roles: vec!["developer".to_string()],
388                persist_snapshot: false,
389                timeline_window: 0,
390                snapshot_ttl_seconds: 0,
391            })
392            .await;
393
394        match result {
395            Err(crate::ApplicationError::NotFound(_)) => {}
396            other => panic!("expected NotFound, got: {other:?}"),
397        }
398    }
399
400    #[tokio::test]
401    async fn rehydrate_session_multi_role_returns_one_bundle_per_role() {
402        let service = QueryApplicationService::new(
403            Arc::new(SeededGraphReader),
404            Arc::new(SeededDetailReader),
405            Arc::new(RecordingSnapshotStore::default()),
406            "0.1.0",
407        );
408
409        let result = service
410            .rehydrate_session(RehydrateSessionQuery {
411                root_node_id: "story-123".to_string(),
412                roles: vec![
413                    "developer".to_string(),
414                    "reviewer".to_string(),
415                    "ops".to_string(),
416                ],
417                persist_snapshot: false,
418                timeline_window: 0,
419                snapshot_ttl_seconds: 0,
420            })
421            .await
422            .expect("multi-role rehydration should succeed");
423
424        assert_eq!(result.bundles.len(), 3);
425        assert_eq!(result.bundles[0].role().as_str(), "developer");
426        assert_eq!(result.bundles[1].role().as_str(), "reviewer");
427        assert_eq!(result.bundles[2].role().as_str(), "ops");
428        // All bundles share the same graph — root and details are identical.
429        for bundle in &result.bundles {
430            assert_eq!(bundle.root_node().node_id(), "story-123");
431            assert_eq!(bundle.node_details().len(), 1);
432        }
433    }
434
435    #[tokio::test]
436    async fn rehydrate_session_multi_role_not_found_returns_error_not_partial() {
437        let service = QueryApplicationService::new(
438            Arc::new(EmptyGraphReader),
439            Arc::new(SeededDetailReader),
440            Arc::new(RecordingSnapshotStore::default()),
441            "0.1.0",
442        );
443
444        let result = service
445            .rehydrate_session(RehydrateSessionQuery {
446                root_node_id: "nonexistent".to_string(),
447                roles: vec!["developer".to_string(), "reviewer".to_string()],
448                persist_snapshot: false,
449                timeline_window: 0,
450                snapshot_ttl_seconds: 0,
451            })
452            .await;
453
454        match result {
455            Err(crate::ApplicationError::NotFound(_)) => {}
456            other => panic!("expected NotFound for all roles, got: {other:?}"),
457        }
458    }
459
460    #[tokio::test]
461    async fn rehydrate_session_multi_role_persists_snapshot_per_bundle() {
462        let snapshot_store = Arc::new(RecordingSnapshotStore::default());
463        let service = QueryApplicationService::new(
464            Arc::new(SeededGraphReader),
465            Arc::new(SeededDetailReader),
466            Arc::clone(&snapshot_store),
467            "0.1.0",
468        );
469
470        let result = service
471            .rehydrate_session(RehydrateSessionQuery {
472                root_node_id: "story-123".to_string(),
473                roles: vec!["developer".to_string(), "reviewer".to_string()],
474                persist_snapshot: true,
475                timeline_window: 0,
476                snapshot_ttl_seconds: 900,
477            })
478            .await
479            .expect("multi-role rehydration should succeed");
480
481        assert!(result.snapshot_persisted);
482        assert_eq!(
483            result.snapshot_id,
484            Some("snapshot:story-123:developer,reviewer".to_string())
485        );
486        // One save per role bundle.
487        assert_eq!(snapshot_store.options.lock().await.len(), 2);
488    }
489
490    #[tokio::test]
491    async fn rehydrate_session_empty_roles_returns_validation_error() {
492        let service = QueryApplicationService::new(
493            Arc::new(SeededGraphReader),
494            Arc::new(SeededDetailReader),
495            Arc::new(RecordingSnapshotStore::default()),
496            "0.1.0",
497        );
498
499        let result = service
500            .rehydrate_session(RehydrateSessionQuery {
501                root_node_id: "story-123".to_string(),
502                roles: vec![],
503                persist_snapshot: false,
504                timeline_window: 0,
505                snapshot_ttl_seconds: 0,
506            })
507            .await;
508
509        match result {
510            Err(crate::ApplicationError::Validation(msg)) => {
511                assert!(msg.contains("roles"), "error should mention roles: {msg}");
512            }
513            other => panic!("expected Validation error, got: {other:?}"),
514        }
515    }
516}