1use crate::root::{
4 count_root, diff_roots, get_root, query_root_with_cache, read_meta, validate_config,
5 validate_root, SearchView,
6};
7use crate::*;
8use git2::{Commit, Oid, Repository, Signature};
9use std::path::{Path, PathBuf};
10use std::sync::{Arc, Mutex, OnceLock};
11
12#[derive(Clone, Copy, Debug)]
13struct ResolvedSnapshot {
14 root: Oid,
15 commit: Option<Oid>,
16}
17
18#[derive(Clone, Debug)]
20pub struct Database {
21 pub(crate) path: PathBuf,
22}
23
24#[derive(Clone, Debug)]
26pub struct Collection {
27 db: Database,
28 name: String,
29 historical: Option<ResolvedSnapshot>,
30 query_cache: Arc<Mutex<CollectionQueryCache>>,
31}
32
33#[derive(Debug)]
34struct CollectionQueryCache {
35 root: Option<Oid>,
36 points: Arc<OnceLock<SearchView>>,
37}
38
39impl Default for CollectionQueryCache {
40 fn default() -> Self {
41 Self {
42 root: None,
43 points: Arc::new(OnceLock::new()),
44 }
45 }
46}
47
48impl Database {
49 pub fn init(path: impl AsRef<Path>) -> Result<Self> {
51 Self::init_with_options(path, false)
52 }
53
54 pub fn init_bare(path: impl AsRef<Path>) -> Result<Self> {
56 Self::init_with_options(path, true)
57 }
58
59 pub fn init_with_options(path: impl AsRef<Path>, bare: bool) -> Result<Self> {
61 let path = path.as_ref();
62 if bare {
63 Repository::init_bare(path)?;
64 } else {
65 Repository::init(path)?;
66 }
67 Self::open(path)
68 }
69
70 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
72 let repository = Repository::open(path.as_ref())?;
73 Ok(Self {
74 path: repository.path().to_path_buf(),
75 })
76 }
77
78 fn repo(&self) -> Result<Repository> {
79 Ok(Repository::open(&self.path)?)
80 }
81
82 pub fn create_collection(
88 &self,
89 name: impl AsRef<str>,
90 config: CollectionConfig,
91 ) -> Result<Collection> {
92 let name = name.as_ref();
93 validate_collection_name(name)?;
94 validate_config(&config)?;
95 let repo = self.repo()?;
96 let ref_name = collection_ref(name);
97 if repo.find_reference(&ref_name).is_ok() {
98 return Err(Error::CollectionExists(name.into()));
99 }
100 let root = SnapshotEngine::open(&self.path)?
101 .build(config, Vec::new())?
102 .oid();
103 let commit = create_commit(&repo, root, None, &format!("create collection {name}"))?;
104 repo.reference(&ref_name, commit, false, "git-vdb create collection")?;
105 Ok(Collection {
106 db: self.clone(),
107 name: name.into(),
108 historical: None,
109 query_cache: Default::default(),
110 })
111 }
112
113 pub fn get_or_create_collection(
117 &self,
118 name: impl AsRef<str>,
119 config: CollectionConfig,
120 ) -> Result<Collection> {
121 match self.collection(name.as_ref()) {
122 Ok(collection) => {
123 let actual = collection.info()?.config;
124 if actual != config {
125 return Err(Error::Invalid(format!(
126 "collection exists with a different configuration: {actual:?}"
127 )));
128 }
129 Ok(collection)
130 }
131 Err(Error::CollectionNotFound(_)) => self.create_collection(name, config),
132 Err(error) => Err(error),
133 }
134 }
135
136 pub fn collection(&self, name: impl AsRef<str>) -> Result<Collection> {
138 let name = name.as_ref();
139 validate_collection_name(name)?;
140 let repo = self.repo()?;
141 if repo.find_reference(&collection_ref(name)).is_err() {
142 return Err(Error::CollectionNotFound(name.into()));
143 }
144 Ok(Collection {
145 db: self.clone(),
146 name: name.into(),
147 historical: None,
148 query_cache: Default::default(),
149 })
150 }
151
152 pub fn list_collections(&self) -> Result<Vec<String>> {
154 let repo = self.repo()?;
155 let mut names = Vec::new();
156 for reference in repo.references_glob("refs/git-vdb/collections/*")? {
157 let reference = reference?;
158 if let Some(name) = reference.name()?.strip_prefix("refs/git-vdb/collections/") {
159 names.push(name.to_owned());
160 }
161 }
162 names.sort();
163 Ok(names)
164 }
165
166 pub fn delete_collection(&self, name: impl AsRef<str>) -> Result<ObjectId> {
171 let collection = self.collection(name.as_ref())?;
172 let root = collection.root()?;
173 let repo = self.repo()?;
174 repo.find_reference(&collection_ref(name.as_ref()))?
175 .delete()?;
176 Ok(root)
177 }
178}
179
180impl Collection {
181 fn repo(&self) -> Result<Repository> {
182 self.db.repo()
183 }
184
185 fn snapshot(&self, repo: &Repository) -> Result<ResolvedSnapshot> {
186 if let Some(snapshot) = self.historical {
187 return Ok(snapshot);
188 }
189 current_snapshot(repo, &self.name)
190 }
191
192 pub fn root(&self) -> Result<ObjectId> {
194 let repo = self.repo()?;
195 Ok(self.snapshot(&repo)?.root.into())
196 }
197
198 pub fn at(&self, revision: impl AsRef<str>) -> Result<Self> {
200 let repo = self.repo()?;
201 let snapshot = resolve_snapshot(&repo, revision.as_ref())?;
202 read_meta(&repo, snapshot.root)?;
203 Ok(Self {
204 db: self.db.clone(),
205 name: self.name.clone(),
206 historical: Some(snapshot),
207 query_cache: Default::default(),
208 })
209 }
210
211 pub fn info(&self) -> Result<CollectionInfo> {
213 let repo = self.repo()?;
214 let snapshot = self.snapshot(&repo)?;
215 let meta = read_meta(&repo, snapshot.root)?;
216 Ok(CollectionInfo {
217 root: snapshot.root.into(),
218 name: self.name.clone(),
219 format_version: meta.format_version(),
220 point_count: meta.point_count(),
221 config: meta.config(),
222 read_only: self.historical.is_some(),
223 })
224 }
225
226 pub fn upsert(&self, points: Vec<Point>) -> Result<WriteResult> {
231 self.upsert_expect(points, None)
232 }
233
234 pub fn upsert_expect(
239 &self,
240 points: Vec<Point>,
241 expected_root: Option<ObjectId>,
242 ) -> Result<WriteResult> {
243 if self.historical.is_some() {
244 return Err(Error::ReadOnly);
245 }
246 if points.is_empty() {
247 return Err(Error::Invalid("upsert batch must not be empty".into()));
248 }
249 let repo = self.repo()?;
250 let snapshot = current_snapshot(&repo, &self.name)?;
251 check_expected_root(snapshot.root, expected_root.as_ref())?;
252 let affected_points = points.len();
253 let mutations = points.into_iter().map(SnapshotMutation::upsert).collect();
254 let root = SnapshotEngine::open(&self.db.path)?
255 .apply(snapshot.root.to_string(), mutations)?
256 .oid();
257 advance_collection(
258 &repo,
259 &self.name,
260 snapshot,
261 root,
262 &format!("upsert {affected_points} points"),
263 )?;
264 Ok(WriteResult {
265 root: root.into(),
266 affected_points,
267 })
268 }
269
270 pub fn delete(&self, selector: DeleteSelector) -> Result<WriteResult> {
274 self.delete_expect(selector, None)
275 }
276
277 pub fn delete_expect(
279 &self,
280 selector: DeleteSelector,
281 expected_root: Option<ObjectId>,
282 ) -> Result<WriteResult> {
283 if self.historical.is_some() {
284 return Err(Error::ReadOnly);
285 }
286 if selector.ids.is_empty() && selector.filter.is_none() {
287 return Err(Error::Invalid("delete selector must not be empty".into()));
288 }
289 let repo = self.repo()?;
290 let snapshot = current_snapshot(&repo, &self.name)?;
291 check_expected_root(snapshot.root, expected_root.as_ref())?;
292 let before = read_meta(&repo, snapshot.root)?.point_count();
293 let mut mutations = Vec::new();
294 if !selector.ids.is_empty() {
295 mutations.push(SnapshotMutation::delete_ids(selector.ids));
296 }
297 if let Some(filter) = selector.filter {
298 mutations.push(SnapshotMutation::delete_filter(filter));
299 }
300 let new_snapshot =
301 SnapshotEngine::open(&self.db.path)?.apply(snapshot.root.to_string(), mutations)?;
302 let root = new_snapshot.oid();
303 let affected_points = before - new_snapshot.info()?.point_count;
304 advance_collection(
305 &repo,
306 &self.name,
307 snapshot,
308 root,
309 &format!("delete {affected_points} points"),
310 )?;
311 Ok(WriteResult {
312 root: root.into(),
313 affected_points,
314 })
315 }
316
317 pub fn get(&self, request: GetRequest) -> Result<GetResult> {
319 let repo = self.repo()?;
320 let snapshot = self.snapshot(&repo)?;
321 get_root(&repo, snapshot.root, request)
322 }
323
324 pub fn count(&self, filter: Option<Filter>) -> Result<CountResult> {
326 let repo = self.repo()?;
327 let snapshot = self.snapshot(&repo)?;
328 count_root(&repo, snapshot.root, filter)
329 }
330
331 pub fn query(&self, query: Query) -> Result<QueryResult> {
336 let repo = self.repo()?;
337 let snapshot = self.snapshot(&repo)?;
338 let points = {
339 let mut cache = self
340 .query_cache
341 .lock()
342 .map_err(|_| Error::Invalid("collection query cache lock is poisoned".into()))?;
343 if cache.root != Some(snapshot.root) {
344 cache.root = Some(snapshot.root);
345 cache.points = Arc::new(OnceLock::new());
346 }
347 cache.points.clone()
348 };
349 query_root_with_cache(&repo, snapshot.root, query, Some(&points))
350 }
351
352 pub fn history(&self, limit: usize) -> Result<Vec<HistoryEntry>> {
354 let repo = self.repo()?;
355 let mut commit_id = self
356 .snapshot(&repo)?
357 .commit
358 .ok_or_else(|| Error::Invalid("a root tree has no commit history".into()))?;
359 let mut history = Vec::new();
360 while history.len() < limit {
361 let commit = repo.find_commit(commit_id)?;
362 let parent = commit.parent_id(0).ok();
363 history.push(HistoryEntry {
364 commit: commit.id().into(),
365 root: commit.tree_id().into(),
366 parent: parent.map(Into::into),
367 message: commit.message().unwrap_or_default().to_owned(),
368 time_seconds: commit.time().seconds(),
369 });
370 let Some(parent) = parent else { break };
371 commit_id = parent;
372 }
373 Ok(history)
374 }
375
376 pub fn diff(
378 &self,
379 left_revision: impl AsRef<str>,
380 right_revision: impl AsRef<str>,
381 ) -> Result<DiffResult> {
382 let repo = self.repo()?;
383 let left = resolve_snapshot(&repo, left_revision.as_ref())?;
384 let right = resolve_snapshot(&repo, right_revision.as_ref())?;
385 diff_roots(&repo, left.root, right.root)
386 }
387
388 pub fn validate(&self, full: bool) -> Result<ValidationReport> {
392 let repo = self.repo()?;
393 let snapshot = self.snapshot(&repo)?;
394 validate_root(&repo, snapshot.root, full)
395 }
396}
397
398fn validate_collection_name(name: &str) -> Result<()> {
399 if name.is_empty()
400 || name.len() > 128
401 || !name
402 .bytes()
403 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
404 || name.starts_with('.')
405 || name.ends_with('.')
406 || name.contains("..")
407 {
408 return Err(Error::Invalid(format!("invalid collection name {name:?}")));
409 }
410 Ok(())
411}
412
413fn collection_ref(name: &str) -> String {
414 format!("refs/git-vdb/collections/{name}")
415}
416
417fn current_snapshot(repo: &Repository, name: &str) -> Result<ResolvedSnapshot> {
418 let reference = repo
419 .find_reference(&collection_ref(name))
420 .map_err(|_| Error::CollectionNotFound(name.into()))?;
421 let commit = reference.peel_to_commit()?;
422 Ok(ResolvedSnapshot {
423 root: commit.tree_id(),
424 commit: Some(commit.id()),
425 })
426}
427
428fn resolve_snapshot(repo: &Repository, revision: &str) -> Result<ResolvedSnapshot> {
429 let object = repo.revparse_single(revision)?;
430 match object.kind() {
431 Some(git2::ObjectType::Commit) => {
432 let commit = object.peel_to_commit()?;
433 Ok(ResolvedSnapshot {
434 root: commit.tree_id(),
435 commit: Some(commit.id()),
436 })
437 }
438 Some(git2::ObjectType::Tree) => Ok(ResolvedSnapshot {
439 root: object.id(),
440 commit: None,
441 }),
442 _ => Err(Error::Invalid("revision is not a commit or tree".into())),
443 }
444}
445
446fn check_expected_root(actual: Oid, expected: Option<&ObjectId>) -> Result<()> {
447 if let Some(expected) = expected {
448 if expected.0 != actual.to_string() {
449 return Err(Error::StaleRoot {
450 expected: expected.clone(),
451 actual: actual.into(),
452 });
453 }
454 }
455 Ok(())
456}
457
458fn signature() -> Result<Signature<'static>> {
459 Ok(Signature::now("git-vdb", "git-vdb@localhost")?)
460}
461
462fn create_commit(
463 repo: &Repository,
464 root: Oid,
465 parent: Option<&Commit<'_>>,
466 message: &str,
467) -> Result<Oid> {
468 let tree = repo.find_tree(root)?;
469 let signature = signature()?;
470 let parents: Vec<&Commit<'_>> = parent.into_iter().collect();
471 Ok(repo.commit(None, &signature, &signature, message, &tree, &parents)?)
472}
473
474fn advance_collection(
475 repo: &Repository,
476 name: &str,
477 old: ResolvedSnapshot,
478 root: Oid,
479 message: &str,
480) -> Result<()> {
481 let old_commit_id = old
482 .commit
483 .ok_or_else(|| Error::Corrupt("current collection ref is not a commit".into()))?;
484 let old_commit = repo.find_commit(old_commit_id)?;
485 let new_commit = create_commit(repo, root, Some(&old_commit), message)?;
486 repo.reference_matching(
487 &collection_ref(name),
488 new_commit,
489 true,
490 old_commit_id,
491 "git-vdb atomic collection update",
492 )
493 .map_err(|error| {
494 if error.code() == git2::ErrorCode::Modified {
495 let actual = current_snapshot(repo, name)
496 .map(|snapshot| snapshot.root.into())
497 .unwrap_or_else(|_| ObjectId("unknown".into()));
498 Error::StaleRoot {
499 expected: old.root.into(),
500 actual,
501 }
502 } else {
503 Error::Git(error)
504 }
505 })?;
506 Ok(())
507}