weavatrix_memory/projection/
memory.rs1mod 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}