Skip to main content

locus_core_rs/application/services/
monthly_rollup_service.rs

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    /// Create a monthly rollup service with storage and validation dependencies.
26    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    /// Build a monthly rollup node from source nodes in the requested date range.
53    ///
54    /// Depending on request settings, this can run as preview-only or persist
55    /// the generated node into the configured store.
56    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__", &timestamp)
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}