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