1use 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#[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 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 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 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 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 pub fn get(&self, root: impl AsRef<str>, request: GetRequest) -> Result<GetResult> {
296 self.open_snapshot(root)?.get(request)
297 }
298
299 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 pub fn validate(&self, root: impl AsRef<str>, full: bool) -> Result<ValidationReport> {
310 self.open_snapshot(root)?.validate(full)
311 }
312
313 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 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 pub fn root(&self) -> ObjectId {
356 self.root.into()
357 }
358
359 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 pub fn get(&self, request: GetRequest) -> Result<GetResult> {
374 get_root(&self.repo()?, self.root, request)
375 }
376
377 pub fn count(&self, filter: Option<crate::Filter>) -> Result<CountResult> {
379 count_root(&self.repo()?, self.root, filter)
380 }
381
382 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 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 pub fn validate(&self, full: bool) -> Result<ValidationReport> {
405 validate_root(&self.repo()?, self.root, full)
406 }
407
408 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}