1use 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#[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 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 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 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 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 pub fn get(&self, root: impl AsRef<str>, request: GetRequest) -> Result<GetResult> {
291 self.open_snapshot(root)?.get(request)
292 }
293
294 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 pub fn validate(&self, root: impl AsRef<str>, full: bool) -> Result<ValidationReport> {
305 self.open_snapshot(root)?.validate(full)
306 }
307
308 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 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 pub fn root(&self) -> ObjectId {
351 self.root.into()
352 }
353
354 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 pub fn get(&self, request: GetRequest) -> Result<GetResult> {
369 get_root(&self.repo()?, self.root, request)
370 }
371
372 pub fn count(&self, filter: Option<crate::Filter>) -> Result<CountResult> {
374 count_root(&self.repo()?, self.root, filter)
375 }
376
377 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 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 pub fn validate(&self, full: bool) -> Result<ValidationReport> {
400 validate_root(&self.repo()?, self.root, full)
401 }
402
403 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}