Skip to main content

weavatrix_memory/projection/
memory.rs

1mod binary;
2mod index;
3mod parts;
4mod state;
5
6use super::Projection;
7use crate::{
8    EntityId, FactId, MemoryError, MemoryEvent, MemoryFact, MemoryNode, MemoryView, Result,
9    StoredEvent, Timestamp,
10};
11use serde::{Deserialize, Serialize};
12use state::{NodeHistory, NodeRevision, Retraction, Supersession};
13use std::collections::BTreeSet;
14
15pub use binary::CompactSnapshotCodec;
16pub use state::MemoryProjection;
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
19pub struct ProjectionClock {
20    pub valid_at: Timestamp,
21    pub known_at: Timestamp,
22}
23
24impl ProjectionClock {
25    #[must_use]
26    pub const fn new(valid_at: Timestamp, known_at: Timestamp) -> Self {
27        Self { valid_at, known_at }
28    }
29}
30
31impl MemoryProjection {
32    #[must_use]
33    pub const fn last_global_position(&self) -> Option<u64> {
34        self.last_global_position
35    }
36
37    #[must_use]
38    pub fn superseded_by(&self, fact: &FactId) -> Option<&FactId> {
39        self.supersessions
40            .get(fact)
41            .map(|change| &change.replacement)
42    }
43
44    pub(crate) fn visible_node(&self, id: &EntityId, known_at: Timestamp) -> Option<&MemoryNode> {
45        self.node_lookup
46            .get(id)
47            .and_then(|index| visible_revision(&self.nodes[*index], known_at))
48    }
49
50    pub(crate) fn fact(&self, id: &FactId) -> Option<&MemoryFact> {
51        self.fact_lookup.get(id).map(|index| &self.facts[*index])
52    }
53
54    pub(crate) fn all_facts(&self) -> &[MemoryFact] {
55        &self.facts
56    }
57
58    pub(crate) fn incident_fact_ids(&self, id: &EntityId) -> impl Iterator<Item = &FactId> + '_ {
59        self.node_lookup
60            .get(id)
61            .into_iter()
62            .flat_map(|index| {
63                let stable = &self.incident_facts
64                    [self.incident_offsets[*index]..self.incident_offsets[*index + 1]];
65                stable
66                    .iter()
67                    .chain(self.incident_delta.get(index).into_iter().flatten())
68            })
69            .map(|index| &self.facts[*index].id)
70    }
71
72    #[must_use]
73    pub fn view(&self, clock: ProjectionClock) -> MemoryView {
74        let nodes = self
75            .nodes
76            .iter()
77            .filter_map(|revisions| visible_revision(revisions, clock.known_at))
78            .cloned()
79            .collect::<Vec<_>>();
80        let visible_ids = nodes
81            .iter()
82            .map(|node| node.id.clone())
83            .collect::<BTreeSet<_>>();
84        let facts = self
85            .facts
86            .iter()
87            .filter(|fact| {
88                visible_ids.contains(&fact.source)
89                    && visible_ids.contains(&fact.target)
90                    && self.fact_is_active(fact, clock)
91            })
92            .cloned()
93            .collect();
94        MemoryView { nodes, facts }
95    }
96
97    pub(crate) fn fact_is_active(&self, fact: &MemoryFact, clock: ProjectionClock) -> bool {
98        if fact.recorded_at > clock.known_at || fact.valid_from > clock.valid_at {
99            return false;
100        }
101        if fact
102            .valid_until
103            .is_some_and(|until| clock.valid_at >= until)
104        {
105            return false;
106        }
107        if self.supersessions.get(&fact.id).is_some_and(|change| {
108            change.recorded_at <= clock.known_at && change.valid_from <= clock.valid_at
109        }) {
110            return false;
111        }
112        !self.retractions.get(&fact.id).is_some_and(|change| {
113            change.recorded_at <= clock.known_at && change.valid_until <= clock.valid_at
114        })
115    }
116
117    fn insert_node(&mut self, revision: NodeRevision) -> Result<()> {
118        revision.node.validate()?;
119        if let Some(index) = self.node_lookup.get(&revision.node.id).copied() {
120            self.nodes[index].later.push(revision);
121        } else {
122            let index = self.nodes.len();
123            self.node_lookup.insert(revision.node.id.clone(), index);
124            self.nodes.push(NodeHistory::new(revision));
125            let offset = *self.incident_offsets.last().unwrap_or(&0);
126            self.incident_offsets.push(offset);
127        }
128        Ok(())
129    }
130
131    fn insert_fact(&mut self, fact: MemoryFact) -> Result<()> {
132        fact.validate()?;
133        let source = self.require_entity(&fact.source)?;
134        let target = self.require_entity(&fact.target)?;
135        if self.fact_lookup.contains_key(&fact.id) {
136            return Err(MemoryError::ConflictingFact {
137                id: fact.id.to_string(),
138            });
139        }
140        if let Some(prior) = &fact.supersedes {
141            self.apply_supersession(prior, &fact)?;
142        }
143        let index = self.facts.len();
144        self.fact_lookup.insert(fact.id.clone(), index);
145        self.incident_delta.entry(source).or_default().push(index);
146        if target != source {
147            self.incident_delta.entry(target).or_default().push(index);
148        }
149        self.facts.push(fact);
150        Ok(())
151    }
152
153    fn apply_supersession(&mut self, prior: &FactId, fact: &MemoryFact) -> Result<()> {
154        if !self.fact_lookup.contains_key(prior) {
155            return Err(MemoryError::MissingFact {
156                id: prior.to_string(),
157            });
158        }
159        if self.supersessions.contains_key(prior) {
160            return Err(MemoryError::ConflictingFact {
161                id: prior.to_string(),
162            });
163        }
164        self.supersessions.insert(
165            prior.clone(),
166            Supersession {
167                replacement: fact.id.clone(),
168                valid_from: fact.valid_from,
169                recorded_at: fact.recorded_at,
170            },
171        );
172        Ok(())
173    }
174
175    fn require_entity(&self, id: &EntityId) -> Result<usize> {
176        self.node_lookup
177            .get(id)
178            .copied()
179            .ok_or_else(|| MemoryError::MissingEntity { id: id.to_string() })
180    }
181
182    fn apply_retraction(
183        &mut self,
184        event: &StoredEvent<MemoryEvent>,
185        fact_id: &FactId,
186        valid_until: Timestamp,
187        evidence: &[crate::Evidence],
188    ) -> Result<()> {
189        let fact = self.fact(fact_id).ok_or_else(|| MemoryError::MissingFact {
190            id: fact_id.to_string(),
191        })?;
192        if valid_until <= fact.valid_from {
193            return Err(MemoryError::InvalidValue {
194                field: "retraction.valid_until",
195                reason: "must be later than the fact valid_from",
196            });
197        }
198        if evidence.is_empty() {
199            return Err(MemoryError::InvalidValue {
200                field: "retraction.evidence",
201                reason: "at least one evidence item is required",
202            });
203        }
204        evidence.iter().try_for_each(crate::Evidence::validate)?;
205        self.retractions.insert(
206            fact_id.clone(),
207            Retraction {
208                valid_until,
209                recorded_at: event.metadata.recorded_at,
210            },
211        );
212        Ok(())
213    }
214}
215
216impl Projection<MemoryEvent> for MemoryProjection {
217    fn apply(&mut self, event: &StoredEvent<MemoryEvent>) -> Result<()> {
218        if event.metadata.event_type != event.payload.event_type() {
219            return Err(MemoryError::InvalidValue {
220                field: "event_type",
221                reason: "must match the memory event payload",
222            });
223        }
224        match &event.payload {
225            MemoryEvent::NodeUpserted { node } => self.insert_node(NodeRevision {
226                node: node.clone(),
227                recorded_at: event.metadata.recorded_at,
228                position: event.metadata.global_position,
229            })?,
230            MemoryEvent::FactRecorded { fact } => {
231                if fact.recorded_at != event.metadata.recorded_at
232                    || fact.agent_id != event.metadata.agent_id
233                    || fact.session_id != event.metadata.session_id
234                {
235                    return Err(MemoryError::InvalidValue {
236                        field: "fact.envelope",
237                        reason: "fact provenance must match its event envelope",
238                    });
239                }
240                self.insert_fact(fact.clone())?;
241            }
242            MemoryEvent::FactRetracted {
243                fact_id,
244                valid_until,
245                evidence,
246            } => self.apply_retraction(event, fact_id, *valid_until, evidence)?,
247        }
248        self.last_global_position = Some(event.metadata.global_position);
249        Ok(())
250    }
251}
252
253fn visible_revision(history: &NodeHistory, known_at: Timestamp) -> Option<&MemoryNode> {
254    core::iter::once(&history.first)
255        .chain(&history.later)
256        .filter(|revision| revision.recorded_at <= known_at)
257        .max_by_key(|revision| (revision.recorded_at, revision.position))
258        .map(|revision| &revision.node)
259}