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    /// Initializes a non-bare Git repository and opens it as a database.
50    pub fn init(path: impl AsRef<Path>) -> Result<Self> {
51        Self::init_with_options(path, false)
52    }
53
54    /// Initializes a bare Git repository and opens it as a database.
55    pub fn init_bare(path: impl AsRef<Path>) -> Result<Self> {
56        Self::init_with_options(path, true)
57    }
58
59    /// Initializes either a bare or non-bare Git repository.
60    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    /// Opens an existing bare or non-bare Git repository.
71    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    /// Creates an empty named collection with canonical configuration.
83    ///
84    /// The method writes the initial immutable root and commit, then creates the
85    /// collection ref. It returns [`Error::CollectionExists`] when the name is
86    /// already present.
87    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    /// Opens a collection or atomically creates it with the supplied config.
114    ///
115    /// An existing collection with different configuration is rejected.
116    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    /// Opens the current mutable view of a named collection.
137    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    /// Returns collection names in ascending byte order.
153    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    /// Deletes only a named collection ref and returns its last root.
167    ///
168    /// Commits and trees remain subject to ordinary Git reachability and
169    /// garbage collection.
170    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    /// Returns the deterministic root resolved by this collection view.
193    pub fn root(&self) -> Result<ObjectId> {
194        let repo = self.repo()?;
195        Ok(self.snapshot(&repo)?.root.into())
196    }
197
198    /// Creates a read-only historical view at a root, commit, or revision.
199    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    /// Returns metadata for the root resolved by this collection view.
212    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    /// Adds or replaces a non-empty point batch at the current collection root.
226    ///
227    /// The ref is advanced with compare-and-swap semantics. Use
228    /// [`Collection::upsert_expect`] to supply an explicit expected root.
229    pub fn upsert(&self, points: Vec<Point>) -> Result<WriteResult> {
230        self.upsert_expect(points, None)
231    }
232
233    /// Adds or replaces points only if the current root matches `expected_root`.
234    ///
235    /// All immutable objects are written before the collection ref is advanced.
236    /// A stale precondition returns [`Error::StaleRoot`] without advancing it.
237    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    /// Deletes the union of selected IDs and filter matches.
270    ///
271    /// The selector must contain at least one ID or a filter.
272    pub fn delete(&self, selector: DeleteSelector) -> Result<WriteResult> {
273        self.delete_expect(selector, None)
274    }
275
276    /// Deletes selected points only when the current root matches a precondition.
277    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    /// Retrieves canonically ordered records without similarity scoring.
317    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    /// Counts all points or those matching a filter.
324    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    /// Executes an exact or deterministic approximate vector query.
331    ///
332    /// The returned result identifies the root actually read. Query caches are
333    /// scoped to that immutable root and are replaced when the collection moves.
334    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    /// Returns at most `limit` collection commits, newest first.
352    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    /// Compares logical points and structural sharing between two revisions.
376    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    /// Validates the resolved root without modifying objects or refs.
388    ///
389    /// Full validation recomputes every approximate-index bucket.
390    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}