1use std::collections::HashSet;
2use std::sync::Arc;
3
4use chrono::Utc;
5
6use crate::domain::contracts::{
7 EmbeddingProvider, NodeStore, NodeValidator, SemanticIndexStore, TagEmbedding,
8};
9use crate::domain::models::{
10 AvecState, ConfidenceBandSummary, MonthlyRollupRequest, MonthlyRollupResult, NodeQuery,
11 NumericRange, SemanticTagNodeRef,
12};
13use crate::parsing::SttpNodeParser;
14use crate::storage::derive_tenant_id_from_session;
15
16pub struct MonthlyRollupService {
17 store: Arc<dyn NodeStore>,
18 validator: Arc<dyn NodeValidator>,
19 parser: SttpNodeParser,
20 semantic_index: Option<Arc<dyn SemanticIndexStore>>,
21 embedding_provider: Option<Arc<dyn EmbeddingProvider>>,
22}
23
24impl MonthlyRollupService {
25 pub fn new(store: Arc<dyn NodeStore>, validator: Arc<dyn NodeValidator>) -> Self {
27 Self {
28 store,
29 validator,
30 parser: SttpNodeParser::new(),
31 semantic_index: None,
32 embedding_provider: None,
33 }
34 }
35
36 pub fn with_semantic_index(
37 mut self,
38 semantic_index: Arc<dyn SemanticIndexStore>,
39 ) -> Self {
40 self.semantic_index = Some(semantic_index);
41 self
42 }
43
44 pub fn with_embedding_provider(
45 mut self,
46 embedding_provider: Arc<dyn EmbeddingProvider>,
47 ) -> Self {
48 self.embedding_provider = Some(embedding_provider);
49 self
50 }
51
52 pub async fn create_async(&self, request: MonthlyRollupRequest) -> MonthlyRollupResult {
57 emit_rollup_trace(
58 &request.session_id,
59 "create_async_start",
60 &format!(
61 "start_utc={} end_utc={} source_session_id={} limit={} persist={}",
62 request.start_utc.to_rfc3339(),
63 request.end_utc.to_rfc3339(),
64 request.source_session_id.as_deref().unwrap_or("all"),
65 request.limit,
66 request.persist
67 ),
68 );
69
70 if request.end_utc < request.start_utc {
71 emit_rollup_trace(
72 &request.session_id,
73 "invalid_range",
74 "end_utc precedes start_utc",
75 );
76 return MonthlyRollupResult {
77 error: Some(
78 "InvalidRange: end must be greater than or equal to start.".to_string(),
79 ),
80 ..MonthlyRollupResult::default()
81 };
82 }
83
84 let nodes = match self
85 .store
86 .query_nodes_async(NodeQuery {
87 session_id: request.source_session_id.clone(),
88 from_utc: Some(request.start_utc),
89 to_utc: Some(request.end_utc),
90 limit: request.limit,
91 tiers: None,
92 })
93 .await
94 {
95 Ok(nodes) => nodes,
96 Err(err) => {
97 emit_rollup_trace(
98 &request.session_id,
99 "query_failure",
100 &format!("error={} content_redacted=true", err),
101 );
102 return MonthlyRollupResult {
103 error: Some(format!("QueryFailure: {err}")),
104 ..MonthlyRollupResult::default()
105 };
106 }
107 };
108
109 emit_rollup_trace(
110 &request.session_id,
111 "query_result",
112 &format!(
113 "source_nodes={} source_session_id={} window={}..{}",
114 nodes.len(),
115 request.source_session_id.as_deref().unwrap_or("all"),
116 request.start_utc.to_rfc3339(),
117 request.end_utc.to_rfc3339()
118 ),
119 );
120
121 if nodes.is_empty() {
122 emit_rollup_trace(
123 &request.session_id,
124 "no_source_nodes",
125 "query returned zero nodes; verify upstream store acceptance and filters",
126 );
127 return MonthlyRollupResult {
128 error: Some("NoSourceNodes: no nodes found in the requested range.".to_string()),
129 ..MonthlyRollupResult::default()
130 };
131 }
132
133 let mut ordered_nodes = nodes;
134 ordered_nodes.sort_by(|a, b| a.timestamp.cmp(&b.timestamp));
135
136 let user_nodes = ordered_nodes
137 .iter()
138 .filter(|n| n.user_avec.psi() > 0.0)
139 .collect::<Vec<_>>();
140 let model_nodes = ordered_nodes
141 .iter()
142 .filter(|n| n.model_avec.psi() > 0.0)
143 .collect::<Vec<_>>();
144 let compression_nodes = ordered_nodes
145 .iter()
146 .filter(|n| n.compression_avec.map(|avec| avec.psi()).unwrap_or(0.0) > 0.0)
147 .collect::<Vec<_>>();
148
149 let user_average = average_avec(user_nodes.iter().map(|n| n.user_avec));
150 let model_average = average_avec(model_nodes.iter().map(|n| n.model_avec));
151 let compression_average =
152 average_avec(compression_nodes.iter().filter_map(|n| n.compression_avec));
153
154 let rho_range = range_for(ordered_nodes.iter().map(|n| n.rho));
155 let kappa_range = range_for(ordered_nodes.iter().map(|n| n.kappa));
156 let psi_range = range_for(ordered_nodes.iter().map(|n| n.psi));
157 let rho_bands = bands_for(ordered_nodes.iter().map(|n| n.rho));
158 let kappa_bands = bands_for(ordered_nodes.iter().map(|n| n.kappa));
159 let active_days = ordered_nodes
160 .iter()
161 .map(|n| n.timestamp.date_naive())
162 .collect::<HashSet<_>>()
163 .len();
164 let parent_reference = request
165 .parent_node_id
166 .clone()
167 .unwrap_or_else(|| ordered_nodes[0].session_id.clone());
168
169 let rollup_tags = aggregate_semantic_tags(&ordered_nodes);
170
171 let raw_node = build_monthly_node(
172 &request,
173 &parent_reference,
174 ordered_nodes.len(),
175 user_nodes.len(),
176 active_days,
177 user_average,
178 model_average,
179 compression_average,
180 rho_range,
181 kappa_range,
182 psi_range,
183 rho_bands,
184 kappa_bands,
185 &rollup_tags,
186 );
187
188 let validation = self.validator.validate(&raw_node);
189 if !validation.is_valid {
190 emit_rollup_trace(
191 &request.session_id,
192 "validation_failure",
193 &format!(
194 "reason={} error={} content_redacted=true",
195 validation.reason,
196 validation.error.as_deref().unwrap_or_default()
197 ),
198 );
199 return MonthlyRollupResult {
200 raw_node,
201 source_nodes: ordered_nodes.len(),
202 parent_reference: Some(parent_reference),
203 error: Some(format!(
204 "ValidationFailure: {}: {}",
205 validation.reason,
206 validation.error.unwrap_or_default()
207 )),
208 ..MonthlyRollupResult::default()
209 };
210 }
211
212 let parse_result = self.parser.try_parse(&raw_node, &request.session_id);
213 if !parse_result.success {
214 emit_rollup_trace(
215 &request.session_id,
216 "parse_failure",
217 &format!(
218 "profile={:?} strict_valid={} diagnostics_count={} error={} content_redacted=true",
219 parse_result.profile,
220 parse_result.strict_valid,
221 parse_result.diagnostics.len(),
222 parse_result.error.as_deref().unwrap_or_default()
223 ),
224 );
225 return MonthlyRollupResult {
226 raw_node,
227 source_nodes: ordered_nodes.len(),
228 parent_reference: Some(parent_reference),
229 error: Some(format!(
230 "ParseFailure: {}",
231 parse_result.error.unwrap_or_default()
232 )),
233 ..MonthlyRollupResult::default()
234 };
235 }
236
237 let mut node_id = String::new();
238 if request.persist {
239 let Some(parsed_node) = parse_result.node else {
240 emit_rollup_trace(
241 &request.session_id,
242 "parse_failure",
243 "missing parsed node after successful parse result",
244 );
245 return MonthlyRollupResult {
246 raw_node,
247 source_nodes: ordered_nodes.len(),
248 parent_reference: Some(parent_reference),
249 error: Some("ParseFailure: missing parsed node".to_string()),
250 ..MonthlyRollupResult::default()
251 };
252 };
253
254 let parsed_snapshot = parsed_node.clone();
255
256 match self.store.upsert_node_async(parsed_node).await {
257 Ok(upsert) => {
258 if let Err(err) = self
259 .sync_semantic_tags_async(&parsed_snapshot, &upsert.node_id, &upsert.sync_key)
260 .await
261 {
262 emit_rollup_trace(
263 &request.session_id,
264 "semantic_index_sync_failure",
265 &format!("error={} content_redacted=true", err),
266 );
267 }
268
269 emit_rollup_trace(
270 &request.session_id,
271 "store_success",
272 &format!(
273 "node_id={} sync_key={} source_nodes={} content_redacted=true",
274 upsert.node_id, upsert.sync_key, ordered_nodes.len()
275 ),
276 );
277 node_id = upsert.node_id
278 }
279 Err(err) => {
280 emit_rollup_trace(
281 &request.session_id,
282 "store_failure",
283 &format!("error={} content_redacted=true", err),
284 );
285 return MonthlyRollupResult {
286 raw_node,
287 source_nodes: ordered_nodes.len(),
288 parent_reference: Some(parent_reference),
289 user_average,
290 model_average,
291 compression_average,
292 rho_range,
293 kappa_range,
294 psi_range,
295 rho_bands,
296 kappa_bands,
297 error: Some(format!("StoreFailure: {err}")),
298 ..MonthlyRollupResult::default()
299 };
300 }
301 }
302 }
303
304 emit_rollup_trace(
305 &request.session_id,
306 "create_async_complete",
307 &format!(
308 "success=true source_nodes={} parent_reference={} persist={}",
309 ordered_nodes.len(),
310 parent_reference,
311 request.persist
312 ),
313 );
314
315 MonthlyRollupResult {
316 success: true,
317 node_id,
318 raw_node,
319 source_nodes: ordered_nodes.len(),
320 parent_reference: Some(parent_reference),
321 user_average,
322 model_average,
323 compression_average,
324 rho_range,
325 kappa_range,
326 psi_range,
327 rho_bands,
328 kappa_bands,
329 ..MonthlyRollupResult::default()
330 }
331 }
332
333 async fn sync_semantic_tags_async(
334 &self,
335 parsed: &crate::domain::models::SttpNode,
336 node_id: &str,
337 sync_key: &str,
338 ) -> anyhow::Result<()> {
339 let Some(index) = self.semantic_index.as_ref() else {
340 return Ok(());
341 };
342
343 sync_semantic_tags_for_rollup(
344 index,
345 self.embedding_provider.as_ref(),
346 parsed,
347 node_id,
348 sync_key,
349 )
350 .await
351 }
352}
353
354async fn sync_semantic_tags_for_rollup(
355 semantic_index: &Arc<dyn SemanticIndexStore>,
356 embedding_provider: Option<&Arc<dyn EmbeddingProvider>>,
357 parsed: &crate::domain::models::SttpNode,
358 node_id: &str,
359 sync_key: &str,
360) -> anyhow::Result<()> {
361 use std::collections::HashMap;
362
363 let tags = parsed.semantic_tags.clone().unwrap_or_default();
364 let tenant_id = derive_tenant_id_from_session(&parsed.session_id);
365 let node_ref = SemanticTagNodeRef {
366 tenant_id,
367 session_id: parsed.session_id.clone(),
368 node_id: node_id.to_string(),
369 sync_key: sync_key.to_string(),
370 };
371
372 let embeddings = if let Some(provider) = embedding_provider {
373 let mut map = HashMap::new();
374 for tag in &tags {
375 let canonical = tag.trim().to_lowercase();
376 if canonical.is_empty() {
377 continue;
378 }
379 if let Ok(vector) = provider.embed_async(tag).await {
380 map.insert(
381 canonical,
382 TagEmbedding {
383 vector,
384 model: provider.model_name().to_string(),
385 },
386 );
387 }
388 }
389 Some(map)
390 } else {
391 None
392 };
393
394 semantic_index
395 .sync_node_tags_async(node_ref, &tags, embeddings.as_ref())
396 .await
397}
398
399fn emit_rollup_trace(session_id: &str, event: &str, detail: &str) {
400 eprintln!(
401 "[sttp_rollup_trace] session_id={} event={} detail={}",
402 session_id, event, detail
403 );
404}
405
406#[allow(clippy::too_many_arguments)]
407fn build_monthly_node(
408 request: &MonthlyRollupRequest,
409 parent_reference: &str,
410 source_nodes: usize,
411 source_user_avec_nodes: usize,
412 active_days: usize,
413 user_average: AvecState,
414 model_average: AvecState,
415 compression_average: AvecState,
416 rho_range: NumericRange,
417 kappa_range: NumericRange,
418 psi_range: NumericRange,
419 rho_bands: ConfidenceBandSummary,
420 kappa_bands: ConfidenceBandSummary,
421 rollup_tags: &[String],
422) -> String {
423 let timestamp = Utc::now().to_rfc3339();
424 let start = request.start_utc.format("%Y-%m-%d").to_string();
425 let end = request.end_utc.format("%Y-%m-%d").to_string();
426 let source_session_token = request
427 .source_session_id
428 .as_deref()
429 .filter(|s| !s.trim().is_empty())
430 .map(slug)
431 .unwrap_or_else(|| "all_sessions".to_string());
432
433 let template = r#"⊕⟨ ⏣0{ trigger: manual, response_format: temporal_node, origin_session: "__SESSION_ID__", compression_depth: 2, parent_node: ref:__PARENT_REFERENCE__, prime: { attractor_config: { stability: __USER_STABILITY__, friction: __USER_FRICTION__, logic: __USER_LOGIC__, autonomy: __USER_AUTONOMY__ }, context_summary: monthly_rollup_across_stored_sttp_nodes_with_average_state_and_confidence_spread, relevant_tier: monthly, retrieval_budget: 16__SEMANTIC_TAGS_BLOCK__ } } ⟩
434⦿⟨ ⏣0{ timestamp: "__TIMESTAMP__", tier: monthly, session_id: "__SESSION_ID__", schema_version: "sttp-1.0", user_avec: { stability: __USER_STABILITY__, friction: __USER_FRICTION__, logic: __USER_LOGIC__, autonomy: __USER_AUTONOMY__, psi: __USER_PSI__ }, model_avec: { stability: __MODEL_STABILITY__, friction: __MODEL_FRICTION__, logic: __MODEL_LOGIC__, autonomy: __MODEL_AUTONOMY__, psi: __MODEL_PSI__ } } ⟩
435◈⟨ ⏣0{ source_nodes(.99): __SOURCE_NODES__, source_user_avec_nodes(.97): __SOURCE_USER_AVEC_NODES__, active_days(.95): __ACTIVE_DAYS__, date_span(.99): __START___to___END__, source_session_filter(.78): __SOURCE_SESSION_TOKEN__, parent_anchor(.99): __PARENT_ANCHOR__, activity_shape(.83): burst_work_pattern_with_gaps_between_deep_sessions, monthly_arc(.86): stabilization_then_design_then_implementation_then_synthesis, behavioral_signature(.84): high_stability_high_logic_high_autonomy_with_low_to_moderate_friction, user_avec_average(.99): { stability: __USER_STABILITY__, friction: __USER_FRICTION__, logic: __USER_LOGIC__, autonomy: __USER_AUTONOMY__, psi: __USER_PSI__ }, model_avec_average(.97): { stability: __MODEL_STABILITY__, friction: __MODEL_FRICTION__, logic: __MODEL_LOGIC__, autonomy: __MODEL_AUTONOMY__, psi: __MODEL_PSI__ }, compression_avec_average(.96): { stability: __COMP_STABILITY__, friction: __COMP_FRICTION__, logic: __COMP_LOGIC__, autonomy: __COMP_AUTONOMY__, psi: __COMP_PSI__ }, confidence_ranges(.94): { rho_avg: __RHO_AVG__, rho_min: __RHO_MIN__, rho_max: __RHO_MAX__, kappa_avg: __KAPPA_AVG__, kappa_min: __KAPPA_MIN__, kappa_max: __KAPPA_MAX__, psi_avg: __PSI_AVG__, psi_min: __PSI_MIN__, psi_max: __PSI_MAX__ }, confidence_bands(.71): { rho_low: __RHO_LOW__, rho_medium: __RHO_MEDIUM__, rho_high: __RHO_HIGH__, kappa_low: __KAPPA_LOW__, kappa_medium: __KAPPA_MEDIUM__, kappa_high: __KAPPA_HIGH__ }, uncertainty(.41): interpretive_fields_carry_lower_confidence_than_numeric_rollups } ⟩
436⍉⟨ ⏣0{ rho: __RHO_AVG__, kappa: __KAPPA_AVG__, psi: __PSI_AVG__, compression_avec: { stability: __COMP_STABILITY__, friction: __COMP_FRICTION__, logic: __COMP_LOGIC__, autonomy: __COMP_AUTONOMY__, psi: __COMP_PSI__ } } ⟩"#;
437
438 template
439 .replace("__SESSION_ID__", &request.session_id)
440 .replace("__PARENT_REFERENCE__", parent_reference)
441 .replace("__TIMESTAMP__", ×tamp)
442 .replace("__SOURCE_NODES__", &source_nodes.to_string())
443 .replace(
444 "__SOURCE_USER_AVEC_NODES__",
445 &source_user_avec_nodes.to_string(),
446 )
447 .replace("__ACTIVE_DAYS__", &active_days.to_string())
448 .replace("__START__", &start)
449 .replace("__END__", &end)
450 .replace("__SOURCE_SESSION_TOKEN__", &source_session_token)
451 .replace("__PARENT_ANCHOR__", &slug(parent_reference))
452 .replace("__USER_STABILITY__", &format_float(user_average.stability))
453 .replace("__USER_FRICTION__", &format_float(user_average.friction))
454 .replace("__USER_LOGIC__", &format_float(user_average.logic))
455 .replace("__USER_AUTONOMY__", &format_float(user_average.autonomy))
456 .replace("__USER_PSI__", &format_float(user_average.psi()))
457 .replace(
458 "__MODEL_STABILITY__",
459 &format_float(model_average.stability),
460 )
461 .replace("__MODEL_FRICTION__", &format_float(model_average.friction))
462 .replace("__MODEL_LOGIC__", &format_float(model_average.logic))
463 .replace("__MODEL_AUTONOMY__", &format_float(model_average.autonomy))
464 .replace("__MODEL_PSI__", &format_float(model_average.psi()))
465 .replace(
466 "__COMP_STABILITY__",
467 &format_float(compression_average.stability),
468 )
469 .replace(
470 "__COMP_FRICTION__",
471 &format_float(compression_average.friction),
472 )
473 .replace("__COMP_LOGIC__", &format_float(compression_average.logic))
474 .replace(
475 "__COMP_AUTONOMY__",
476 &format_float(compression_average.autonomy),
477 )
478 .replace("__COMP_PSI__", &format_float(compression_average.psi()))
479 .replace("__RHO_AVG__", &format_float(rho_range.average))
480 .replace("__RHO_MIN__", &format_float(rho_range.min))
481 .replace("__RHO_MAX__", &format_float(rho_range.max))
482 .replace("__KAPPA_AVG__", &format_float(kappa_range.average))
483 .replace("__KAPPA_MIN__", &format_float(kappa_range.min))
484 .replace("__KAPPA_MAX__", &format_float(kappa_range.max))
485 .replace("__PSI_AVG__", &format_float(psi_range.average))
486 .replace("__PSI_MIN__", &format_float(psi_range.min))
487 .replace("__PSI_MAX__", &format_float(psi_range.max))
488 .replace("__RHO_LOW__", &rho_bands.low.to_string())
489 .replace("__RHO_MEDIUM__", &rho_bands.medium.to_string())
490 .replace("__RHO_HIGH__", &rho_bands.high.to_string())
491 .replace("__KAPPA_LOW__", &kappa_bands.low.to_string())
492 .replace("__KAPPA_MEDIUM__", &kappa_bands.medium.to_string())
493 .replace("__KAPPA_HIGH__", &kappa_bands.high.to_string())
494 .replace(
495 "__SEMANTIC_TAGS_BLOCK__",
496 &format_semantic_tags_block(rollup_tags),
497 )
498}
499
500fn aggregate_semantic_tags(nodes: &[crate::domain::models::SttpNode]) -> Vec<String> {
501 let mut tags = HashSet::new();
502 for node in nodes {
503 if let Some(node_tags) = &node.semantic_tags {
504 for tag in node_tags {
505 tags.insert(tag.clone());
506 }
507 }
508 }
509
510 let mut sorted = tags.into_iter().collect::<Vec<_>>();
511 sorted.sort();
512 sorted
513}
514
515fn format_semantic_tags_block(tags: &[String]) -> String {
516 if tags.is_empty() {
517 return String::new();
518 }
519
520 let formatted = tags
521 .iter()
522 .map(|tag| format!("\"{}\"", tag.replace('"', "\\\"")))
523 .collect::<Vec<_>>()
524 .join(", ");
525
526 format!(", semantic_tags: [{formatted}]")
527}
528
529fn average_avec<I>(states: I) -> AvecState
530where
531 I: IntoIterator<Item = AvecState>,
532{
533 let values = states.into_iter().collect::<Vec<_>>();
534 if values.is_empty() {
535 return AvecState::zero();
536 }
537
538 let len = values.len() as f32;
539 let stability = values.iter().map(|s| s.stability).sum::<f32>() / len;
540 let friction = values.iter().map(|s| s.friction).sum::<f32>() / len;
541 let logic = values.iter().map(|s| s.logic).sum::<f32>() / len;
542 let autonomy = values.iter().map(|s| s.autonomy).sum::<f32>() / len;
543
544 AvecState {
545 stability,
546 friction,
547 logic,
548 autonomy,
549 }
550}
551
552fn range_for<I>(values: I) -> NumericRange
553where
554 I: IntoIterator<Item = f32>,
555{
556 let values = values.into_iter().collect::<Vec<_>>();
557 if values.is_empty() {
558 return NumericRange::default();
559 }
560
561 let min = values
562 .iter()
563 .fold(f32::INFINITY, |acc, value| acc.min(*value));
564 let max = values
565 .iter()
566 .fold(f32::NEG_INFINITY, |acc, value| acc.max(*value));
567 let average = values.iter().sum::<f32>() / values.len() as f32;
568
569 NumericRange { min, max, average }
570}
571
572fn bands_for<I>(values: I) -> ConfidenceBandSummary
573where
574 I: IntoIterator<Item = f32>,
575{
576 let values = values.into_iter().collect::<Vec<_>>();
577 ConfidenceBandSummary {
578 low: values.iter().filter(|v| **v < 0.5).count(),
579 medium: values.iter().filter(|v| **v >= 0.5 && **v < 0.85).count(),
580 high: values.iter().filter(|v| **v >= 0.85).count(),
581 }
582}
583
584fn format_float(value: f32) -> String {
585 let mut s = format!("{value:.10}");
586 while s.contains('.') && s.ends_with('0') {
587 s.pop();
588 }
589 if s.ends_with('.') {
590 s.pop();
591 }
592 if s.is_empty() { "0".to_string() } else { s }
593}
594
595fn slug(value: &str) -> String {
596 let mut output = String::with_capacity(value.len());
597 for ch in value.chars() {
598 if ch.is_alphanumeric() {
599 for lower in ch.to_lowercase() {
600 output.push(lower);
601 }
602 } else {
603 output.push('_');
604 }
605 }
606
607 output.trim_matches('_').to_string()
608}