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 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 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 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}