1use crate::{
9 BlobStore, Body, Identity, KeyPrefix, MergeOutcome, Meta, PutResult, Record, RecordKey, Result,
10 Revision, Store, merge,
11};
12use std::collections::BTreeSet;
13
14#[derive(Clone, Debug, PartialEq, Eq)]
16pub struct SyncConflict {
17 pub key: RecordKey,
18 pub a: Box<Record>,
19 pub b: Box<Record>,
20}
21
22#[derive(Clone, Debug, Default, PartialEq, Eq)]
24#[must_use = "a SyncReport may contain unresolved conflicts that must be handled"]
25pub struct SyncReport {
26 pub copied_to_a: Vec<RecordKey>,
28 pub copied_to_b: Vec<RecordKey>,
30 pub merged: Vec<RecordKey>,
32 pub conflicts: Vec<SyncConflict>,
34}
35
36const MAX_SYNC_PASSES: usize = 16;
43
44pub async fn sync(a: &dyn Store, b: &dyn Store) -> Result<SyncReport> {
52 sync_with_ancestry(a, b, None).await
53}
54
55pub async fn sync_with_ancestry(
62 a: &dyn Store,
63 b: &dyn Store,
64 ancestry: Option<&dyn BlobStore>,
65) -> Result<SyncReport> {
66 let mut report = SyncReport::default();
67 for _ in 0..MAX_SYNC_PASSES {
68 let (pass, raced) = sync_pass(a, b, ancestry).await?;
69 report = pass;
70 if !raced {
71 break; }
73 }
74 Ok(report)
75}
76
77async fn ancestry_base(ancestry: Option<&dyn BlobStore>, rec_a: &Record, rec_b: &Record) -> Body {
81 if let Some(anc) = ancestry
82 && let (Some(pa), Some(pb)) = (&rec_a.parent, &rec_b.parent)
83 && pa == pb
84 && let Ok(Some(bytes)) = anc.get_blob(&pa.hash).await
85 {
86 return Body::Inline(bytes);
87 }
88 Body::Inline(Vec::new())
89}
90
91async fn sync_pass(
97 a: &dyn Store,
98 b: &dyn Store,
99 ancestry: Option<&dyn BlobStore>,
100) -> Result<(SyncReport, bool)> {
101 let mut report = SyncReport::default();
102 let mut raced = false;
103
104 let mut keys: BTreeSet<RecordKey> = BTreeSet::new();
105 keys.extend(a.list(&KeyPrefix::default()).await?);
106 keys.extend(b.list(&KeyPrefix::default()).await?);
107
108 for key in keys {
109 let ra = a.get(&key).await?;
110 let rb = b.get(&key).await?;
111 match (ra, rb) {
112 (Some(rec), None) => {
113 if copy(b, &rec).await? {
114 report.copied_to_b.push(key);
115 } else {
116 raced = true;
117 }
118 }
119 (None, Some(rec)) => {
120 if copy(a, &rec).await? {
121 report.copied_to_a.push(key);
122 } else {
123 raced = true;
124 }
125 }
126 (Some(rec_a), Some(rec_b)) => {
127 if rec_a.revision == rec_b.revision {
128 continue; }
130 let base = ancestry_base(ancestry, &rec_a, &rec_b).await;
131 match merge(rec_a.kind.merge_class(), &base, &rec_a.body, &rec_b.body) {
132 MergeOutcome::Merged(body) => {
133 let merged = build_merged(&key, &rec_a, &rec_b, body);
134 let la = overwrite(a, &merged, &rec_a.revision).await?;
135 let lb = overwrite(b, &merged, &rec_b.revision).await?;
136 if la && lb {
137 report.merged.push(key);
138 } else {
139 raced = true;
142 }
143 }
144 MergeOutcome::NeedsResolution => {
145 report.conflicts.push(SyncConflict {
146 key,
147 a: Box::new(rec_a),
148 b: Box::new(rec_b),
149 });
150 }
151 }
152 }
153 (None, None) => {}
154 }
155 }
156 Ok((report, raced))
157}
158
159async fn copy(dst: &dyn Store, rec: &Record) -> Result<bool> {
162 Ok(matches!(
163 dst.put(rec.clone(), None).await?,
164 PutResult::Committed(_)
165 ))
166}
167
168async fn overwrite(dst: &dyn Store, rec: &Record, expected: &Revision) -> Result<bool> {
172 Ok(matches!(
173 dst.put(rec.clone(), Some(expected.clone())).await?,
174 PutResult::Committed(_)
175 ))
176}
177
178fn build_merged(key: &RecordKey, a: &Record, b: &Record, body: Body) -> Record {
179 let counter = a.revision.counter.max(b.revision.counter) + 1;
180 let mut labels = a.meta.labels.clone();
181 labels.extend(b.meta.labels.clone());
182 let mut links = a.links.clone();
183 for l in &b.links {
184 if !links.contains(l) {
185 links.push(l.clone());
186 }
187 }
188 Record {
189 key: key.clone(),
190 kind: a.kind,
191 revision: Revision {
192 counter,
193 hash: crate::ContentHash::of(body.bytes()),
194 },
195 parent: Some(if a.revision.counter >= b.revision.counter {
196 a.revision.clone()
197 } else {
198 b.revision.clone()
199 }),
200 body,
201 meta: Meta {
202 author: Identity::new("gonzalo-sync"),
203 origin_system: "sync".into(),
204 created: a.meta.created.min(b.meta.created),
205 updated: a.meta.updated.max(b.meta.updated),
206 labels,
207 },
208 links,
209 }
210}
211
212#[cfg(test)]
213mod tests {
214 use super::*;
215 use crate::{DeleteResult, PutResult, RecordKind, store::Conflict};
216 use async_trait::async_trait;
217 use std::collections::BTreeMap;
218 use std::sync::Mutex;
219
220 #[derive(Default)]
221 struct MemStore(Mutex<BTreeMap<RecordKey, Record>>);
222
223 #[async_trait]
224 impl Store for MemStore {
225 async fn get(&self, key: &RecordKey) -> Result<Option<Record>> {
226 Ok(self.0.lock().unwrap().get(key).cloned())
227 }
228 async fn put(&self, record: Record, expected: Option<Revision>) -> Result<PutResult> {
229 let mut g = self.0.lock().unwrap();
230 let current = g.get(&record.key).map(|r| r.revision.clone());
231 if current != expected {
232 if let Some(cur) = g.get(&record.key).cloned() {
233 return Ok(PutResult::Conflict(Box::new(Conflict {
234 key: record.key.clone(),
235 expected,
236 current: cur,
237 })));
238 }
239 return Err(crate::CoreError::NotFound(record.key.clone()));
240 }
241 let rev = record.revision.clone();
242 g.insert(record.key.clone(), record);
243 Ok(PutResult::Committed(rev))
244 }
245 async fn list(&self, prefix: &KeyPrefix) -> Result<Vec<RecordKey>> {
246 Ok(self
247 .0
248 .lock()
249 .unwrap()
250 .keys()
251 .filter(|k| prefix.matches(k))
252 .cloned()
253 .collect())
254 }
255 async fn delete(
256 &self,
257 key: &RecordKey,
258 expected: Option<Revision>,
259 ) -> Result<DeleteResult> {
260 let mut g = self.0.lock().unwrap();
261 match g.get(key) {
262 None => Ok(DeleteResult::Deleted),
263 Some(cur) if expected.is_none() || expected.as_ref() == Some(&cur.revision) => {
264 g.remove(key);
265 Ok(DeleteResult::Deleted)
266 }
267 Some(cur) => Ok(DeleteResult::Conflict(Box::new(Conflict {
268 key: key.clone(),
269 expected,
270 current: cur.clone(),
271 }))),
272 }
273 }
274 }
275
276 #[derive(Default)]
281 struct FlakyOnceStore {
282 inner: Mutex<BTreeMap<RecordKey, Record>>,
283 tripped: Mutex<std::collections::HashSet<RecordKey>>,
284 }
285
286 #[async_trait]
287 impl Store for FlakyOnceStore {
288 async fn get(&self, key: &RecordKey) -> Result<Option<Record>> {
289 Ok(self.inner.lock().unwrap().get(key).cloned())
290 }
291 async fn put(&self, record: Record, expected: Option<Revision>) -> Result<PutResult> {
292 if expected.is_some() && self.tripped.lock().unwrap().insert(record.key.clone()) {
293 let current = self
294 .inner
295 .lock()
296 .unwrap()
297 .get(&record.key)
298 .cloned()
299 .unwrap();
300 return Ok(PutResult::Conflict(Box::new(Conflict {
301 key: record.key.clone(),
302 expected,
303 current,
304 })));
305 }
306 let mut g = self.inner.lock().unwrap();
307 let current = g.get(&record.key).map(|r| r.revision.clone());
308 if current != expected {
309 if let Some(cur) = g.get(&record.key).cloned() {
310 return Ok(PutResult::Conflict(Box::new(Conflict {
311 key: record.key.clone(),
312 expected,
313 current: cur,
314 })));
315 }
316 return Err(crate::CoreError::NotFound(record.key.clone()));
317 }
318 let rev = record.revision.clone();
319 g.insert(record.key.clone(), record);
320 Ok(PutResult::Committed(rev))
321 }
322 async fn list(&self, prefix: &KeyPrefix) -> Result<Vec<RecordKey>> {
323 Ok(self
324 .inner
325 .lock()
326 .unwrap()
327 .keys()
328 .filter(|k| prefix.matches(k))
329 .cloned()
330 .collect())
331 }
332 async fn delete(
333 &self,
334 key: &RecordKey,
335 expected: Option<Revision>,
336 ) -> Result<DeleteResult> {
337 let mut g = self.inner.lock().unwrap();
338 match g.get(key) {
339 None => Ok(DeleteResult::Deleted),
340 Some(cur) if expected.is_none() || expected.as_ref() == Some(&cur.revision) => {
341 g.remove(key);
342 Ok(DeleteResult::Deleted)
343 }
344 Some(cur) => Ok(DeleteResult::Conflict(Box::new(Conflict {
345 key: key.clone(),
346 expected,
347 current: cur.clone(),
348 }))),
349 }
350 }
351 }
352
353 #[derive(Default)]
357 struct AlwaysRacyStore(Mutex<BTreeMap<RecordKey, Record>>);
358
359 #[async_trait]
360 impl Store for AlwaysRacyStore {
361 async fn get(&self, key: &RecordKey) -> Result<Option<Record>> {
362 Ok(self.0.lock().unwrap().get(key).cloned())
363 }
364 async fn put(&self, record: Record, expected: Option<Revision>) -> Result<PutResult> {
365 let mut g = self.0.lock().unwrap();
366 match expected {
367 None if !g.contains_key(&record.key) => {
368 let rev = record.revision.clone();
369 g.insert(record.key.clone(), record);
370 Ok(PutResult::Committed(rev))
371 }
372 _ => {
373 let current = g.get(&record.key).cloned().unwrap();
374 Ok(PutResult::Conflict(Box::new(Conflict {
375 key: record.key.clone(),
376 expected,
377 current,
378 })))
379 }
380 }
381 }
382 async fn list(&self, prefix: &KeyPrefix) -> Result<Vec<RecordKey>> {
383 Ok(self
384 .0
385 .lock()
386 .unwrap()
387 .keys()
388 .filter(|k| prefix.matches(k))
389 .cloned()
390 .collect())
391 }
392 async fn delete(
393 &self,
394 key: &RecordKey,
395 expected: Option<Revision>,
396 ) -> Result<DeleteResult> {
397 let mut g = self.0.lock().unwrap();
398 match g.get(key) {
399 None => Ok(DeleteResult::Deleted),
400 Some(cur) if expected.is_none() || expected.as_ref() == Some(&cur.revision) => {
401 g.remove(key);
402 Ok(DeleteResult::Deleted)
403 }
404 Some(cur) => Ok(DeleteResult::Conflict(Box::new(Conflict {
405 key: key.clone(),
406 expected,
407 current: cur.clone(),
408 }))),
409 }
410 }
411 }
412
413 fn rec(id: &str, kind: RecordKind, payload: &str) -> Record {
414 let body = Body::Inline(payload.as_bytes().to_vec());
415 Record {
416 revision: Revision::initial(body.bytes()),
417 parent: None,
418 body,
419 kind,
420 key: RecordKey::new("ns", "col", id),
421 meta: Meta {
422 author: Identity::new("t"),
423 origin_system: "test".into(),
424 created: 0,
425 updated: 0,
426 labels: BTreeMap::new(),
427 },
428 links: Vec::new(),
429 }
430 }
431
432 #[tokio::test]
433 async fn copies_one_sided_records_both_directions() {
434 let a = MemStore::default();
435 let b = MemStore::default();
436 let _ = a
437 .put(rec("only_a", RecordKind::Topic, "x"), None)
438 .await
439 .unwrap();
440 let _ = b
441 .put(rec("only_b", RecordKind::Topic, "y"), None)
442 .await
443 .unwrap();
444
445 let report = sync(&a, &b).await.unwrap();
446 assert_eq!(
447 report.copied_to_b,
448 vec![RecordKey::new("ns", "col", "only_a")]
449 );
450 assert_eq!(
451 report.copied_to_a,
452 vec![RecordKey::new("ns", "col", "only_b")]
453 );
454 assert!(
455 a.get(&RecordKey::new("ns", "col", "only_b"))
456 .await
457 .unwrap()
458 .is_some()
459 );
460 assert!(
461 b.get(&RecordKey::new("ns", "col", "only_a"))
462 .await
463 .unwrap()
464 .is_some()
465 );
466 }
467
468 #[tokio::test]
469 async fn append_only_divergence_auto_merges() {
470 let a = MemStore::default();
471 let b = MemStore::default();
472 let _ = a
473 .put(rec("t", RecordKind::Topic, "base\nfrom_a\n"), None)
474 .await
475 .unwrap();
476 let _ = b
477 .put(rec("t", RecordKind::Topic, "base\nfrom_b\n"), None)
478 .await
479 .unwrap();
480
481 let report = sync(&a, &b).await.unwrap();
482 assert_eq!(report.merged, vec![RecordKey::new("ns", "col", "t")]);
483 assert!(report.conflicts.is_empty());
484 let merged = a
485 .get(&RecordKey::new("ns", "col", "t"))
486 .await
487 .unwrap()
488 .unwrap();
489 let text = String::from_utf8(merged.body.bytes().to_vec()).unwrap();
490 assert!(text.contains("from_a") && text.contains("from_b") && text.contains("base"));
491 let mb = b
493 .get(&RecordKey::new("ns", "col", "t"))
494 .await
495 .unwrap()
496 .unwrap();
497 assert_eq!(merged.revision, mb.revision);
498 }
499
500 #[tokio::test]
501 async fn checkpoint_divergence_surfaces_conflict() {
502 let a = MemStore::default();
503 let b = MemStore::default();
504 let _ = a
505 .put(rec("c", RecordKind::Checkpoint, "a"), None)
506 .await
507 .unwrap();
508 let _ = b
509 .put(rec("c", RecordKind::Checkpoint, "b"), None)
510 .await
511 .unwrap();
512
513 let report = sync(&a, &b).await.unwrap();
514 assert_eq!(report.conflicts.len(), 1);
515 assert_eq!(report.conflicts[0].key, RecordKey::new("ns", "col", "c"));
516 assert!(report.merged.is_empty());
517 }
518
519 #[tokio::test]
520 async fn memory_tier_divergence_surfaces_conflict() {
521 let a = MemStore::default();
522 let b = MemStore::default();
523 let _ = a
524 .put(rec("m", RecordKind::MemoryTier, "a"), None)
525 .await
526 .unwrap();
527 let _ = b
528 .put(rec("m", RecordKind::MemoryTier, "b"), None)
529 .await
530 .unwrap();
531
532 let report = sync(&a, &b).await.unwrap();
533 assert_eq!(report.conflicts.len(), 1);
534 assert!(report.merged.is_empty());
535 }
536
537 #[tokio::test]
538 async fn session_divergence_auto_merges() {
539 let a = MemStore::default();
540 let b = MemStore::default();
541 let _ = a
542 .put(rec("s", RecordKind::Session, "base\nfrom_a\n"), None)
543 .await
544 .unwrap();
545 let _ = b
546 .put(rec("s", RecordKind::Session, "base\nfrom_b\n"), None)
547 .await
548 .unwrap();
549
550 let report = sync(&a, &b).await.unwrap();
551 assert_eq!(report.merged, vec![RecordKey::new("ns", "col", "s")]);
552 assert!(report.conflicts.is_empty());
553 }
554
555 #[tokio::test]
556 async fn session_with_blank_and_duplicate_lines_survives_divergent_sync() {
557 let a = MemStore::default();
564 let b = MemStore::default();
565 let _ = a
566 .put(
567 rec("s", RecordKind::Session, "a\n\nyes\nyes\nfrom_a\n"),
568 None,
569 )
570 .await
571 .unwrap();
572 let _ = b
573 .put(
574 rec("s", RecordKind::Session, "a\n\nyes\nyes\nfrom_b\n"),
575 None,
576 )
577 .await
578 .unwrap();
579
580 let report = sync(&a, &b).await.unwrap();
581 let key = RecordKey::new("ns", "col", "s");
582 assert_eq!(report.merged, vec![key.clone()]);
583 assert!(report.conflicts.is_empty());
584
585 let ma = a.get(&key).await.unwrap().unwrap();
586 let text = String::from_utf8(ma.body.bytes().to_vec()).unwrap();
587 assert_eq!(text, "a\n\nyes\nyes\nfrom_a\nfrom_b\n");
589 let mb = b.get(&key).await.unwrap().unwrap();
591 assert_eq!(ma.revision, mb.revision);
592 assert_eq!(ma.body, mb.body);
593 }
594
595 #[tokio::test]
596 async fn re_loops_until_a_racing_store_converges() {
597 let a = MemStore::default();
601 let b = FlakyOnceStore::default();
602 let _ = a
603 .put(rec("t", RecordKind::Topic, "base\nfrom_a\n"), None)
604 .await
605 .unwrap();
606 let _ = b
607 .put(rec("t", RecordKind::Topic, "base\nfrom_b\n"), None)
608 .await
609 .unwrap();
610
611 let report = sync(&a, &b).await.unwrap();
612
613 assert!(report.conflicts.is_empty());
614 let key = RecordKey::new("ns", "col", "t");
615 let ra = a.get(&key).await.unwrap().unwrap();
616 let rb = b.get(&key).await.unwrap().unwrap();
617 assert_eq!(ra.revision, rb.revision);
619 let text = String::from_utf8(rb.body.bytes().to_vec()).unwrap();
620 assert!(text.contains("from_a") && text.contains("from_b"));
621 assert_eq!(report.merged, vec![key]);
622 }
623
624 fn structured_divergence() -> (Record, Record, &'static str) {
627 let base = rec("m", RecordKind::MemoryTier, r#"{"name":"a","content":"x"}"#);
628 let base_rev = base.revision.clone();
629 let mut ours = rec("m", RecordKind::MemoryTier, r#"{"name":"b","content":"x"}"#);
630 ours.parent = Some(base_rev.clone());
631 let mut theirs = rec("m", RecordKind::MemoryTier, r#"{"name":"a","content":"y"}"#);
632 theirs.parent = Some(base_rev);
633 (ours, theirs, r#"{"name":"a","content":"x"}"#)
634 }
635
636 #[tokio::test]
637 async fn structured_divergence_merges_with_ancestry() {
638 use crate::ancestry::tests::Mem;
639 let a = Mem::default();
640 let b = Mem::default();
641 let ancestry = Mem::default();
642 let (ours, theirs, base_body) = structured_divergence();
643 ancestry.put_blob(base_body.as_bytes()).await.unwrap();
645 let _ = a.put(ours, None).await.unwrap();
646 let _ = b.put(theirs, None).await.unwrap();
647
648 let report = sync_with_ancestry(&a, &b, Some(&ancestry)).await.unwrap();
649
650 let key = RecordKey::new("ns", "col", "m");
651 assert_eq!(report.merged, vec![key.clone()], "3-way merged");
652 assert!(report.conflicts.is_empty());
653 let merged = a.get(&key).await.unwrap().unwrap();
655 let v: serde_json::Value = serde_json::from_slice(merged.body.bytes()).unwrap();
656 assert_eq!(v, serde_json::json!({"name": "b", "content": "y"}));
657 }
658
659 #[tokio::test]
660 async fn structured_divergence_conflicts_without_ancestry() {
661 use crate::ancestry::tests::Mem;
662 let a = Mem::default();
663 let b = Mem::default();
664 let (ours, theirs, _) = structured_divergence();
665 let _ = a.put(ours, None).await.unwrap();
666 let _ = b.put(theirs, None).await.unwrap();
667
668 let report = sync(&a, &b).await.unwrap();
671 assert_eq!(report.conflicts.len(), 1);
672 assert!(report.merged.is_empty());
673 }
674
675 #[tokio::test]
676 async fn bounded_retry_terminates_under_continuous_writes() {
677 let a = AlwaysRacyStore::default();
681 let b = AlwaysRacyStore::default();
682 let _ = a
683 .put(rec("t", RecordKind::Topic, "base\nfrom_a\n"), None)
684 .await
685 .unwrap();
686 let _ = b
687 .put(rec("t", RecordKind::Topic, "base\nfrom_b\n"), None)
688 .await
689 .unwrap();
690
691 let report = sync(&a, &b).await.unwrap();
692
693 let key = RecordKey::new("ns", "col", "t");
695 let ra = a.get(&key).await.unwrap().unwrap();
696 let rb = b.get(&key).await.unwrap().unwrap();
697 assert_ne!(ra.revision, rb.revision);
698 assert!(report.merged.is_empty());
699 }
700}