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 validate_collection_name(name: &str) -> Result<()> {
51 validate_collection_name(name)
52 }
53
54 pub fn init(path: impl AsRef<Path>) -> Result<Self> {
56 Self::init_with_options(path, false)
57 }
58
59 pub fn init_bare(path: impl AsRef<Path>) -> Result<Self> {
61 Self::init_with_options(path, true)
62 }
63
64 pub fn init_with_options(path: impl AsRef<Path>, bare: bool) -> Result<Self> {
66 let path = path.as_ref();
67 if bare {
68 Repository::init_bare(path)?;
69 } else {
70 Repository::init(path)?;
71 }
72 Self::open(path)
73 }
74
75 pub fn open(path: impl AsRef<Path>) -> Result<Self> {
77 let repository = Repository::open(path.as_ref())?;
78 Ok(Self {
79 path: repository.path().to_path_buf(),
80 })
81 }
82
83 fn repo(&self) -> Result<Repository> {
84 Ok(Repository::open(&self.path)?)
85 }
86
87 pub fn path(&self) -> &Path {
89 &self.path
90 }
91
92 pub fn create_collection(
98 &self,
99 name: impl AsRef<str>,
100 config: CollectionConfig,
101 ) -> Result<Collection> {
102 let name = name.as_ref();
103 validate_collection_name(name)?;
104 validate_config(&config)?;
105 let repo = self.repo()?;
106 let ref_name = collection_ref(name);
107 if repo.find_reference(&ref_name).is_ok() {
108 return Err(Error::CollectionExists(name.into()));
109 }
110 let root = SnapshotEngine::open(&self.path)?
111 .build(config, Vec::new())?
112 .oid();
113 let commit = create_commit(&repo, root, None, &format!("create collection {name}"))?;
114 repo.reference(&ref_name, commit, false, "git-vdb create collection")?;
115 Ok(Collection {
116 db: self.clone(),
117 name: name.into(),
118 historical: None,
119 query_cache: Default::default(),
120 })
121 }
122
123 pub fn get_or_create_collection(
127 &self,
128 name: impl AsRef<str>,
129 config: CollectionConfig,
130 ) -> Result<Collection> {
131 match self.collection(name.as_ref()) {
132 Ok(collection) => {
133 let actual = collection.info()?.config;
134 if actual != config {
135 return Err(Error::Invalid(format!(
136 "collection exists with a different configuration: {actual:?}"
137 )));
138 }
139 Ok(collection)
140 }
141 Err(Error::CollectionNotFound(_)) => self.create_collection(name, config),
142 Err(error) => Err(error),
143 }
144 }
145
146 pub fn collection(&self, name: impl AsRef<str>) -> Result<Collection> {
148 let name = name.as_ref();
149 validate_collection_name(name)?;
150 let repo = self.repo()?;
151 if repo.find_reference(&collection_ref(name)).is_err() {
152 return Err(Error::CollectionNotFound(name.into()));
153 }
154 Ok(Collection {
155 db: self.clone(),
156 name: name.into(),
157 historical: None,
158 query_cache: Default::default(),
159 })
160 }
161
162 pub fn list_collections(&self) -> Result<Vec<String>> {
164 let repo = self.repo()?;
165 let mut names = Vec::new();
166 for reference in repo.references_glob("refs/git-vdb/collections/*")? {
167 let reference = reference?;
168 if let Some(name) = reference.name()?.strip_prefix("refs/git-vdb/collections/") {
169 names.push(name.to_owned());
170 }
171 }
172 names.sort();
173 Ok(names)
174 }
175
176 pub fn delete_collection(&self, name: impl AsRef<str>) -> Result<ObjectId> {
181 let collection = self.collection(name.as_ref())?;
182 let root = collection.root()?;
183 let repo = self.repo()?;
184 repo.find_reference(&collection_ref(name.as_ref()))?
185 .delete()?;
186 Ok(root)
187 }
188}
189
190impl Collection {
191 fn repo(&self) -> Result<Repository> {
192 self.db.repo()
193 }
194
195 fn snapshot(&self, repo: &Repository) -> Result<ResolvedSnapshot> {
196 if let Some(snapshot) = self.historical {
197 return Ok(snapshot);
198 }
199 current_snapshot(repo, &self.name)
200 }
201
202 pub fn root(&self) -> Result<ObjectId> {
204 let repo = self.repo()?;
205 Ok(self.snapshot(&repo)?.root.into())
206 }
207
208 pub fn at(&self, revision: impl AsRef<str>) -> Result<Self> {
210 let repo = self.repo()?;
211 let snapshot = resolve_snapshot(&repo, revision.as_ref())?;
212 read_meta(&repo, snapshot.root)?;
213 Ok(Self {
214 db: self.db.clone(),
215 name: self.name.clone(),
216 historical: Some(snapshot),
217 query_cache: Default::default(),
218 })
219 }
220
221 pub fn info(&self) -> Result<CollectionInfo> {
223 let repo = self.repo()?;
224 let snapshot = self.snapshot(&repo)?;
225 let meta = read_meta(&repo, snapshot.root)?;
226 Ok(CollectionInfo {
227 root: snapshot.root.into(),
228 name: self.name.clone(),
229 format_version: meta.format_version(),
230 point_count: meta.point_count(),
231 config: meta.config(),
232 read_only: self.historical.is_some(),
233 })
234 }
235
236 pub fn upsert(&self, points: Vec<Point>) -> Result<WriteResult> {
241 self.upsert_expect(points, None)
242 }
243
244 pub fn upsert_expect(
249 &self,
250 points: Vec<Point>,
251 expected_root: Option<ObjectId>,
252 ) -> Result<WriteResult> {
253 if self.historical.is_some() {
254 return Err(Error::ReadOnly);
255 }
256 if points.is_empty() {
257 return Err(Error::Invalid("upsert batch must not be empty".into()));
258 }
259 let repo = self.repo()?;
260 let snapshot = current_snapshot(&repo, &self.name)?;
261 check_expected_root(snapshot.root, expected_root.as_ref())?;
262 let affected_points = points.len();
263 let mutations = points.into_iter().map(SnapshotMutation::upsert).collect();
264 let root = SnapshotEngine::open(&self.db.path)?
265 .apply(snapshot.root.to_string(), mutations)?
266 .oid();
267 advance_collection(
268 &repo,
269 &self.name,
270 snapshot,
271 root,
272 &format!("upsert {affected_points} points"),
273 )?;
274 Ok(WriteResult {
275 root: root.into(),
276 affected_points,
277 })
278 }
279
280 pub fn apply(&self, mutations: Vec<SnapshotMutation>) -> Result<MutationResult> {
285 if self.historical.is_some() {
286 return Err(Error::ReadOnly);
287 }
288 if mutations.is_empty() {
289 return Err(Error::Invalid("mutation batch must not be empty".into()));
290 }
291 let repo = self.repo()?;
292 let snapshot = current_snapshot(&repo, &self.name)?;
293 let points_before = read_meta(&repo, snapshot.root)?.point_count();
294 let operations = mutations.len();
295 let next =
296 SnapshotEngine::open(&self.db.path)?.apply(snapshot.root.to_string(), mutations)?;
297 let root = next.oid();
298 let points_after = next.info()?.point_count;
299 advance_collection(
300 &repo,
301 &self.name,
302 snapshot,
303 root,
304 &format!("apply {operations} mutations"),
305 )?;
306 Ok(MutationResult {
307 root: root.into(),
308 points_before,
309 points_after,
310 operations,
311 })
312 }
313
314 pub fn restore(&self, revision: impl AsRef<str>) -> Result<WriteResult> {
318 if self.historical.is_some() {
319 return Err(Error::ReadOnly);
320 }
321 let repo = self.repo()?;
322 let current = current_snapshot(&repo, &self.name)?;
323 let target = resolve_snapshot(&repo, revision.as_ref())?;
324 let diff = diff_roots(&repo, current.root, target.root)?;
325 let affected_points = diff.added.len() + diff.removed.len() + diff.changed.len();
326 advance_collection(
327 &repo,
328 &self.name,
329 current,
330 target.root,
331 &format!("restore {}", revision.as_ref()),
332 )?;
333 Ok(WriteResult {
334 root: target.root.into(),
335 affected_points,
336 })
337 }
338
339 pub fn delete(&self, selector: DeleteSelector) -> Result<WriteResult> {
343 self.delete_expect(selector, None)
344 }
345
346 pub fn delete_expect(
348 &self,
349 selector: DeleteSelector,
350 expected_root: Option<ObjectId>,
351 ) -> Result<WriteResult> {
352 if self.historical.is_some() {
353 return Err(Error::ReadOnly);
354 }
355 if selector.ids.is_empty() && selector.filter.is_none() {
356 return Err(Error::Invalid("delete selector must not be empty".into()));
357 }
358 let repo = self.repo()?;
359 let snapshot = current_snapshot(&repo, &self.name)?;
360 check_expected_root(snapshot.root, expected_root.as_ref())?;
361 let before = read_meta(&repo, snapshot.root)?.point_count();
362 let mut mutations = Vec::new();
363 if !selector.ids.is_empty() {
364 mutations.push(SnapshotMutation::delete_ids(selector.ids));
365 }
366 if let Some(filter) = selector.filter {
367 mutations.push(SnapshotMutation::delete_filter(filter));
368 }
369 let new_snapshot =
370 SnapshotEngine::open(&self.db.path)?.apply(snapshot.root.to_string(), mutations)?;
371 let root = new_snapshot.oid();
372 let affected_points = before - new_snapshot.info()?.point_count;
373 advance_collection(
374 &repo,
375 &self.name,
376 snapshot,
377 root,
378 &format!("delete {affected_points} points"),
379 )?;
380 Ok(WriteResult {
381 root: root.into(),
382 affected_points,
383 })
384 }
385
386 pub fn get(&self, request: GetRequest) -> Result<GetResult> {
388 let repo = self.repo()?;
389 let snapshot = self.snapshot(&repo)?;
390 get_root(&repo, snapshot.root, request)
391 }
392
393 pub fn count(&self, filter: Option<Filter>) -> Result<CountResult> {
395 let repo = self.repo()?;
396 let snapshot = self.snapshot(&repo)?;
397 count_root(&repo, snapshot.root, filter)
398 }
399
400 pub fn query(&self, query: Query) -> Result<QueryResult> {
405 let repo = self.repo()?;
406 let snapshot = self.snapshot(&repo)?;
407 let points = {
408 let mut cache = self
409 .query_cache
410 .lock()
411 .map_err(|_| Error::Invalid("collection query cache lock is poisoned".into()))?;
412 if cache.root != Some(snapshot.root) {
413 cache.root = Some(snapshot.root);
414 cache.points = Arc::new(OnceLock::new());
415 }
416 cache.points.clone()
417 };
418 query_root_with_cache(&repo, snapshot.root, query, Some(&points))
419 }
420
421 pub fn history(&self, limit: usize) -> Result<Vec<HistoryEntry>> {
423 let repo = self.repo()?;
424 let mut commit_id = self
425 .snapshot(&repo)?
426 .commit
427 .ok_or_else(|| Error::Invalid("a root tree has no commit history".into()))?;
428 let mut history = Vec::new();
429 while history.len() < limit {
430 let commit = repo.find_commit(commit_id)?;
431 let parent = commit.parent_id(0).ok();
432 history.push(HistoryEntry {
433 commit: commit.id().into(),
434 root: commit.tree_id().into(),
435 parent: parent.map(Into::into),
436 message: commit.message().unwrap_or_default().to_owned(),
437 time_seconds: commit.time().seconds(),
438 });
439 let Some(parent) = parent else { break };
440 commit_id = parent;
441 }
442 Ok(history)
443 }
444
445 pub fn diff(
447 &self,
448 left_revision: impl AsRef<str>,
449 right_revision: impl AsRef<str>,
450 ) -> Result<DiffResult> {
451 let repo = self.repo()?;
452 let left = resolve_snapshot(&repo, left_revision.as_ref())?;
453 let right = resolve_snapshot(&repo, right_revision.as_ref())?;
454 diff_roots(&repo, left.root, right.root)
455 }
456
457 pub fn validate(&self, full: bool) -> Result<ValidationReport> {
461 let repo = self.repo()?;
462 let snapshot = self.snapshot(&repo)?;
463 validate_root(&repo, snapshot.root, full)
464 }
465}
466
467fn validate_collection_name(name: &str) -> Result<()> {
468 if name.is_empty()
469 || name.len() > 128
470 || !name
471 .bytes()
472 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
473 || name.starts_with('.')
474 || name.ends_with('.')
475 || name.contains("..")
476 {
477 return Err(Error::Invalid(format!("invalid collection name {name:?}")));
478 }
479 Ok(())
480}
481
482fn collection_ref(name: &str) -> String {
483 format!("refs/git-vdb/collections/{name}")
484}
485
486fn current_snapshot(repo: &Repository, name: &str) -> Result<ResolvedSnapshot> {
487 let reference = repo
488 .find_reference(&collection_ref(name))
489 .map_err(|_| Error::CollectionNotFound(name.into()))?;
490 let commit = reference.peel_to_commit()?;
491 Ok(ResolvedSnapshot {
492 root: commit.tree_id(),
493 commit: Some(commit.id()),
494 })
495}
496
497fn resolve_snapshot(repo: &Repository, revision: &str) -> Result<ResolvedSnapshot> {
498 let object = repo.revparse_single(revision)?;
499 match object.kind() {
500 Some(git2::ObjectType::Commit) => {
501 let commit = object.peel_to_commit()?;
502 Ok(ResolvedSnapshot {
503 root: commit.tree_id(),
504 commit: Some(commit.id()),
505 })
506 }
507 Some(git2::ObjectType::Tree) => Ok(ResolvedSnapshot {
508 root: object.id(),
509 commit: None,
510 }),
511 _ => Err(Error::Invalid("revision is not a commit or tree".into())),
512 }
513}
514
515fn check_expected_root(actual: Oid, expected: Option<&ObjectId>) -> Result<()> {
516 if let Some(expected) = expected {
517 if expected.0 != actual.to_string() {
518 return Err(Error::StaleRoot {
519 expected: expected.clone(),
520 actual: actual.into(),
521 });
522 }
523 }
524 Ok(())
525}
526
527fn signature() -> Result<Signature<'static>> {
528 Ok(Signature::now("git-vdb", "git-vdb@localhost")?)
529}
530
531fn create_commit(
532 repo: &Repository,
533 root: Oid,
534 parent: Option<&Commit<'_>>,
535 message: &str,
536) -> Result<Oid> {
537 let tree = repo.find_tree(root)?;
538 let signature = signature()?;
539 let parents: Vec<&Commit<'_>> = parent.into_iter().collect();
540 Ok(repo.commit(None, &signature, &signature, message, &tree, &parents)?)
541}
542
543fn advance_collection(
544 repo: &Repository,
545 name: &str,
546 old: ResolvedSnapshot,
547 root: Oid,
548 message: &str,
549) -> Result<()> {
550 let old_commit_id = old
551 .commit
552 .ok_or_else(|| Error::Corrupt("current collection ref is not a commit".into()))?;
553 let old_commit = repo.find_commit(old_commit_id)?;
554 let new_commit = create_commit(repo, root, Some(&old_commit), message)?;
555 repo.reference_matching(
556 &collection_ref(name),
557 new_commit,
558 true,
559 old_commit_id,
560 "git-vdb atomic collection update",
561 )
562 .map_err(|error| {
563 if error.code() == git2::ErrorCode::Modified {
564 let actual = current_snapshot(repo, name)
565 .map(|snapshot| snapshot.root.into())
566 .unwrap_or_else(|_| ObjectId("unknown".into()));
567 Error::StaleRoot {
568 expected: old.root.into(),
569 actual,
570 }
571 } else {
572 Error::Git(error)
573 }
574 })?;
575 Ok(())
576}