1use 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#[derive(Clone)]
31pub struct SnapshotEngine {
32 object_database: PathBuf,
33 temporary: Option<Arc<TempDir>>,
34}
35
36#[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 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 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 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 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 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 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 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 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 pub fn get(&self, root: impl AsRef<str>, request: GetRequest) -> Result<GetResult> {
282 self.open_snapshot(root)?.get(request)
283 }
284
285 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 pub fn validate(&self, root: impl AsRef<str>, full: bool) -> Result<ValidationReport> {
296 self.open_snapshot(root)?.validate(full)
297 }
298
299 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 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 pub fn root(&self) -> ObjectId {
342 self.root.into()
343 }
344
345 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 pub fn get(&self, request: GetRequest) -> Result<GetResult> {
359 get_root(&self.repo()?, self.root, request)
360 }
361
362 pub fn count(&self, filter: Option<crate::Filter>) -> Result<CountResult> {
364 count_root(&self.repo()?, self.root, filter)
365 }
366
367 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 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 pub fn validate(&self, full: bool) -> Result<ValidationReport> {
390 validate_root(&self.repo()?, self.root, full)
391 }
392
393 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}