Skip to main content

weavatrix_memory/analytics/
consolidation.rs

1use super::{ConsolidationAction, ConsolidationKind, ConsolidationPlan, MemoryAnalytics};
2use crate::{
3    domain::MemoryFact,
4    error::{MemoryError, Result},
5    graph_projection::project_graph,
6    id::EntityId,
7    projection::{MemoryProjection, ProjectionClock},
8};
9use std::collections::BTreeMap;
10
11impl MemoryAnalytics {
12    /// Produces a deterministic, non-mutating maintenance plan.
13    ///
14    /// Event history is never deleted. The caller may translate proposed
15    /// duplicate actions into explicit supersession events.
16    ///
17    /// # Errors
18    ///
19    /// Rejects a zero action limit or graph projection failures.
20    pub fn consolidation_plan(
21        projection: &MemoryProjection,
22        clock: ProjectionClock,
23        max_actions: usize,
24    ) -> Result<ConsolidationPlan> {
25        if max_actions == 0 {
26            return Err(MemoryError::InvalidValue {
27                field: "consolidation.max_actions",
28                reason: "must be greater than zero",
29            });
30        }
31        let view = projection.view(clock);
32        let graph = project_graph(&view)?;
33        let mut actions = duplicate_actions(&view.facts);
34        actions.extend(orphan_actions(&graph));
35        actions.extend(revision_actions(projection, clock));
36        actions.sort_by(|left, right| {
37            left.kind
38                .cmp(&right.kind)
39                .then_with(|| left.affected_entities.cmp(&right.affected_entities))
40                .then_with(|| left.affected_facts.cmp(&right.affected_facts))
41        });
42        actions.truncate(max_actions);
43        let projected_savings = actions
44            .iter()
45            .map(|action| match action.kind {
46                ConsolidationKind::ReviewOrphan => 0,
47                _ => action.affected_facts.len().saturating_sub(1),
48            })
49            .sum();
50        Ok(ConsolidationPlan {
51            actions,
52            projected_savings,
53            source_position: projection.last_global_position(),
54        })
55    }
56}
57
58fn duplicate_actions(facts: &[MemoryFact]) -> Vec<ConsolidationAction> {
59    let mut groups = BTreeMap::<(EntityId, String, EntityId), Vec<&MemoryFact>>::new();
60    for fact in facts {
61        groups
62            .entry((
63                fact.source.clone(),
64                fact.relation.clone(),
65                fact.target.clone(),
66            ))
67            .or_default()
68            .push(fact);
69    }
70    let mut actions = Vec::new();
71    for ((source, _, target), mut duplicates) in groups {
72        if duplicates.len() < 2 {
73            continue;
74        }
75        duplicates.sort_by(|left, right| {
76            right
77                .confidence
78                .cmp(&left.confidence)
79                .then_with(|| right.recorded_at.cmp(&left.recorded_at))
80                .then_with(|| left.id.cmp(&right.id))
81        });
82        actions.push(ConsolidationAction {
83            kind: ConsolidationKind::SupersedeDuplicate,
84            keep: Some(duplicates[0].id.clone()),
85            affected_facts: duplicates.iter().map(|fact| fact.id.clone()).collect(),
86            affected_entities: vec![source, target],
87            rationale: "same active source, relation, and target; keep strongest evidence"
88                .to_owned(),
89        });
90    }
91    actions
92}
93
94fn orphan_actions(graph: &weavatrix_graph::Graph) -> Vec<ConsolidationAction> {
95    graph
96        .nodes()
97        .iter()
98        .filter_map(|node| {
99            let index = graph.node_index(node.id.as_str())?;
100            let isolated = graph.in_degree(index) == Some(0) && graph.out_degree(index) == Some(0);
101            isolated.then(|| ConsolidationAction {
102                kind: ConsolidationKind::ReviewOrphan,
103                keep: None,
104                affected_facts: Vec::new(),
105                affected_entities: EntityId::new(node.id.as_str()).into_iter().collect(),
106                rationale: "entity has no active incoming or outgoing evidence".to_owned(),
107            })
108        })
109        .collect()
110}
111
112fn revision_actions(
113    projection: &MemoryProjection,
114    clock: ProjectionClock,
115) -> Vec<ConsolidationAction> {
116    let mut groups = BTreeMap::<(EntityId, String), Vec<&MemoryFact>>::new();
117    for fact in projection
118        .all_facts()
119        .iter()
120        .filter(|fact| fact.recorded_at <= clock.known_at)
121    {
122        groups
123            .entry((fact.source.clone(), fact.relation.clone()))
124            .or_default()
125            .push(fact);
126    }
127    groups
128        .into_iter()
129        .filter_map(|((entity, _), mut facts)| {
130            facts.sort_by_key(|fact| (fact.recorded_at, fact.id.clone()));
131            let chain = facts
132                .iter()
133                .filter(|fact| fact.supersedes.is_some())
134                .count();
135            (chain >= 3).then(|| ConsolidationAction {
136                kind: ConsolidationKind::CompactRevisionChain,
137                keep: facts.last().map(|fact| fact.id.clone()),
138                affected_facts: facts.iter().map(|fact| fact.id.clone()).collect(),
139                affected_entities: vec![entity],
140                rationale: "retain event history but checkpoint a long revision chain".to_owned(),
141            })
142        })
143        .collect()
144}