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_meta, read_point_by_id,
6    read_stored_points, update_root, validate_config, validate_point, validate_root, PointChange,
7    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 mut states = BTreeMap::<PointId, (Option<Point>, Option<Point>)>::new();
133            let mut upsert_ids = BTreeSet::new();
134            for mutation in mutations {
135                match mutation {
136                    SnapshotMutation::Upsert { point } => {
137                        validate_point(&point, &config)?;
138                        let id = point.id.clone();
139                        if !upsert_ids.insert(id.clone()) {
140                            return Err(Error::Invalid(format!(
141                                "snapshot mutation batch contains duplicate upsert ID {}",
142                                id
143                            )));
144                        }
145                        if !states.contains_key(&id) {
146                            let old = read_point_by_id(&repo, previous_root, &id)?;
147                            states.insert(id.clone(), (old.clone(), old));
148                        }
149                        states.get_mut(&id).expect("point state exists").1 = Some(point);
150                    }
151                    SnapshotMutation::DeleteIds { ids } => {
152                        if ids.is_empty() {
153                            return Err(Error::Invalid(
154                                "snapshot delete_ids mutation must not be empty".into(),
155                            ));
156                        }
157                        for id in ids {
158                            if !states.contains_key(&id) {
159                                let old = read_point_by_id(&repo, previous_root, &id)?;
160                                states.insert(id.clone(), (old.clone(), old));
161                            }
162                            states.get_mut(&id).expect("point state exists").1 = None;
163                        }
164                    }
165                    SnapshotMutation::DeleteFilter { .. } => unreachable!(),
166                }
167            }
168            let mut final_point_count = meta.point_count();
169            let changes = states
170                .into_iter()
171                .filter_map(|(id, (old, new))| {
172                    if old == new {
173                        return None;
174                    }
175                    match (old.is_some(), new.is_some()) {
176                        (false, true) => final_point_count += 1,
177                        (true, false) => final_point_count -= 1,
178                        _ => {}
179                    }
180                    Some((id, PointChange { old, new }))
181                })
182                .collect();
183            let root = update_root(&repo, previous_root, &config, final_point_count, &changes)?;
184            return Ok(self.snapshot(root, None));
185        }
186        let stored_points = read_stored_points(&repo, previous_root)?;
187        let mut points = BTreeMap::new();
188        for (id, stored) in stored_points {
189            points.insert(id, stored.point);
190        }
191        let mut upsert_ids = BTreeSet::new();
192        let mut originals = BTreeMap::new();
193
194        for mutation in mutations {
195            match mutation {
196                SnapshotMutation::Upsert { point } => {
197                    validate_point(&point, &config)?;
198                    if !upsert_ids.insert(point.id.clone()) {
199                        return Err(Error::Invalid(format!(
200                            "snapshot mutation batch contains duplicate upsert ID {}",
201                            point.id
202                        )));
203                    }
204                    if !originals.contains_key(&point.id) {
205                        originals.insert(point.id.clone(), points.get(&point.id).cloned());
206                    }
207                    points.insert(point.id.clone(), point);
208                }
209                SnapshotMutation::DeleteIds { ids } => {
210                    if ids.is_empty() {
211                        return Err(Error::Invalid(
212                            "snapshot delete_ids mutation must not be empty".into(),
213                        ));
214                    }
215                    for id in ids {
216                        if !originals.contains_key(&id) {
217                            originals.insert(id.clone(), points.get(&id).cloned());
218                        }
219                        points.remove(&id);
220                    }
221                }
222                SnapshotMutation::DeleteFilter { filter } => {
223                    let ids = points
224                        .iter()
225                        .filter(|(id, point)| matches_filter(&filter, id, &point.payload))
226                        .map(|(id, _)| id.clone())
227                        .collect::<Vec<_>>();
228                    for id in ids {
229                        if !originals.contains_key(&id) {
230                            originals.insert(id.clone(), points.get(&id).cloned());
231                        }
232                        points.remove(&id);
233                    }
234                }
235            }
236        }
237
238        let changes = originals
239            .into_iter()
240            .filter_map(|(id, old)| {
241                let new = points.get(&id).cloned();
242                (old != new).then_some((id, PointChange { old, new }))
243            })
244            .collect();
245        let root = update_root(&repo, previous_root, &config, points.len(), &changes)?;
246        Ok(self.snapshot(root, Some(points)))
247    }
248
249    /// Opens an exact tree object ID without resolving refs or commits.
250    pub fn open_snapshot(&self, root: impl AsRef<str>) -> Result<Snapshot> {
251        let repo = self.repo()?;
252        let root = exact_root(&repo, root.as_ref())?;
253        read_meta(&repo, root)?;
254        Ok(self.snapshot(root, None))
255    }
256
257    /// Imports a materialized canonical tree into this engine's object database.
258    pub fn import_directory(&self, path: impl AsRef<Path>) -> Result<Snapshot> {
259        let path = path.as_ref();
260        if !path.is_dir() {
261            return Err(Error::Invalid(format!(
262                "materialized snapshot is not a directory: {}",
263                path.display()
264            )));
265        }
266        let repo = self.repo()?;
267        let root = import_directory(&repo, path)?;
268        read_meta(&repo, root)?;
269        Ok(self.snapshot(root, None))
270    }
271
272    /// Queries an exact root ID without first constructing a named collection.
273    pub fn query(&self, root: impl AsRef<str>, query: Query) -> Result<QueryResult> {
274        let repo = self.repo()?;
275        let root = exact_root(&repo, root.as_ref())?;
276        read_meta(&repo, root)?;
277        query_root_with_cache(&repo, root, query, None)
278    }
279
280    /// Retrieves records from an exact root ID.
281    pub fn get(&self, root: impl AsRef<str>, request: GetRequest) -> Result<GetResult> {
282        self.open_snapshot(root)?.get(request)
283    }
284
285    /// Counts records in an exact root ID.
286    pub fn count(
287        &self,
288        root: impl AsRef<str>,
289        filter: Option<crate::Filter>,
290    ) -> Result<CountResult> {
291        self.open_snapshot(root)?.count(filter)
292    }
293
294    /// Validates an exact root ID.
295    pub fn validate(&self, root: impl AsRef<str>, full: bool) -> Result<ValidationReport> {
296        self.open_snapshot(root)?.validate(full)
297    }
298
299    /// Builds a snapshot without a caller-provided Git repository and writes its
300    /// canonical files into a new materialized directory.
301    pub fn build_directory(
302        path: impl AsRef<Path>,
303        config: CollectionConfig,
304        points: Vec<Point>,
305    ) -> Result<Snapshot> {
306        let engine = Self::ephemeral()?;
307        let snapshot = engine.build(config, points)?;
308        snapshot.materialize(path.as_ref())?;
309        Snapshot::open_directory(path)
310    }
311
312    fn snapshot(&self, root: Oid, points: Option<BTreeMap<PointId, Point>>) -> Snapshot {
313        let cache = OnceLock::new();
314        if let Some(points) = points {
315            cache
316                .set(SearchView::new(points.into_values().collect()))
317                .expect("new snapshot point cache must be empty");
318        }
319        Snapshot {
320            object_database: self.object_database.clone(),
321            root,
322            points: Arc::new(cache),
323            temporary: self.temporary.clone(),
324        }
325    }
326}
327
328impl Snapshot {
329    /// Imports a materialized canonical tree into an isolated temporary object
330    /// database and opens the computed root. The source directory is not changed.
331    pub fn open_directory(path: impl AsRef<Path>) -> Result<Self> {
332        let engine = SnapshotEngine::ephemeral()?;
333        engine.import_directory(path)
334    }
335
336    fn repo(&self) -> Result<Repository> {
337        Ok(Repository::open(&self.object_database)?)
338    }
339
340    /// Returns the deterministic Git tree ID that identifies this snapshot.
341    pub fn root(&self) -> ObjectId {
342        self.root.into()
343    }
344
345    /// Returns configuration and point-count metadata for this root.
346    pub fn info(&self) -> Result<SnapshotInfo> {
347        let repo = self.repo()?;
348        let meta = read_meta(&repo, self.root)?;
349        let point_count = count_root(&repo, self.root, None)?.count;
350        Ok(SnapshotInfo {
351            root: self.root(),
352            point_count,
353            config: meta.config(),
354        })
355    }
356
357    /// Retrieves canonically ordered records from this immutable root.
358    pub fn get(&self, request: GetRequest) -> Result<GetResult> {
359        get_root(&self.repo()?, self.root, request)
360    }
361
362    /// Counts all points or those matching a filter at this immutable root.
363    pub fn count(&self, filter: Option<crate::Filter>) -> Result<CountResult> {
364        count_root(&self.repo()?, self.root, filter)
365    }
366
367    /// Executes an exact or deterministic approximate query at this root.
368    ///
369    /// The first exact query may populate a root-scoped immutable search view;
370    /// this cache does not alter persisted objects or the root ID.
371    pub fn query(&self, query: Query) -> Result<QueryResult> {
372        let repo = self.repo()?;
373        query_root_with_cache(&repo, self.root, query, Some(&self.points))
374    }
375
376    /// Applies mutations using this snapshot's object database without creating
377    /// collection history or refs.
378    pub fn apply(&self, mutations: Vec<SnapshotMutation>) -> Result<Snapshot> {
379        SnapshotEngine {
380            object_database: self.object_database.clone(),
381            temporary: self.temporary.clone(),
382        }
383        .apply(self.root.to_string(), mutations)
384    }
385
386    /// Validates this root without changing its object database.
387    ///
388    /// Full validation recomputes every approximate-index bucket.
389    pub fn validate(&self, full: bool) -> Result<ValidationReport> {
390        validate_root(&self.repo()?, self.root, full)
391    }
392
393    /// Writes this exact Git tree as ordinary files and directories.
394    ///
395    /// The target must not already exist. A sibling staging directory is renamed
396    /// into place only after every object has been read successfully.
397    pub fn materialize(&self, target: impl AsRef<Path>) -> Result<()> {
398        let target = target.as_ref();
399        if target.exists() {
400            return Err(Error::Invalid(format!(
401                "materialization target already exists: {}",
402                target.display()
403            )));
404        }
405        let parent = target
406            .parent()
407            .filter(|path| !path.as_os_str().is_empty())
408            .unwrap_or_else(|| Path::new("."));
409        fs::create_dir_all(parent)?;
410        let staging = tempfile::Builder::new()
411            .prefix(".git-vdb-snapshot-")
412            .tempdir_in(parent)?;
413        materialize_tree(&self.repo()?, self.root, staging.path())?;
414        fs::rename(staging.path(), target).map_err(|error| {
415            Error::Invalid(format!(
416                "cannot publish materialized snapshot {} as {}: {error}",
417                staging.path().display(),
418                target.display()
419            ))
420        })?;
421        Ok(())
422    }
423
424    pub(crate) fn oid(&self) -> Oid {
425        self.root
426    }
427}
428
429fn canonical_point_set(
430    points: Vec<Point>,
431    config: &CollectionConfig,
432) -> Result<BTreeMap<PointId, Point>> {
433    let mut canonical = BTreeMap::new();
434    for point in points {
435        validate_point(&point, config)?;
436        if canonical.insert(point.id.clone(), point).is_some() {
437            return Err(Error::Invalid(
438                "snapshot build contains a duplicate typed point ID".into(),
439            ));
440        }
441    }
442    Ok(canonical)
443}
444
445fn exact_root(repo: &Repository, root: &str) -> Result<Oid> {
446    let oid = Oid::from_str(root)
447        .map_err(|_| Error::Invalid(format!("invalid snapshot root object ID {root:?}")))?;
448    repo.find_tree(oid)
449        .map_err(|_| Error::Invalid(format!("snapshot root {root} is not a tree object")))?;
450    Ok(oid)
451}
452
453fn materialize_tree(repo: &Repository, tree_oid: Oid, path: &Path) -> Result<()> {
454    fs::create_dir_all(path)?;
455    let tree = repo.find_tree(tree_oid)?;
456    for entry in &tree {
457        let name = entry
458            .name()
459            .map_err(|_| Error::Corrupt("Git tree entry name is not UTF-8".into()))?;
460        validate_tree_name(name)?;
461        let destination = path.join(name);
462        match entry.kind() {
463            Some(ObjectType::Tree) => materialize_tree(repo, entry.id(), &destination)?,
464            Some(ObjectType::Blob) => {
465                fs::write(destination, repo.find_blob(entry.id())?.content())?
466            }
467            kind => {
468                return Err(Error::Corrupt(format!(
469                    "unsupported object kind {kind:?} in snapshot tree"
470                )));
471            }
472        }
473    }
474    Ok(())
475}
476
477fn validate_tree_name(name: &str) -> Result<()> {
478    if name.is_empty() || matches!(name, "." | "..") || name.contains(['/', '\\']) {
479        return Err(Error::Corrupt(format!(
480            "unsafe path name in snapshot tree: {name:?}"
481        )));
482    }
483    Ok(())
484}
485
486fn import_directory(repo: &Repository, path: &Path) -> Result<Oid> {
487    let mut entries = fs::read_dir(path)?
488        .map(|entry| entry.map(|entry| entry.path()))
489        .collect::<std::result::Result<Vec<PathBuf>, std::io::Error>>()?;
490    entries.sort_by(|left, right| left.file_name().cmp(&right.file_name()));
491
492    let mut tree = repo.treebuilder(None)?;
493    for path in entries {
494        let name = path
495            .file_name()
496            .and_then(|name| name.to_str())
497            .ok_or_else(|| Error::Invalid("snapshot path names must be valid UTF-8".into()))?;
498        let file_type = fs::symlink_metadata(&path)?.file_type();
499        if file_type.is_symlink() {
500            return Err(Error::Invalid(format!(
501                "materialized snapshots cannot contain symlinks: {}",
502                path.display()
503            )));
504        }
505        if file_type.is_dir() {
506            tree.insert(name, import_directory(repo, &path)?, TREE_MODE)?;
507        } else if file_type.is_file() {
508            tree.insert(name, repo.blob(&fs::read(&path)?)?, BLOB_MODE)?;
509        } else {
510            return Err(Error::Invalid(format!(
511                "unsupported materialized snapshot entry: {}",
512                path.display()
513            )));
514        }
515    }
516    Ok(tree.write()?)
517}