Skip to main content

lora_store/memory/
snapshot.rs

1//! Snapshot payload helpers for the in-memory graph.
2//!
3//! `lora-store` no longer ships its own on-disk codec. The byte-level
4//! columnar format lives in `lora-snapshot`; this module just bridges
5//! between [`InMemoryGraph`] and the portable [`SnapshotPayload`]
6//! vocabulary.
7
8use std::collections::BTreeSet;
9
10use crate::{SnapshotError, SnapshotMeta, SnapshotPayload};
11
12use super::distinct_stats::DistinctBuilder;
13use super::index_catalog::StoredIndexEntity;
14use super::InMemoryGraph;
15
16/// Format-version stamp surfaced through [`SnapshotMeta::format_version`]
17/// for payloads produced via the inherent helpers below. Kept stable
18/// across `lora-snapshot` codec versions because the payload shape
19/// itself has not changed; only the on-disk encoding has.
20pub(super) const PAYLOAD_FORMAT_VERSION: u32 = 1;
21
22impl InMemoryGraph {
23    /// Return the portable graph-state payload. Callers downstream of
24    /// `lora-store` (typically `lora-database`) feed this into
25    /// `lora-snapshot` for byte-level encoding.
26    pub fn snapshot_payload(&self) -> SnapshotPayload {
27        let mut vector_indexes = self
28            .vector_indexes_read(StoredIndexEntity::Node)
29            .to_snapshots(StoredIndexEntity::Node);
30        vector_indexes.extend(
31            self.vector_indexes_read(StoredIndexEntity::Relationship)
32                .to_snapshots(StoredIndexEntity::Relationship),
33        );
34        SnapshotPayload {
35            next_node_id: self.next_node_id,
36            next_rel_id: self.next_rel_id,
37            nodes: self.iter_node_refs().map(|r| r.to_record()).collect(),
38            relationships: self.iter_rel_refs().map(|r| r.to_record()).collect(),
39            indexes: self.index_catalog_read().list(),
40            constraints: self.constraint_catalog_read().list(),
41            vector_indexes,
42        }
43    }
44
45    /// Replace the graph from a portable graph-state payload, preserving the
46    /// currently installed mutation recorder across the swap.
47    pub fn load_snapshot_payload(
48        &mut self,
49        payload: SnapshotPayload,
50    ) -> Result<SnapshotMeta, SnapshotError> {
51        let meta = SnapshotMeta {
52            format_version: PAYLOAD_FORMAT_VERSION,
53            node_count: payload.nodes.len(),
54            relationship_count: payload.relationships.len(),
55            wal_lsn: None,
56        };
57
58        validate_payload_ids(&payload)?;
59
60        // Build the restored graph in a fresh local instance and only
61        // commit it into `self` at the very end. Capacity is based on live
62        // entity count, not `next_*_id`: snapshots may contain tombstone gaps,
63        // and hostile next-id values must not force huge allocations before
64        // the checked slab-growth path validates each concrete record id.
65        let node_capacity = payload.nodes.len();
66        let relationship_capacity = payload.relationships.len();
67        let mut rebuilt = Self::with_capacity_hint(node_capacity, relationship_capacity);
68        rebuilt.next_node_id = payload.next_node_id;
69        rebuilt.next_rel_id = payload.next_rel_id;
70
71        let mut node_distinct = DistinctBuilder::default();
72        let mut rel_distinct = DistinctBuilder::default();
73        for node in payload.nodes {
74            let id = node.id;
75            let labels = node.labels.clone();
76            // Distinct-value sketches are not persisted: count the record
77            // as a live create would, through the bulk builder.
78            for label in &labels {
79                node_distinct.add_all(label, &node.properties);
80            }
81            if rebuilt.node_at(id).is_some() {
82                return Err(SnapshotError::Decode(format!(
83                    "duplicate node id {id} in snapshot payload"
84                )));
85            }
86            rebuilt
87                .put_node_checked(id, &node)
88                .map_err(SnapshotError::Decode)?;
89            for label in &labels {
90                rebuilt.insert_node_label_index(id, label);
91            }
92        }
93
94        for rel in payload.relationships {
95            if rebuilt.rel_at(rel.id).is_some() {
96                return Err(SnapshotError::Decode(format!(
97                    "duplicate relationship id {} in snapshot payload",
98                    rel.id
99                )));
100            }
101            if rebuilt.node_at(rel.src).is_none() {
102                return Err(SnapshotError::Decode(format!(
103                    "relationship {} references missing source node {}",
104                    rel.id, rel.src
105                )));
106            }
107            if rebuilt.node_at(rel.dst).is_none() {
108                return Err(SnapshotError::Decode(format!(
109                    "relationship {} references missing target node {}",
110                    rel.id, rel.dst
111                )));
112            }
113            let id = rel.id;
114            rebuilt
115                .put_rel_checked(id, &rel)
116                .map_err(SnapshotError::Decode)?;
117            rebuilt.attach_relationship(&rel);
118            rel_distinct.add_all(&rel.rel_type, &rel.properties);
119        }
120        rebuilt.distinct_stats.nodes = node_distinct.finish();
121        rebuilt.distinct_stats.relationships = rel_distinct.finish();
122        // No hash property index is built here. Restart reproduces the
123        // state of a fresh process: the declared RANGE indexes and the
124        // uniqueness/key constraints below activate (and backfill) the
125        // hash indexes they need through `register_index` /
126        // `register_constraint`. Implicit, lookup-activated indexes are
127        // rebuilt lazily by the first lookup that needs them.
128
129        let constraint_owned_indexes: BTreeSet<String> = payload
130            .constraints
131            .iter()
132            .filter_map(|def| {
133                def.owned_index
134                    .clone()
135                    .or_else(|| def.kind.requires_backing_index().then(|| def.name.clone()))
136            })
137            .collect();
138
139        // Re-register every user-visible index in the catalog. Going through
140        // `register_index` re-populates RANGE buckets and keeps the
141        // `populate_index_data` invariant aligned with the catalog —
142        // skipping it would leave RANGE indexes registered but never populated.
143        // Constraint-owned backing indexes are restored by re-registering the
144        // owning constraint below, which keeps catalog ownership explicit.
145        for def in payload.indexes {
146            if constraint_owned_indexes.contains(&def.name) {
147                continue;
148            }
149            // Errors here would mean the snapshot itself is corrupt or
150            // ambiguous; map them into Decode rather than panicking.
151            rebuilt
152                .register_index(
153                    crate::memory::IndexRequest {
154                        explicit_name: Some(def.name.clone()),
155                        kind: def.kind,
156                        entity: def.entity,
157                        label: def.label.clone(),
158                        additional_labels: def.additional_labels.clone(),
159                        properties: def.properties.clone(),
160                        options: def.options.clone(),
161                    },
162                    /*if_not_exists*/ true,
163                )
164                .map_err(|e| SnapshotError::Decode(format!("index `{}`: {e}", def.name)))?;
165        }
166
167        // Overlay persisted HNSW snapshots over the freshly-registered
168        // (and freshly-backfilled) vector indexes. This is the
169        // post-step that gives Phase 5 its raison d'être: instead of
170        // paying O(n log n) to re-insert every vector through the
171        // HNSW algorithm, we install the graph topology byte-for-byte.
172        // Snapshots from versions before this trailer round-trip with
173        // `vector_indexes = []` so the fallback path (the backfill
174        // that already ran inside `register_index`) handles them
175        // correctly.
176        for snap in payload.vector_indexes {
177            let entity = snap.entity;
178            let mut registry = rebuilt.vector_indexes_write(entity);
179            if !registry.restore_snapshot(snap) {
180                // Catalog/snapshot mismatch — registry already
181                // contains the populate-built backend, which is the
182                // safe fallback. No further action.
183            }
184        }
185
186        // Re-register constraints. Uniqueness / key constraints recreate
187        // their own backing indexes as part of registration.
188        for def in payload.constraints {
189            rebuilt
190                .register_constraint(
191                    crate::memory::ConstraintRequest {
192                        name: def.name.clone(),
193                        kind: def.kind.clone(),
194                        entity: def.entity,
195                        label: def.label.clone(),
196                        properties: def.properties.clone(),
197                    },
198                    /*if_not_exists*/ true,
199                )
200                .map_err(|e| SnapshotError::Decode(format!("constraint `{}`: {e}", def.name)))?;
201        }
202
203        // Preserve the existing recorder across the swap — observers of the
204        // store's identity should not be silently detached by a restore,
205        // same policy as `clear()`.
206        rebuilt.recorder = self.recorder.take();
207        *self = rebuilt;
208
209        Ok(meta)
210    }
211}
212
213fn validate_payload_ids(payload: &SnapshotPayload) -> Result<(), SnapshotError> {
214    validate_next_id("node", payload.next_node_id)?;
215    validate_next_id("relationship", payload.next_rel_id)?;
216
217    for node in &payload.nodes {
218        validate_entity_id("node", node.id, payload.next_node_id)?;
219    }
220    for rel in &payload.relationships {
221        validate_entity_id("relationship", rel.id, payload.next_rel_id)?;
222        validate_slot_id("relationship source node", rel.src)?;
223        validate_slot_id("relationship target node", rel.dst)?;
224    }
225
226    Ok(())
227}
228
229fn validate_next_id(kind: &str, next_id: u64) -> Result<(), SnapshotError> {
230    validate_slot_id(&format!("next {kind} id"), next_id)?;
231    if next_id == u64::MAX {
232        return Err(SnapshotError::Decode(format!(
233            "next {kind} id {next_id} leaves no allocatable id"
234        )));
235    }
236    Ok(())
237}
238
239fn validate_entity_id(kind: &str, id: u64, next_id: u64) -> Result<(), SnapshotError> {
240    validate_slot_id(kind, id)?;
241    if id >= next_id {
242        return Err(SnapshotError::Decode(format!(
243            "{kind} id {id} is not below next {kind} id {next_id}"
244        )));
245    }
246    Ok(())
247}
248
249fn validate_slot_id(label: &str, id: u64) -> Result<(), SnapshotError> {
250    let idx = usize::try_from(id).map_err(|_| {
251        SnapshotError::Decode(format!(
252            "{label} {id} does not fit in usize on this platform"
253        ))
254    })?;
255    idx.checked_add(1)
256        .ok_or_else(|| SnapshotError::Decode(format!("{label} {id} leaves no valid slab slot")))?;
257    Ok(())
258}