Skip to main content

git_vdb/
adapter.rs

1//! Named collections backed by Git commits and compare-and-swap refs.
2
3use 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/// A Git repository that manages named, mutable collection refs.
19#[derive(Clone, Debug)]
20pub struct Database {
21    pub(crate) path: PathBuf,
22}
23
24/// A named collection view backed by the Git commit/ref adapter.
25#[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    /// Validates a collection name without opening or creating it.
50    pub fn validate_collection_name(name: &str) -> Result<()> {
51        validate_collection_name(name)
52    }
53
54    /// Initializes a non-bare Git repository and opens it as a database.
55    pub fn init(path: impl AsRef<Path>) -> Result<Self> {
56        Self::init_with_options(path, false)
57    }
58
59    /// Initializes a bare Git repository and opens it as a database.
60    pub fn init_bare(path: impl AsRef<Path>) -> Result<Self> {
61        Self::init_with_options(path, true)
62    }
63
64    /// Initializes either a bare or non-bare Git repository.
65    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    /// Opens an existing bare or non-bare Git repository.
76    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    /// Returns the resolved Git directory used by this database.
88    pub fn path(&self) -> &Path {
89        &self.path
90    }
91
92    /// Creates an empty named collection with canonical configuration.
93    ///
94    /// The method writes the initial immutable root and commit, then creates the
95    /// collection ref. It returns [`Error::CollectionExists`] when the name is
96    /// already present.
97    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    /// Opens a collection or atomically creates it with the supplied config.
124    ///
125    /// An existing collection with different configuration is rejected.
126    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    /// Opens the current mutable view of a named collection.
147    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    /// Returns collection names in ascending byte order.
163    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    /// Deletes only a named collection ref and returns its last root.
177    ///
178    /// Commits and trees remain subject to ordinary Git reachability and
179    /// garbage collection.
180    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    /// Returns the deterministic root resolved by this collection view.
203    pub fn root(&self) -> Result<ObjectId> {
204        let repo = self.repo()?;
205        Ok(self.snapshot(&repo)?.root.into())
206    }
207
208    /// Creates a read-only historical view at a root, commit, or revision.
209    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    /// Returns metadata for the root resolved by this collection view.
222    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    /// Adds or replaces a non-empty point batch at the current collection root.
237    ///
238    /// The ref is advanced with compare-and-swap semantics. Use
239    /// [`Collection::upsert_expect`] to supply an explicit expected root.
240    pub fn upsert(&self, points: Vec<Point>) -> Result<WriteResult> {
241        self.upsert_expect(points, None)
242    }
243
244    /// Adds or replaces points only if the current root matches `expected_root`.
245    ///
246    /// All immutable objects are written before the collection ref is advanced.
247    /// A stale precondition returns [`Error::StaleRoot`] without advancing it.
248    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    /// Applies an ordered non-empty batch of upserts and deletions atomically.
281    ///
282    /// Every mutation observes the result of the preceding mutation. Immutable
283    /// objects are built first and the collection ref advances exactly once.
284    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    /// Restores a historical root as a new commit at the collection tip.
315    ///
316    /// History is preserved: this never rewinds or deletes existing commits.
317    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    /// Deletes the union of selected IDs and filter matches.
340    ///
341    /// The selector must contain at least one ID or a filter.
342    pub fn delete(&self, selector: DeleteSelector) -> Result<WriteResult> {
343        self.delete_expect(selector, None)
344    }
345
346    /// Deletes selected points only when the current root matches a precondition.
347    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    /// Retrieves canonically ordered records without similarity scoring.
387    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    /// Counts all points or those matching a filter.
394    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    /// Executes an exact or deterministic approximate vector query.
401    ///
402    /// The returned result identifies the root actually read. Query caches are
403    /// scoped to that immutable root and are replaced when the collection moves.
404    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    /// Returns at most `limit` collection commits, newest first.
422    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    /// Compares logical points and structural sharing between two revisions.
446    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    /// Validates the resolved root without modifying objects or refs.
458    ///
459    /// Full validation recomputes every approximate-index bucket.
460    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}