Skip to main content

lora_store/
mutation.rs

1//! Mutation events and the optional recorder hook.
2//!
3//! [`MutationEvent`] is the vocabulary a write-ahead log (or any observer —
4//! replication, audit, change-data-capture) appends to a durable stream.
5//! The enum covers every method on [`GraphStorageMut`]: each event carries
6//! exactly the information needed to deterministically re-apply the mutation
7//! against an empty store (or a snapshot) and recover the same state.
8//!
9//! [`MutationRecorder`] is the observer trait. Backends that want to emit
10//! events install a recorder via [`InMemoryGraph::set_mutation_recorder`].
11//! The default is `None` so zero-WAL workloads pay only a null-pointer check
12//! per mutation — no allocation, no clone.
13//!
14//! The persistent WAL implementation lives in the `lora-wal` crate, which
15//! supplies a `WalRecorder` that implements `MutationRecorder` by
16//! appending each event to an on-disk log. The snapshot header's
17//! `wal_lsn` field is what makes the checkpoint hybrid expressible
18//! across crate boundaries without `lora-store` learning about the WAL.
19
20use serde::{Deserialize, Serialize};
21
22use crate::memory::{ConstraintRequest, IndexRequest};
23use crate::{NodeId, NodeRecord, Properties, PropertyValue, RelationshipId, RelationshipRecord};
24
25/// A durable, replayable mutation against a graph store.
26///
27/// Each variant mirrors a method on `GraphStorageMut`. Applying every event
28/// in order against a store initialised from the snapshot whose `wal_lsn`
29/// immediately precedes the first event reproduces the committed state.
30///
31/// The enum derives `Serialize`/`Deserialize` for non-WAL observers and
32/// tooling; the production WAL uses its own compact tagged codec.
33#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
34pub enum MutationEvent {
35    CreateNode {
36        /// Id the backend allocated for the new node. Captured so replay
37        /// against a clean store produces the same id assignment as the
38        /// original (`next_node_id` advances deterministically).
39        id: NodeId,
40        labels: Vec<String>,
41        properties: Properties,
42    },
43    CreateRelationship {
44        id: RelationshipId,
45        src: NodeId,
46        dst: NodeId,
47        rel_type: String,
48        properties: Properties,
49    },
50    SetNodeProperty {
51        node_id: NodeId,
52        key: String,
53        value: PropertyValue,
54    },
55    RemoveNodeProperty {
56        node_id: NodeId,
57        key: String,
58    },
59    AddNodeLabel {
60        node_id: NodeId,
61        label: String,
62    },
63    RemoveNodeLabel {
64        node_id: NodeId,
65        label: String,
66    },
67    SetRelationshipProperty {
68        rel_id: RelationshipId,
69        key: String,
70        value: PropertyValue,
71    },
72    RemoveRelationshipProperty {
73        rel_id: RelationshipId,
74        key: String,
75    },
76    DeleteRelationship {
77        rel_id: RelationshipId,
78    },
79    DeleteNode {
80        node_id: NodeId,
81    },
82    DetachDeleteNode {
83        node_id: NodeId,
84    },
85    Clear,
86    /// Catalog-level mutation: register an index in the catalog. Replay
87    /// re-applies via the same `register_index` path with the recorder
88    /// detached, so events do not duplicate themselves.
89    CreateIndex {
90        request: IndexRequest,
91        if_not_exists: bool,
92    },
93    /// Catalog-level mutation: drop an index by name.
94    DropIndex {
95        name: String,
96        if_exists: bool,
97    },
98    /// Catalog-level mutation: register a constraint in the catalog.
99    CreateConstraint {
100        request: ConstraintRequest,
101        if_not_exists: bool,
102    },
103    /// Catalog-level mutation: drop a constraint by name. The store
104    /// cascades to the backing index when the constraint owned one.
105    DropConstraint {
106        name: String,
107        if_exists: bool,
108    },
109}
110
111/// Observer that receives every successful mutation in the order the store
112/// applied it.
113///
114/// The recorder sees events *after* the mutation has been applied to the
115/// in-memory state, so it never observes a mutation that the store
116/// rejected (invalid id, empty relationship type, …). This matches the
117/// classic write-ahead-log convention of logging committed changes only.
118///
119/// Implementations must be `Send + Sync` so a shared recorder can be driven
120/// from any thread holding the store's write lock.
121pub trait MutationRecorder: Send + Sync + 'static {
122    fn record(&self, event: MutationEvent);
123
124    /// Sticky failure flag for durability-shaped recorders.
125    ///
126    /// `record` itself is infallible — non-WAL observers (audit taps,
127    /// replication shadows, CDC sinks) should not abort a write because
128    /// their downstream queue is full. Recorders that *do* care about
129    /// durability — most importantly the WAL adapter — flip a flag when
130    /// an append fails and surface it here. The host (typically
131    /// `Database::execute_with_params`) polls this once per critical
132    /// section while still holding the store write lock; if poisoned, the
133    /// query fails loudly and the caller observes the durability error
134    /// rather than a silently-lost write.
135    ///
136    /// The default returns `None`, so existing recorders compile
137    /// unchanged.
138    fn poisoned(&self) -> Option<String> {
139        None
140    }
141}
142
143/// Observer for records a delete is about to drop.
144///
145/// Installed with [`InMemoryGraph::set_deleted_record_sink`]. Change feeds
146/// use it to describe deleted entities (labels, properties, endpoints)
147/// when the write mutates the live graph in place and no pre-write copy
148/// exists.
149///
150/// [`InMemoryGraph::set_deleted_record_sink`]: crate::InMemoryGraph::set_deleted_record_sink
151pub trait DeletedRecordSink: Send + Sync + 'static {
152    fn node_deleted(&self, record: &NodeRecord);
153    fn relationship_deleted(&self, record: &RelationshipRecord);
154}
155
156/// Convenience adapter that turns any `Fn(MutationEvent) + Send + Sync`
157/// into a `MutationRecorder` — useful in tests and for quick wiring.
158pub struct ClosureRecorder<F>(pub F)
159where
160    F: Fn(MutationEvent) + Send + Sync + 'static;
161
162impl<F> MutationRecorder for ClosureRecorder<F>
163where
164    F: Fn(MutationEvent) + Send + Sync + 'static,
165{
166    fn record(&self, event: MutationEvent) {
167        (self.0)(event)
168    }
169}
170
171/// Set of record ids touched by a buffered [`MutationEvent`] stream.
172///
173/// Built incrementally as events buffer (or in one pass at commit
174/// time) by [`MutationWriteSet::extend_from_events`]. Used by the OCC
175/// auto-commit path to (a) sort lock-acquire on commit, (b) validate
176/// per-record Arc identity against the snapshot.
177#[derive(Debug, Default, Clone)]
178pub struct MutationWriteSet {
179    /// Nodes whose record was created, modified, or deleted.
180    pub nodes: std::collections::BTreeSet<NodeId>,
181    /// Relationships whose record was created, modified, or deleted.
182    pub rels: std::collections::BTreeSet<RelationshipId>,
183    /// `true` if the stream contained a `MutationEvent::Clear`. A
184    /// clear invalidates any per-record check — the writer must
185    /// fall back to a full-graph commit (or fail under OCC).
186    pub cleared: bool,
187}
188
189impl MutationWriteSet {
190    pub fn new() -> Self {
191        Self::default()
192    }
193
194    /// Walk a `MutationEvent` stream and accumulate every touched
195    /// record id. Variants that touch two records (e.g.
196    /// `CreateRelationship` mentions both `src` and `dst` plus the
197    /// new relationship) record both nodes — the writer's view of
198    /// those nodes' adjacency changed too.
199    pub fn extend_from_events<'a>(&mut self, events: impl IntoIterator<Item = &'a MutationEvent>) {
200        for event in events {
201            match event {
202                MutationEvent::CreateNode { id, .. } => {
203                    self.nodes.insert(*id);
204                }
205                MutationEvent::CreateRelationship { id, src, dst, .. } => {
206                    self.rels.insert(*id);
207                    self.nodes.insert(*src);
208                    self.nodes.insert(*dst);
209                }
210                MutationEvent::SetNodeProperty { node_id, .. }
211                | MutationEvent::RemoveNodeProperty { node_id, .. }
212                | MutationEvent::AddNodeLabel { node_id, .. }
213                | MutationEvent::RemoveNodeLabel { node_id, .. } => {
214                    self.nodes.insert(*node_id);
215                }
216                MutationEvent::SetRelationshipProperty { rel_id, .. }
217                | MutationEvent::RemoveRelationshipProperty { rel_id, .. } => {
218                    self.rels.insert(*rel_id);
219                }
220                MutationEvent::DeleteRelationship { rel_id } => {
221                    self.rels.insert(*rel_id);
222                }
223                MutationEvent::DeleteNode { node_id } => {
224                    self.nodes.insert(*node_id);
225                }
226                MutationEvent::DetachDeleteNode { node_id } => {
227                    // Detach-delete also touches every incident
228                    // relationship, but those fire as
229                    // `DeleteRelationship` events of their own and
230                    // the surrounding loop will pick them up.
231                    self.nodes.insert(*node_id);
232                }
233                MutationEvent::Clear => {
234                    self.cleared = true;
235                }
236                MutationEvent::CreateIndex { .. }
237                | MutationEvent::DropIndex { .. }
238                | MutationEvent::CreateConstraint { .. }
239                | MutationEvent::DropConstraint { .. } => {
240                    // Catalog mutations don't touch node/rel write
241                    // sets — they live next to the graph slabs but
242                    // don't share record locks. The single-writer
243                    // `writer` mutex on the database serialises them
244                    // against everything else.
245                }
246            }
247        }
248    }
249
250    pub fn is_empty(&self) -> bool {
251        !self.cleared && self.nodes.is_empty() && self.rels.is_empty()
252    }
253}
254
255#[cfg(test)]
256mod tests {
257    use super::*;
258
259    #[test]
260    fn write_set_extracts_ids_from_events() {
261        let events = [
262            MutationEvent::CreateNode {
263                id: 1,
264                labels: vec!["A".into()],
265                properties: Default::default(),
266            },
267            MutationEvent::CreateRelationship {
268                id: 10,
269                src: 1,
270                dst: 2,
271                rel_type: "R".into(),
272                properties: Default::default(),
273            },
274            MutationEvent::SetNodeProperty {
275                node_id: 3,
276                key: "x".into(),
277                value: PropertyValue::Int(5),
278            },
279            MutationEvent::DeleteRelationship { rel_id: 11 },
280        ];
281
282        let mut ws = MutationWriteSet::new();
283        ws.extend_from_events(events.iter());
284
285        // CreateRelationship pulls in src=1, dst=2 alongside its own rel id.
286        assert_eq!(ws.nodes.iter().copied().collect::<Vec<_>>(), vec![1, 2, 3]);
287        assert_eq!(ws.rels.iter().copied().collect::<Vec<_>>(), vec![10, 11]);
288        assert!(!ws.cleared);
289    }
290
291    #[test]
292    fn write_set_clear_event_is_sticky() {
293        let mut ws = MutationWriteSet::new();
294        ws.extend_from_events([&MutationEvent::Clear]);
295        assert!(ws.cleared);
296        assert!(!ws.is_empty()); // cleared counts as non-empty
297    }
298}