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