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 point_count: meta.point_count(),
220 config: meta.config(),
221 read_only: self.historical.is_some(),
222 })
223 }
224
225 pub fn upsert(&self, points: Vec<Point>) -> Result<WriteResult> {
230 self.upsert_expect(points, None)
231 }
232
233 pub fn upsert_expect(
238 &self,
239 points: Vec<Point>,
240 expected_root: Option<ObjectId>,
241 ) -> Result<WriteResult> {
242 if self.historical.is_some() {
243 return Err(Error::ReadOnly);
244 }
245 if points.is_empty() {
246 return Err(Error::Invalid("upsert batch must not be empty".into()));
247 }
248 let repo = self.repo()?;
249 let snapshot = current_snapshot(&repo, &self.name)?;
250 check_expected_root(snapshot.root, expected_root.as_ref())?;
251 let affected_points = points.len();
252 let mutations = points.into_iter().map(SnapshotMutation::upsert).collect();
253 let root = SnapshotEngine::open(&self.db.path)?
254 .apply(snapshot.root.to_string(), mutations)?
255 .oid();
256 advance_collection(
257 &repo,
258 &self.name,
259 snapshot,
260 root,
261 &format!("upsert {affected_points} points"),
262 )?;
263 Ok(WriteResult {
264 root: root.into(),
265 affected_points,
266 })
267 }
268
269 pub fn delete(&self, selector: DeleteSelector) -> Result<WriteResult> {
273 self.delete_expect(selector, None)
274 }
275
276 pub fn delete_expect(
278 &self,
279 selector: DeleteSelector,
280 expected_root: Option<ObjectId>,
281 ) -> Result<WriteResult> {
282 if self.historical.is_some() {
283 return Err(Error::ReadOnly);
284 }
285 if selector.ids.is_empty() && selector.filter.is_none() {
286 return Err(Error::Invalid("delete selector must not be empty".into()));
287 }
288 let repo = self.repo()?;
289 let snapshot = current_snapshot(&repo, &self.name)?;
290 check_expected_root(snapshot.root, expected_root.as_ref())?;
291 let before = read_meta(&repo, snapshot.root)?.point_count();
292 let mut mutations = Vec::new();
293 if !selector.ids.is_empty() {
294 mutations.push(SnapshotMutation::delete_ids(selector.ids));
295 }
296 if let Some(filter) = selector.filter {
297 mutations.push(SnapshotMutation::delete_filter(filter));
298 }
299 let new_snapshot =
300 SnapshotEngine::open(&self.db.path)?.apply(snapshot.root.to_string(), mutations)?;
301 let root = new_snapshot.oid();
302 let affected_points = before - new_snapshot.info()?.point_count;
303 advance_collection(
304 &repo,
305 &self.name,
306 snapshot,
307 root,
308 &format!("delete {affected_points} points"),
309 )?;
310 Ok(WriteResult {
311 root: root.into(),
312 affected_points,
313 })
314 }
315
316 pub fn get(&self, request: GetRequest) -> Result<GetResult> {
318 let repo = self.repo()?;
319 let snapshot = self.snapshot(&repo)?;
320 get_root(&repo, snapshot.root, request)
321 }
322
323 pub fn count(&self, filter: Option<Filter>) -> Result<CountResult> {
325 let repo = self.repo()?;
326 let snapshot = self.snapshot(&repo)?;
327 count_root(&repo, snapshot.root, filter)
328 }
329
330 pub fn query(&self, query: Query) -> Result<QueryResult> {
335 let repo = self.repo()?;
336 let snapshot = self.snapshot(&repo)?;
337 let points = {
338 let mut cache = self
339 .query_cache
340 .lock()
341 .map_err(|_| Error::Invalid("collection query cache lock is poisoned".into()))?;
342 if cache.root != Some(snapshot.root) {
343 cache.root = Some(snapshot.root);
344 cache.points = Arc::new(OnceLock::new());
345 }
346 cache.points.clone()
347 };
348 query_root_with_cache(&repo, snapshot.root, query, Some(&points))
349 }
350
351 pub fn history(&self, limit: usize) -> Result<Vec<HistoryEntry>> {
353 let repo = self.repo()?;
354 let mut commit_id = self
355 .snapshot(&repo)?
356 .commit
357 .ok_or_else(|| Error::Invalid("a root tree has no commit history".into()))?;
358 let mut history = Vec::new();
359 while history.len() < limit {
360 let commit = repo.find_commit(commit_id)?;
361 let parent = commit.parent_id(0).ok();
362 history.push(HistoryEntry {
363 commit: commit.id().into(),
364 root: commit.tree_id().into(),
365 parent: parent.map(Into::into),
366 message: commit.message().unwrap_or_default().to_owned(),
367 time_seconds: commit.time().seconds(),
368 });
369 let Some(parent) = parent else { break };
370 commit_id = parent;
371 }
372 Ok(history)
373 }
374
375 pub fn diff(
377 &self,
378 left_revision: impl AsRef<str>,
379 right_revision: impl AsRef<str>,
380 ) -> Result<DiffResult> {
381 let repo = self.repo()?;
382 let left = resolve_snapshot(&repo, left_revision.as_ref())?;
383 let right = resolve_snapshot(&repo, right_revision.as_ref())?;
384 diff_roots(&repo, left.root, right.root)
385 }
386
387 pub fn validate(&self, full: bool) -> Result<ValidationReport> {
391 let repo = self.repo()?;
392 let snapshot = self.snapshot(&repo)?;
393 validate_root(&repo, snapshot.root, full)
394 }
395}
396
397fn validate_collection_name(name: &str) -> Result<()> {
398 if name.is_empty()
399 || name.len() > 128
400 || !name
401 .bytes()
402 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
403 || name.starts_with('.')
404 || name.ends_with('.')
405 || name.contains("..")
406 {
407 return Err(Error::Invalid(format!("invalid collection name {name:?}")));
408 }
409 Ok(())
410}
411
412fn collection_ref(name: &str) -> String {
413 format!("refs/git-vdb/collections/{name}")
414}
415
416fn current_snapshot(repo: &Repository, name: &str) -> Result<ResolvedSnapshot> {
417 let reference = repo
418 .find_reference(&collection_ref(name))
419 .map_err(|_| Error::CollectionNotFound(name.into()))?;
420 let commit = reference.peel_to_commit()?;
421 Ok(ResolvedSnapshot {
422 root: commit.tree_id(),
423 commit: Some(commit.id()),
424 })
425}
426
427fn resolve_snapshot(repo: &Repository, revision: &str) -> Result<ResolvedSnapshot> {
428 let object = repo.revparse_single(revision)?;
429 match object.kind() {
430 Some(git2::ObjectType::Commit) => {
431 let commit = object.peel_to_commit()?;
432 Ok(ResolvedSnapshot {
433 root: commit.tree_id(),
434 commit: Some(commit.id()),
435 })
436 }
437 Some(git2::ObjectType::Tree) => Ok(ResolvedSnapshot {
438 root: object.id(),
439 commit: None,
440 }),
441 _ => Err(Error::Invalid("revision is not a commit or tree".into())),
442 }
443}
444
445fn check_expected_root(actual: Oid, expected: Option<&ObjectId>) -> Result<()> {
446 if let Some(expected) = expected {
447 if expected.0 != actual.to_string() {
448 return Err(Error::StaleRoot {
449 expected: expected.clone(),
450 actual: actual.into(),
451 });
452 }
453 }
454 Ok(())
455}
456
457fn signature() -> Result<Signature<'static>> {
458 Ok(Signature::now("git-vdb", "git-vdb@localhost")?)
459}
460
461fn create_commit(
462 repo: &Repository,
463 root: Oid,
464 parent: Option<&Commit<'_>>,
465 message: &str,
466) -> Result<Oid> {
467 let tree = repo.find_tree(root)?;
468 let signature = signature()?;
469 let parents: Vec<&Commit<'_>> = parent.into_iter().collect();
470 Ok(repo.commit(None, &signature, &signature, message, &tree, &parents)?)
471}
472
473fn advance_collection(
474 repo: &Repository,
475 name: &str,
476 old: ResolvedSnapshot,
477 root: Oid,
478 message: &str,
479) -> Result<()> {
480 let old_commit_id = old
481 .commit
482 .ok_or_else(|| Error::Corrupt("current collection ref is not a commit".into()))?;
483 let old_commit = repo.find_commit(old_commit_id)?;
484 let new_commit = create_commit(repo, root, Some(&old_commit), message)?;
485 repo.reference_matching(
486 &collection_ref(name),
487 new_commit,
488 true,
489 old_commit_id,
490 "git-vdb atomic collection update",
491 )
492 .map_err(|error| {
493 if error.code() == git2::ErrorCode::Modified {
494 let actual = current_snapshot(repo, name)
495 .map(|snapshot| snapshot.root.into())
496 .unwrap_or_else(|_| ObjectId("unknown".into()));
497 Error::StaleRoot {
498 expected: old.root.into(),
499 actual,
500 }
501 } else {
502 Error::Git(error)
503 }
504 })?;
505 Ok(())
506}