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            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    /// Adds or replaces a non-empty point batch at the current collection root.
227    ///
228    /// The ref is advanced with compare-and-swap semantics. Use
229    /// [`Collection::upsert_expect`] to supply an explicit expected root.
230    pub fn upsert(&self, points: Vec<Point>) -> Result<WriteResult> {
231        self.upsert_expect(points, None)
232    }
233
234    /// Adds or replaces points only if the current root matches `expected_root`.
235    ///
236    /// All immutable objects are written before the collection ref is advanced.
237    /// A stale precondition returns [`Error::StaleRoot`] without advancing it.
238    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    /// Deletes the union of selected IDs and filter matches.
271    ///
272    /// The selector must contain at least one ID or a filter.
273    pub fn delete(&self, selector: DeleteSelector) -> Result<WriteResult> {
274        self.delete_expect(selector, None)
275    }
276
277    /// Deletes selected points only when the current root matches a precondition.
278    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    /// Retrieves canonically ordered records without similarity scoring.
318    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    /// Counts all points or those matching a filter.
325    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    /// Executes an exact or deterministic approximate vector query.
332    ///
333    /// The returned result identifies the root actually read. Query caches are
334    /// scoped to that immutable root and are replaced when the collection moves.
335    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    /// Returns at most `limit` collection commits, newest first.
353    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    /// Compares logical points and structural sharing between two revisions.
377    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    /// Validates the resolved root without modifying objects or refs.
389    ///
390    /// Full validation recomputes every approximate-index bucket.
391    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}