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