Skip to main content

git_vdb/
snapshot.rs

1//! Ref-free immutable snapshot construction, mutation, query, and validation.
2
3use crate::filter::matches_filter;
4use crate::root::{
5    build_root, count_root, get_root, query_root_with_cache, read_all_points, read_meta,
6    read_point_by_id, read_stored_points, update_root, validate_config, validate_point,
7    validate_root, PointChange, SearchView,
8};
9use crate::{
10    CollectionConfig, CountResult, Error, GetRequest, GetResult, ObjectId, Point, PointId, Query,
11    QueryResult, Result, SnapshotInfo, SnapshotMutation, ValidationReport,
12};
13use git2::{ObjectType, Oid, Repository};
14use std::collections::{BTreeMap, BTreeSet};
15use std::fmt;
16use std::fs;
17use std::path::{Path, PathBuf};
18use std::sync::{Arc, OnceLock};
19use tempfile::TempDir;
20
21const TREE_MODE: i32 = 0o040000;
22const BLOB_MODE: i32 = 0o100644;
23
24/// A ref-free engine for deterministic immutable collection roots.
25///
26/// The engine writes and reads Git objects but never creates a commit, updates a
27/// ref, consults repository history, or reads the clock. Callers are responsible
28/// for retaining returned roots, for example through an external content store
29/// or the named collection adapter.
30#[derive(Clone)]
31pub struct SnapshotEngine {
32    object_database: PathBuf,
33    temporary: Option<Arc<TempDir>>,
34}
35
36/// An immutable collection root opened independently of collection refs.
37#[derive(Clone)]
38pub struct Snapshot {
39    object_database: PathBuf,
40    root: Oid,
41    points: Arc<OnceLock<SearchView>>,
42    temporary: Option<Arc<TempDir>>,
43}
44
45impl fmt::Debug for SnapshotEngine {
46    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
47        formatter
48            .debug_struct("SnapshotEngine")
49            .field("object_database", &self.object_database)
50            .field("temporary", &self.temporary.is_some())
51            .finish()
52    }
53}
54
55impl fmt::Debug for Snapshot {
56    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
57        formatter
58            .debug_struct("Snapshot")
59            .field("root", &self.root)
60            .field("object_database", &self.object_database)
61            .field("temporary", &self.temporary.is_some())
62            .finish()
63    }
64}
65
66impl SnapshotEngine {
67    /// Opens an existing bare or non-bare Git object database without reading a
68    /// collection ref.
69    pub fn open(path: impl AsRef<Path>) -> Result<Self> {
70        let repository = Repository::open(path)?;
71        Ok(Self {
72            object_database: repository.path().to_path_buf(),
73            temporary: None,
74        })
75    }
76
77    /// Initializes a bare Git object database for immutable snapshots.
78    pub fn init(path: impl AsRef<Path>) -> Result<Self> {
79        let repository = Repository::init_bare(path)?;
80        Ok(Self {
81            object_database: repository.path().to_path_buf(),
82            temporary: None,
83        })
84    }
85
86    /// Creates an isolated temporary object database whose lifetime is retained
87    /// by the engine and snapshots returned from it.
88    pub fn ephemeral() -> Result<Self> {
89        let temporary = Arc::new(TempDir::new()?);
90        let repository = Repository::init_bare(temporary.path())?;
91        Ok(Self {
92            object_database: repository.path().to_path_buf(),
93            temporary: Some(temporary),
94        })
95    }
96
97    fn repo(&self) -> Result<Repository> {
98        Ok(Repository::open(&self.object_database)?)
99    }
100
101    /// Builds a canonical root from a complete point set.
102    pub fn build(&self, config: CollectionConfig, points: Vec<Point>) -> Result<Snapshot> {
103        validate_config(&config)?;
104        let points = canonical_point_set(points, &config)?;
105        let repo = self.repo()?;
106        let root = build_root(&repo, &config, &points)?;
107        Ok(self.snapshot(root, Some(points)))
108    }
109
110    /// Applies an ordered mutation batch to a root and returns the new root.
111    ///
112    /// Mutation order is significant. Duplicate upserts for the same typed ID
113    /// within one call are rejected. No ref or commit is created.
114    pub fn apply(
115        &self,
116        previous_root: impl AsRef<str>,
117        mutations: Vec<SnapshotMutation>,
118    ) -> Result<Snapshot> {
119        if mutations.is_empty() {
120            return Err(Error::Invalid(
121                "snapshot mutation batch must not be empty".into(),
122            ));
123        }
124        let repo = self.repo()?;
125        let previous_root = exact_root(&repo, previous_root.as_ref())?;
126        let meta = read_meta(&repo, previous_root)?;
127        let config = meta.config();
128        if mutations
129            .iter()
130            .all(|mutation| !matches!(mutation, SnapshotMutation::DeleteFilter { .. }))
131        {
132            let v2_points = (meta.format_version() == 2)
133                .then(|| read_all_points(&repo, previous_root))
134                .transpose()?;
135            let mut states = BTreeMap::<PointId, (Option<Point>, Option<Point>)>::new();
136            let mut upsert_ids = BTreeSet::new();
137            for mutation in mutations {
138                match mutation {
139                    SnapshotMutation::Upsert { point } => {
140                        validate_point(&point, &config)?;
141                        let id = point.id.clone();
142                        if !upsert_ids.insert(id.clone()) {
143                            return Err(Error::Invalid(format!(
144                                "snapshot mutation batch contains duplicate upsert ID {}",
145                                id
146                            )));
147                        }
148                        if !states.contains_key(&id) {
149                            let old = match &v2_points {
150                                Some(points) => points.get(&id).cloned(),
151                                None => read_point_by_id(&repo, previous_root, &id)?,
152                            };
153                            states.insert(id.clone(), (old.clone(), old));
154                        }
155                        states.get_mut(&id).expect("point state exists").1 = Some(point);
156                    }
157                    SnapshotMutation::DeleteIds { ids } => {
158                        if ids.is_empty() {
159                            return Err(Error::Invalid(
160                                "snapshot delete_ids mutation must not be empty".into(),
161                            ));
162                        }
163                        for id in ids {
164                            if !states.contains_key(&id) {
165                                let old = match &v2_points {
166                                    Some(points) => points.get(&id).cloned(),
167                                    None => read_point_by_id(&repo, previous_root, &id)?,
168                                };
169                                states.insert(id.clone(), (old.clone(), old));
170                            }
171                            states.get_mut(&id).expect("point state exists").1 = None;
172                        }
173                    }
174                    SnapshotMutation::DeleteFilter { .. } => unreachable!(),
175                }
176            }
177            let mut final_point_count = meta.point_count();
178            let changes = states
179                .into_iter()
180                .filter_map(|(id, (old, new))| {
181                    if old == new {
182                        return None;
183                    }
184                    match (old.is_some(), new.is_some()) {
185                        (false, true) => final_point_count += 1,
186                        (true, false) => final_point_count -= 1,
187                        _ => {}
188                    }
189                    Some((id, PointChange { old, new }))
190                })
191                .collect();
192            let root = update_root(&repo, previous_root, &config, final_point_count, &changes)?;
193            return Ok(self.snapshot(root, None));
194        }
195        let stored_points = read_stored_points(&repo, previous_root)?;
196        let mut points = BTreeMap::new();
197        for (id, stored) in stored_points {
198            points.insert(id, stored.point);
199        }
200        let mut upsert_ids = BTreeSet::new();
201        let mut originals = BTreeMap::new();
202
203        for mutation in mutations {
204            match mutation {
205                SnapshotMutation::Upsert { point } => {
206                    validate_point(&point, &config)?;
207                    if !upsert_ids.insert(point.id.clone()) {
208                        return Err(Error::Invalid(format!(
209                            "snapshot mutation batch contains duplicate upsert ID {}",
210                            point.id
211                        )));
212                    }
213                    if !originals.contains_key(&point.id) {
214                        originals.insert(point.id.clone(), points.get(&point.id).cloned());
215                    }
216                    points.insert(point.id.clone(), point);
217                }
218                SnapshotMutation::DeleteIds { ids } => {
219                    if ids.is_empty() {
220                        return Err(Error::Invalid(
221                            "snapshot delete_ids mutation must not be empty".into(),
222                        ));
223                    }
224                    for id in ids {
225                        if !originals.contains_key(&id) {
226                            originals.insert(id.clone(), points.get(&id).cloned());
227                        }
228                        points.remove(&id);
229                    }
230                }
231                SnapshotMutation::DeleteFilter { filter } => {
232                    let ids = points
233                        .iter()
234                        .filter(|(id, point)| matches_filter(&filter, id, &point.payload))
235                        .map(|(id, _)| id.clone())
236                        .collect::<Vec<_>>();
237                    for id in ids {
238                        if !originals.contains_key(&id) {
239                            originals.insert(id.clone(), points.get(&id).cloned());
240                        }
241                        points.remove(&id);
242                    }
243                }
244            }
245        }
246
247        let changes = originals
248            .into_iter()
249            .filter_map(|(id, old)| {
250                let new = points.get(&id).cloned();
251                (old != new).then_some((id, PointChange { old, new }))
252            })
253            .collect();
254        let root = update_root(&repo, previous_root, &config, points.len(), &changes)?;
255        Ok(self.snapshot(root, Some(points)))
256    }
257
258    /// Opens an exact tree object ID without resolving refs or commits.
259    pub fn open_snapshot(&self, root: impl AsRef<str>) -> Result<Snapshot> {
260        let repo = self.repo()?;
261        let root = exact_root(&repo, root.as_ref())?;
262        read_meta(&repo, root)?;
263        Ok(self.snapshot(root, None))
264    }
265
266    /// Imports a materialized canonical tree into this engine's object database.
267    pub fn import_directory(&self, path: impl AsRef<Path>) -> Result<Snapshot> {
268        let path = path.as_ref();
269        if !path.is_dir() {
270            return Err(Error::Invalid(format!(
271                "materialized snapshot is not a directory: {}",
272                path.display()
273            )));
274        }
275        let repo = self.repo()?;
276        let root = import_directory(&repo, path)?;
277        read_meta(&repo, root)?;
278        Ok(self.snapshot(root, None))
279    }
280
281    /// Queries an exact root ID without first constructing a named collection.
282    pub fn query(&self, root: impl AsRef<str>, query: Query) -> Result<QueryResult> {
283        let repo = self.repo()?;
284        let root = exact_root(&repo, root.as_ref())?;
285        read_meta(&repo, root)?;
286        query_root_with_cache(&repo, root, query, None)
287    }
288
289    /// Retrieves records from an exact root ID.
290    pub fn get(&self, root: impl AsRef<str>, request: GetRequest) -> Result<GetResult> {
291        self.open_snapshot(root)?.get(request)
292    }
293
294    /// Counts records in an exact root ID.
295    pub fn count(
296        &self,
297        root: impl AsRef<str>,
298        filter: Option<crate::Filter>,
299    ) -> Result<CountResult> {
300        self.open_snapshot(root)?.count(filter)
301    }
302
303    /// Validates an exact root ID.
304    pub fn validate(&self, root: impl AsRef<str>, full: bool) -> Result<ValidationReport> {
305        self.open_snapshot(root)?.validate(full)
306    }
307
308    /// Builds a snapshot without a caller-provided Git repository and writes its
309    /// canonical files into a new materialized directory.
310    pub fn build_directory(
311        path: impl AsRef<Path>,
312        config: CollectionConfig,
313        points: Vec<Point>,
314    ) -> Result<Snapshot> {
315        let engine = Self::ephemeral()?;
316        let snapshot = engine.build(config, points)?;
317        snapshot.materialize(path.as_ref())?;
318        Snapshot::open_directory(path)
319    }
320
321    fn snapshot(&self, root: Oid, points: Option<BTreeMap<PointId, Point>>) -> Snapshot {
322        let cache = OnceLock::new();
323        if let Some(points) = points {
324            cache
325                .set(SearchView::new(points.into_values().collect()))
326                .expect("new snapshot point cache must be empty");
327        }
328        Snapshot {
329            object_database: self.object_database.clone(),
330            root,
331            points: Arc::new(cache),
332            temporary: self.temporary.clone(),
333        }
334    }
335}
336
337impl Snapshot {
338    /// Imports a materialized canonical tree into an isolated temporary object
339    /// database and opens the computed root. The source directory is not changed.
340    pub fn open_directory(path: impl AsRef<Path>) -> Result<Self> {
341        let engine = SnapshotEngine::ephemeral()?;
342        engine.import_directory(path)
343    }
344
345    fn repo(&self) -> Result<Repository> {
346        Ok(Repository::open(&self.object_database)?)
347    }
348
349    /// Returns the deterministic Git tree ID that identifies this snapshot.
350    pub fn root(&self) -> ObjectId {
351        self.root.into()
352    }
353
354    /// Returns configuration and point-count metadata for this root.
355    pub fn info(&self) -> Result<SnapshotInfo> {
356        let repo = self.repo()?;
357        let meta = read_meta(&repo, self.root)?;
358        let point_count = count_root(&repo, self.root, None)?.count;
359        Ok(SnapshotInfo {
360            root: self.root(),
361            format_version: meta.format_version(),
362            point_count,
363            config: meta.config(),
364        })
365    }
366
367    /// Retrieves canonically ordered records from this immutable root.
368    pub fn get(&self, request: GetRequest) -> Result<GetResult> {
369        get_root(&self.repo()?, self.root, request)
370    }
371
372    /// Counts all points or those matching a filter at this immutable root.
373    pub fn count(&self, filter: Option<crate::Filter>) -> Result<CountResult> {
374        count_root(&self.repo()?, self.root, filter)
375    }
376
377    /// Executes an exact or deterministic approximate query at this root.
378    ///
379    /// The first exact query may populate a root-scoped immutable search view;
380    /// this cache does not alter persisted objects or the root ID.
381    pub fn query(&self, query: Query) -> Result<QueryResult> {
382        let repo = self.repo()?;
383        query_root_with_cache(&repo, self.root, query, Some(&self.points))
384    }
385
386    /// Applies mutations using this snapshot's object database without creating
387    /// collection history or refs.
388    pub fn apply(&self, mutations: Vec<SnapshotMutation>) -> Result<Snapshot> {
389        SnapshotEngine {
390            object_database: self.object_database.clone(),
391            temporary: self.temporary.clone(),
392        }
393        .apply(self.root.to_string(), mutations)
394    }
395
396    /// Validates this root without changing its object database.
397    ///
398    /// Full validation recomputes every approximate-index bucket.
399    pub fn validate(&self, full: bool) -> Result<ValidationReport> {
400        validate_root(&self.repo()?, self.root, full)
401    }
402
403    /// Writes this exact Git tree as ordinary files and directories.
404    ///
405    /// The target must not already exist. A sibling staging directory is renamed
406    /// into place only after every object has been read successfully.
407    pub fn materialize(&self, target: impl AsRef<Path>) -> Result<()> {
408        let target = target.as_ref();
409        if target.exists() {
410            return Err(Error::Invalid(format!(
411                "materialization target already exists: {}",
412                target.display()
413            )));
414        }
415        let parent = target
416            .parent()
417            .filter(|path| !path.as_os_str().is_empty())
418            .unwrap_or_else(|| Path::new("."));
419        fs::create_dir_all(parent)?;
420        let staging = tempfile::Builder::new()
421            .prefix(".git-vdb-snapshot-")
422            .tempdir_in(parent)?;
423        materialize_tree(&self.repo()?, self.root, staging.path())?;
424        fs::rename(staging.path(), target).map_err(|error| {
425            Error::Invalid(format!(
426                "cannot publish materialized snapshot {} as {}: {error}",
427                staging.path().display(),
428                target.display()
429            ))
430        })?;
431        Ok(())
432    }
433
434    pub(crate) fn oid(&self) -> Oid {
435        self.root
436    }
437}
438
439fn canonical_point_set(
440    points: Vec<Point>,
441    config: &CollectionConfig,
442) -> Result<BTreeMap<PointId, Point>> {
443    let mut canonical = BTreeMap::new();
444    for point in points {
445        validate_point(&point, config)?;
446        if canonical.insert(point.id.clone(), point).is_some() {
447            return Err(Error::Invalid(
448                "snapshot build contains a duplicate typed point ID".into(),
449            ));
450        }
451    }
452    Ok(canonical)
453}
454
455fn exact_root(repo: &Repository, root: &str) -> Result<Oid> {
456    let oid = Oid::from_str(root)
457        .map_err(|_| Error::Invalid(format!("invalid snapshot root object ID {root:?}")))?;
458    repo.find_tree(oid)
459        .map_err(|_| Error::Invalid(format!("snapshot root {root} is not a tree object")))?;
460    Ok(oid)
461}
462
463fn materialize_tree(repo: &Repository, tree_oid: Oid, path: &Path) -> Result<()> {
464    fs::create_dir_all(path)?;
465    let tree = repo.find_tree(tree_oid)?;
466    for entry in &tree {
467        let name = entry
468            .name()
469            .map_err(|_| Error::Corrupt("Git tree entry name is not UTF-8".into()))?;
470        validate_tree_name(name)?;
471        let destination = path.join(name);
472        match entry.kind() {
473            Some(ObjectType::Tree) => materialize_tree(repo, entry.id(), &destination)?,
474            Some(ObjectType::Blob) => {
475                fs::write(destination, repo.find_blob(entry.id())?.content())?
476            }
477            kind => {
478                return Err(Error::Corrupt(format!(
479                    "unsupported object kind {kind:?} in snapshot tree"
480                )));
481            }
482        }
483    }
484    Ok(())
485}
486
487fn validate_tree_name(name: &str) -> Result<()> {
488    if name.is_empty() || matches!(name, "." | "..") || name.contains(['/', '\\']) {
489        return Err(Error::Corrupt(format!(
490            "unsafe path name in snapshot tree: {name:?}"
491        )));
492    }
493    Ok(())
494}
495
496fn import_directory(repo: &Repository, path: &Path) -> Result<Oid> {
497    let mut entries = fs::read_dir(path)?
498        .map(|entry| entry.map(|entry| entry.path()))
499        .collect::<std::result::Result<Vec<PathBuf>, std::io::Error>>()?;
500    entries.sort_by(|left, right| left.file_name().cmp(&right.file_name()));
501
502    let mut tree = repo.treebuilder(None)?;
503    for path in entries {
504        let name = path
505            .file_name()
506            .and_then(|name| name.to_str())
507            .ok_or_else(|| Error::Invalid("snapshot path names must be valid UTF-8".into()))?;
508        let file_type = fs::symlink_metadata(&path)?.file_type();
509        if file_type.is_symlink() {
510            return Err(Error::Invalid(format!(
511                "materialized snapshots cannot contain symlinks: {}",
512                path.display()
513            )));
514        }
515        if file_type.is_dir() {
516            tree.insert(name, import_directory(repo, &path)?, TREE_MODE)?;
517        } else if file_type.is_file() {
518            tree.insert(name, repo.blob(&fs::read(&path)?)?, BLOB_MODE)?;
519        } else {
520            return Err(Error::Invalid(format!(
521                "unsupported materialized snapshot entry: {}",
522                path.display()
523            )));
524        }
525    }
526    Ok(tree.write()?)
527}