Skip to main content

git_vdb/
snapshot.rs

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