1use async_trait::async_trait;
12use gonzalo_core::{
13 Body, ContentHash, CoreError, DeleteResult, Identity, KeyPrefix, MergeOutcome, Meta, PutResult,
14 Record, RecordKey, Result, Revision, decode_segment, merge, record_components, store::Conflict,
15};
16use rustix::fs::{FlockOperation, flock};
17use std::cell::RefCell;
18use std::collections::HashSet;
19use std::path::{Path, PathBuf};
20use std::rc::Rc;
21use std::sync::Arc;
22
23mod diff;
24pub use diff::{ChangedPaths, changed_paths, head_commit, is_git_repo};
25
26#[derive(Debug, Clone)]
29pub struct PullConflict {
30 pub key: RecordKey,
31 pub local: Box<Record>,
32 pub remote: Box<Record>,
33}
34
35#[derive(Debug, Default)]
37#[must_use = "a PullReport may contain unresolved conflicts that must be handled"]
38pub struct PullReport {
39 pub fast_forwarded: bool,
41 pub merged: Vec<RecordKey>,
43 pub conflicts: Vec<PullConflict>,
45}
46
47pub struct GitStore {
48 root: PathBuf,
49}
50
51impl GitStore {
52 pub fn open(root: impl Into<PathBuf>) -> Result<Self> {
54 let root = root.into();
55 std::fs::create_dir_all(&root).map_err(|e| CoreError::Backend(e.to_string()))?;
56 match git2::Repository::open(&root) {
57 Ok(_) => {}
58 Err(_) => {
59 git2::Repository::init(&root).map_err(|e| CoreError::Backend(e.to_string()))?;
60 }
61 }
62 Ok(Self { root })
63 }
64
65 fn path_for(&self, key: &RecordKey) -> PathBuf {
66 let (ns, col, file) = record_components(key);
67 self.root.join(ns).join(col).join(file)
68 }
69
70 fn read(&self, key: &RecordKey) -> Result<Option<Record>> {
71 let path = self.path_for(key);
72 match std::fs::read(&path) {
73 Ok(bytes) => Ok(Some(
74 serde_json::from_slice(&bytes).map_err(|e| CoreError::Serde(e.to_string()))?,
75 )),
76 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
77 Err(e) => Err(CoreError::Backend(e.to_string())),
78 }
79 }
80
81 fn commit_file(&self, rel: &Path, message: &str) -> Result<()> {
82 let repo =
83 git2::Repository::open(&self.root).map_err(|e| CoreError::Backend(e.to_string()))?;
84 let mut index = repo
85 .index()
86 .map_err(|e| CoreError::Backend(e.to_string()))?;
87 index
88 .add_path(rel)
89 .map_err(|e| CoreError::Backend(e.to_string()))?;
90 index
91 .write()
92 .map_err(|e| CoreError::Backend(e.to_string()))?;
93 let tree_oid = index
94 .write_tree()
95 .map_err(|e| CoreError::Backend(e.to_string()))?;
96 let tree = repo
97 .find_tree(tree_oid)
98 .map_err(|e| CoreError::Backend(e.to_string()))?;
99 let sig = git2::Signature::now("gonzalo", "gonzalo@localhost")
100 .map_err(|e| CoreError::Backend(e.to_string()))?;
101 let parent = repo
102 .head()
103 .ok()
104 .and_then(|h| h.target())
105 .and_then(|oid| repo.find_commit(oid).ok());
106 let parents: Vec<&git2::Commit> = parent.iter().collect();
107 repo.commit(Some("HEAD"), &sig, &sig, message, &tree, &parents)
108 .map_err(|e| CoreError::Backend(e.to_string()))?;
109 Ok(())
110 }
111
112 fn commit_removal(&self, rel: &Path, message: &str) -> Result<()> {
116 let repo =
117 git2::Repository::open(&self.root).map_err(|e| CoreError::Backend(e.to_string()))?;
118 let mut index = repo
119 .index()
120 .map_err(|e| CoreError::Backend(e.to_string()))?;
121 index
122 .remove_path(rel)
123 .map_err(|e| CoreError::Backend(e.to_string()))?;
124 index
125 .write()
126 .map_err(|e| CoreError::Backend(e.to_string()))?;
127 let tree_oid = index
128 .write_tree()
129 .map_err(|e| CoreError::Backend(e.to_string()))?;
130 let tree = repo
131 .find_tree(tree_oid)
132 .map_err(|e| CoreError::Backend(e.to_string()))?;
133 let sig = git2::Signature::now("gonzalo", "gonzalo@localhost")
134 .map_err(|e| CoreError::Backend(e.to_string()))?;
135 let parent = repo
136 .head()
137 .ok()
138 .and_then(|h| h.target())
139 .and_then(|oid| repo.find_commit(oid).ok());
140 let parents: Vec<&git2::Commit> = parent.iter().collect();
141 repo.commit(Some("HEAD"), &sig, &sig, message, &tree, &parents)
142 .map_err(|e| CoreError::Backend(e.to_string()))?;
143 Ok(())
144 }
145
146 pub async fn pull(&self, remote: &str, branch: &str) -> Result<PullReport> {
151 let root = self.root.clone();
152 let remote = remote.to_string();
153 let branch = branch.to_string();
154 run_blocking(move || git_pull(&root, &remote, &branch)).await
155 }
156
157 pub async fn push(&self, remote: &str, branch: &str) -> Result<()> {
159 let root = self.root.clone();
160 let remote = remote.to_string();
161 let branch = branch.to_string();
162 run_blocking(move || git_push(&root, &remote, &branch)).await
163 }
164}
165
166fn be<E: std::fmt::Display>(e: E) -> CoreError {
167 CoreError::Backend(e.to_string())
168}
169
170fn lock_repo(root: &Path) -> Result<std::fs::File> {
180 let lock_path = root.join(".gonzalo-git.lock");
181 let lock = std::fs::OpenOptions::new()
182 .create(true)
183 .truncate(false)
184 .write(true)
185 .open(&lock_path)
186 .map_err(be)?;
187 flock(&lock, FlockOperation::LockExclusive).map_err(be)?;
188 Ok(lock)
189}
190
191fn git_pull(root: &Path, remote: &str, branch: &str) -> Result<PullReport> {
192 let repo = git2::Repository::open(root).map_err(be)?;
193 let mut rem = repo.find_remote(remote).map_err(be)?;
194 rem.fetch(&[branch], None, None).map_err(be)?;
195 let fetch_head = repo.find_reference("FETCH_HEAD").map_err(be)?;
196 let fetch_commit = repo
197 .reference_to_annotated_commit(&fetch_head)
198 .map_err(be)?;
199 let (analysis, _) = repo.merge_analysis(&[&fetch_commit]).map_err(be)?;
200
201 if analysis.is_up_to_date() {
202 return Ok(PullReport::default());
203 }
204 if analysis.is_fast_forward() {
205 let refname = format!("refs/heads/{branch}");
206 let mut reference = repo.find_reference(&refname).map_err(be)?;
207 reference
208 .set_target(fetch_commit.id(), "fast-forward")
209 .map_err(be)?;
210 repo.set_head(&refname).map_err(be)?;
211 repo.checkout_head(Some(git2::build::CheckoutBuilder::default().force()))
212 .map_err(be)?;
213 return Ok(PullReport {
214 fast_forwarded: true,
215 ..Default::default()
216 });
217 }
218
219 merge_non_ff(&repo, remote, branch, fetch_commit.id())
220}
221
222fn merge_non_ff(
227 repo: &git2::Repository,
228 remote: &str,
229 branch: &str,
230 remote_oid: git2::Oid,
231) -> Result<PullReport> {
232 let local_oid = repo
233 .head()
234 .map_err(be)?
235 .target()
236 .ok_or_else(|| CoreError::Backend("local HEAD is unborn".into()))?;
237 let local_commit = repo.find_commit(local_oid).map_err(be)?;
238 let remote_commit = repo.find_commit(remote_oid).map_err(be)?;
239 let local_tree = local_commit.tree().map_err(be)?;
240 let remote_tree = remote_commit.tree().map_err(be)?;
241 let base_tree = match repo.merge_base(local_oid, remote_oid) {
244 Ok(base_oid) => Some(repo.find_commit(base_oid).map_err(be)?.tree().map_err(be)?),
245 Err(_) => None,
246 };
247
248 let mut index = repo.index().map_err(be)?;
250 index.read_tree(&local_tree).map_err(be)?;
251
252 let local_changed = changed_paths_set(repo, base_tree.as_ref(), &local_tree)?;
253 let remote_diff = repo
254 .diff_tree_to_tree(base_tree.as_ref(), Some(&remote_tree), None)
255 .map_err(be)?;
256
257 let mut report = PullReport::default();
258 for delta in remote_diff.deltas() {
259 let Some(path) = delta
260 .new_file()
261 .path()
262 .or_else(|| delta.old_file().path())
263 .map(Path::to_path_buf)
264 else {
265 continue;
266 };
267
268 if !local_changed.contains(&path) {
269 if delta.status() == git2::Delta::Deleted {
271 index.remove_path(&path).map_err(be)?;
272 } else if let Some(bytes) = tree_blob(repo, &remote_tree, &path)? {
273 index
274 .add_frombuffer(&blob_entry(&path), &bytes)
275 .map_err(be)?;
276 }
277 continue;
278 }
279
280 let Some(key) = key_from_path(&path) else {
282 continue; };
284 let local_rec = record_at(repo, Some(&local_tree), &path)?;
285 let remote_rec = record_at(repo, Some(&remote_tree), &path)?;
286 match (local_rec, remote_rec) {
287 (Some(local), Some(remote)) if local.body != remote.body => {
288 let base_body = record_at(repo, base_tree.as_ref(), &path)?
289 .map(|r| r.body)
290 .unwrap_or(Body::Inline(Vec::new()));
291 match merge(
292 local.kind.merge_class(),
293 &base_body,
294 &local.body,
295 &remote.body,
296 ) {
297 MergeOutcome::Merged(body) => {
298 let merged = merged_record(&key, &local, &remote, body);
299 let bytes = serde_json::to_vec_pretty(&merged)
300 .map_err(|e| CoreError::Serde(e.to_string()))?;
301 index
302 .add_frombuffer(&blob_entry(&path), &bytes)
303 .map_err(be)?;
304 report.merged.push(key);
305 }
306 MergeOutcome::NeedsResolution => {
307 report.conflicts.push(PullConflict {
309 key,
310 local: Box::new(local),
311 remote: Box::new(remote),
312 });
313 }
314 }
315 }
316 _ => {}
319 }
320 }
321
322 let tree_oid = index.write_tree_to(repo).map_err(be)?;
324 let tree = repo.find_tree(tree_oid).map_err(be)?;
325 let sig = git2::Signature::now("gonzalo", "gonzalo@localhost").map_err(be)?;
326 let refname = format!("refs/heads/{branch}");
327 repo.commit(
328 Some(&refname),
329 &sig,
330 &sig,
331 &format!("merge {remote}/{branch}"),
332 &tree,
333 &[&local_commit, &remote_commit],
334 )
335 .map_err(be)?;
336 repo.set_head(&refname).map_err(be)?;
337 repo.checkout_head(Some(git2::build::CheckoutBuilder::default().force()))
338 .map_err(be)?;
339
340 Ok(report)
341}
342
343fn changed_paths_set(
345 repo: &git2::Repository,
346 base: Option<&git2::Tree>,
347 tree: &git2::Tree,
348) -> Result<HashSet<PathBuf>> {
349 let diff = repo.diff_tree_to_tree(base, Some(tree), None).map_err(be)?;
350 let mut set = HashSet::new();
351 for delta in diff.deltas() {
352 if let Some(path) = delta.new_file().path().or_else(|| delta.old_file().path()) {
353 set.insert(path.to_path_buf());
354 }
355 }
356 Ok(set)
357}
358
359fn record_at(
362 repo: &git2::Repository,
363 tree: Option<&git2::Tree>,
364 path: &Path,
365) -> Result<Option<Record>> {
366 let Some(tree) = tree else {
367 return Ok(None);
368 };
369 match tree.get_path(path) {
370 Ok(entry) => {
371 let obj = entry.to_object(repo).map_err(be)?;
372 let blob = obj
373 .as_blob()
374 .ok_or_else(|| CoreError::Backend("record path is not a blob".into()))?;
375 let rec = serde_json::from_slice(blob.content())
376 .map_err(|e| CoreError::Serde(e.to_string()))?;
377 Ok(Some(rec))
378 }
379 Err(_) => Ok(None),
380 }
381}
382
383fn tree_blob(repo: &git2::Repository, tree: &git2::Tree, path: &Path) -> Result<Option<Vec<u8>>> {
385 match tree.get_path(path) {
386 Ok(entry) => {
387 let obj = entry.to_object(repo).map_err(be)?;
388 Ok(obj.as_blob().map(|b| b.content().to_vec()))
389 }
390 Err(_) => Ok(None),
391 }
392}
393
394fn key_from_path(path: &Path) -> Option<RecordKey> {
396 let comps: Vec<String> = path
397 .components()
398 .map(|c| c.as_os_str().to_string_lossy().into_owned())
399 .collect();
400 if comps.len() != 3 {
401 return None;
402 }
403 let id = comps[2].strip_suffix(".json")?;
404 Some(RecordKey::new(
405 decode_segment(&comps[0]),
406 decode_segment(&comps[1]),
407 decode_segment(id),
408 ))
409}
410
411fn blob_entry(path: &Path) -> git2::IndexEntry {
414 git2::IndexEntry {
415 ctime: git2::IndexTime::new(0, 0),
416 mtime: git2::IndexTime::new(0, 0),
417 dev: 0,
418 ino: 0,
419 mode: 0o100644,
420 uid: 0,
421 gid: 0,
422 file_size: 0,
423 id: git2::Oid::zero(),
424 flags: 0,
425 flags_extended: 0,
426 path: path.to_string_lossy().into_owned().into_bytes(),
427 }
428}
429
430fn merged_record(key: &RecordKey, local: &Record, remote: &Record, body: Body) -> Record {
433 let counter = local.revision.counter.max(remote.revision.counter) + 1;
434 let mut labels = local.meta.labels.clone();
435 labels.extend(remote.meta.labels.clone());
436 let mut links = local.links.clone();
437 for l in &remote.links {
438 if !links.contains(l) {
439 links.push(l.clone());
440 }
441 }
442 let parent = if local.revision.counter >= remote.revision.counter {
443 local.revision.clone()
444 } else {
445 remote.revision.clone()
446 };
447 Record {
448 key: key.clone(),
449 kind: local.kind,
450 revision: Revision {
451 counter,
452 hash: ContentHash::of(body.bytes()),
453 },
454 parent: Some(parent),
455 body,
456 meta: Meta {
457 author: Identity::new("gonzalo-merge"),
458 origin_system: "git-pull".into(),
459 created: local.meta.created.min(remote.meta.created),
460 updated: local.meta.updated.max(remote.meta.updated),
461 labels,
462 },
463 links,
464 }
465}
466
467fn git_push(root: &Path, remote: &str, branch: &str) -> Result<()> {
468 let repo = git2::Repository::open(root).map_err(be)?;
469 let mut rem = repo.find_remote(remote).map_err(be)?;
470 let refspec = format!("refs/heads/{branch}:refs/heads/{branch}");
471
472 let rejected: Rc<RefCell<Vec<String>>> = Rc::new(RefCell::new(Vec::new()));
478 let mut callbacks = git2::RemoteCallbacks::new();
479 {
480 let rejected = Rc::clone(&rejected);
481 callbacks.push_update_reference(move |refname, status| {
482 if let Some(msg) = status {
483 rejected.borrow_mut().push(format!("{refname}: {msg}"));
484 }
485 Ok(())
486 });
487 }
488 let mut opts = git2::PushOptions::new();
489 opts.remote_callbacks(callbacks);
490 rem.push(&[refspec.as_str()], Some(&mut opts)).map_err(be)?;
491 drop(opts); let rejected = rejected.borrow();
494 if !rejected.is_empty() {
495 return Err(CoreError::Backend(format!(
496 "push rejected by remote '{remote}': {}",
497 rejected.join(", ")
498 )));
499 }
500 Ok(())
501}
502
503async fn run_blocking<F, T>(f: F) -> Result<T>
504where
505 F: FnOnce() -> Result<T> + Send + 'static,
506 T: Send + 'static,
507{
508 tokio::task::spawn_blocking(f)
509 .await
510 .map_err(|e| CoreError::Backend(e.to_string()))?
511}
512
513#[async_trait]
514impl gonzalo_core::Store for GitStore {
515 async fn get(&self, key: &RecordKey) -> Result<Option<Record>> {
516 let this = Arc::new(self.root.clone());
517 let key = key.clone();
518 run_blocking(move || {
519 let store = GitStore {
520 root: (*this).clone(),
521 };
522 store.read(&key)
523 })
524 .await
525 }
526
527 async fn put(&self, record: Record, expected: Option<Revision>) -> Result<PutResult> {
528 let root = self.root.clone();
529 run_blocking(move || {
530 let store = GitStore { root: root.clone() };
531 let _lock = lock_repo(&root)?;
534 let current = store.read(&record.key)?;
535 let current_rev = current.as_ref().map(|r| r.revision.clone());
536 if current_rev != expected {
537 if let Some(cur) = current {
538 return Ok(PutResult::Conflict(Box::new(Conflict {
539 key: record.key.clone(),
540 expected,
541 current: cur,
542 })));
543 }
544 return Err(CoreError::NotFound(record.key.clone()));
545 }
546 let (ns, col, file) = record_components(&record.key);
547 let rel = Path::new(&ns).join(&col).join(&file);
548 let abs = root.join(&rel);
549 if let Some(parent) = abs.parent() {
550 std::fs::create_dir_all(parent).map_err(|e| CoreError::Backend(e.to_string()))?;
551 }
552 let bytes =
553 serde_json::to_vec_pretty(&record).map_err(|e| CoreError::Serde(e.to_string()))?;
554 std::fs::write(&abs, &bytes).map_err(|e| CoreError::Backend(e.to_string()))?;
555 store.commit_file(&rel, &format!("put {}", record.key))?;
556 Ok(PutResult::Committed(record.revision))
557 })
558 .await
559 }
560
561 async fn list(&self, prefix: &KeyPrefix) -> Result<Vec<RecordKey>> {
562 let root = self.root.clone();
563 let prefix = prefix.clone();
564 run_blocking(move || {
565 let mut out = Vec::new();
566 collect_keys(&root, &prefix, &mut out)?;
567 Ok(out)
568 })
569 .await
570 }
571
572 async fn delete(&self, key: &RecordKey, expected: Option<Revision>) -> Result<DeleteResult> {
573 let root = self.root.clone();
574 let key = key.clone();
575 run_blocking(move || {
576 let store = GitStore { root: root.clone() };
577 let _lock = lock_repo(&root)?;
581 let current = store.read(&key)?;
582 match (current, &expected) {
583 (None, _) => Ok(DeleteResult::Deleted),
585 (Some(cur), exp) if exp.is_none() || exp.as_ref() == Some(&cur.revision) => {
587 let (ns, col, file) = record_components(&key);
588 let rel = Path::new(&ns).join(&col).join(&file);
589 let abs = root.join(&rel);
590 match std::fs::remove_file(&abs) {
591 Ok(()) => {}
592 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
593 Err(e) => return Err(CoreError::Backend(e.to_string())),
594 }
595 store.commit_removal(&rel, &format!("delete {key}"))?;
596 Ok(DeleteResult::Deleted)
597 }
598 (Some(cur), _) => Ok(DeleteResult::Conflict(Box::new(Conflict {
600 key: key.clone(),
601 expected,
602 current: cur,
603 }))),
604 }
605 })
606 .await
607 }
608}
609
610fn collect_keys(root: &Path, prefix: &KeyPrefix, out: &mut Vec<RecordKey>) -> Result<()> {
611 let namespaces = match std::fs::read_dir(root) {
612 Ok(rd) => rd,
613 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
614 Err(e) => return Err(CoreError::Backend(e.to_string())),
615 };
616 for ns in namespaces {
617 let ns = ns.map_err(|e| CoreError::Backend(e.to_string()))?;
618 let ns_name = ns.file_name().to_string_lossy().to_string();
619 if ns_name == ".git" || !ns.path().is_dir() {
620 continue;
621 }
622 for col in std::fs::read_dir(ns.path()).map_err(|e| CoreError::Backend(e.to_string()))? {
623 let col = col.map_err(|e| CoreError::Backend(e.to_string()))?;
624 if !col.path().is_dir() {
625 continue;
626 }
627 let col_name = col.file_name().to_string_lossy().to_string();
628 for f in std::fs::read_dir(col.path()).map_err(|e| CoreError::Backend(e.to_string()))? {
629 let f = f.map_err(|e| CoreError::Backend(e.to_string()))?;
630 let fname = f.file_name().to_string_lossy().to_string();
631 if let Some(id) = fname.strip_suffix(".json") {
632 let key = RecordKey::new(
635 decode_segment(&ns_name),
636 decode_segment(&col_name),
637 decode_segment(id),
638 );
639 if prefix.matches(&key) {
640 out.push(key);
641 }
642 }
643 }
644 }
645 }
646 Ok(())
647}