Skip to main content

gonzalo_core/
sync.rs

1//! Reconcile two `Store`s. Any store can be a sync peer. Append-only kinds
2//! auto-merge by union; structured/opaque divergences are surfaced as
3//! conflicts. [`sync_with_ancestry`] 3-way-merges structured bodies against
4//! their real common ancestor when an [`AncestryStore`](crate::AncestryStore)
5//! retains it; [`sync`] uses an empty base (correct for append-only union, safe
6//! otherwise) — see ADR 0016.
7
8use crate::{
9    BlobStore, Body, Identity, KeyPrefix, MergeOutcome, Meta, PutResult, Record, RecordKey, Result,
10    Revision, Store, merge,
11};
12use std::collections::BTreeSet;
13
14/// A divergence that could not be auto-merged and needs caller/CLI resolution.
15#[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/// What a sync run did.
23#[derive(Clone, Debug, Default, PartialEq, Eq)]
24#[must_use = "a SyncReport may contain unresolved conflicts that must be handled"]
25pub struct SyncReport {
26    /// Keys copied into store A (were only in B).
27    pub copied_to_a: Vec<RecordKey>,
28    /// Keys copied into store B (were only in A).
29    pub copied_to_b: Vec<RecordKey>,
30    /// Keys auto-merged (append-only) and written to both stores.
31    pub merged: Vec<RecordKey>,
32    /// Divergences needing manual resolution.
33    pub conflicts: Vec<SyncConflict>,
34}
35
36/// Upper bound on sync passes before giving up on a non-quiescent pair.
37///
38/// Each pass re-reads both stores, so a store that settles converges within
39/// one extra pass. The cap only bites when writers never stop racing the merge
40/// window (livelock guard): rather than spin forever, sync returns the last
41/// pass's best-effort report.
42const MAX_SYNC_PASSES: usize = 16;
43
44/// Reconcile stores `a` and `b`. After a clean run (no `conflicts`), both
45/// stores hold the same set of records for every key.
46///
47/// Stores need not be quiescent. A single pass can lose a write that lands in
48/// the read→merge→write window (the OCC `put` returns `Conflict`); sync re-runs
49/// the pass until one completes without any such race (a fixpoint), bounded by
50/// [`MAX_SYNC_PASSES`] so continuous concurrent writes can't livelock it.
51pub async fn sync(a: &dyn Store, b: &dyn Store) -> Result<SyncReport> {
52    sync_with_ancestry(a, b, None).await
53}
54
55/// As [`sync`], but 3-way-merges divergent `Structured` bodies against their
56/// real common ancestor when one is available. `ancestry` is a content-addressed
57/// store of past bodies keyed by revision hash (see
58/// [`AncestryStore`](crate::AncestryStore)): when two records diverge from a
59/// shared parent revision whose body it holds, that body is the merge base;
60/// otherwise sync falls back to the empty base (ADR 0016).
61pub 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; // quiescent: this pass landed cleanly, stores have converged.
72        }
73    }
74    Ok(report)
75}
76
77/// The merge base for a divergence: the body of `a`/`b`'s shared parent revision
78/// when `ancestry` retains it, else an empty base (the base-agnostic fallback,
79/// correct for `AppendOnly` and safe for the rest).
80async 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
91/// One reconciliation pass over the union of keys. Returns the pass's report
92/// and whether any write lost an OCC race (`true` ⇒ a store changed mid-pass,
93/// so the caller should re-loop). A `NeedsResolution` merge conflict is a
94/// terminal divergence (surfaced in the report), not a race, and does not
95/// trigger a re-loop.
96async 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; // already in sync
129                }
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                            // At least one side raced; re-loop to reconcile the
140                            // store that moved against the now-merged peer.
141                            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
159/// Create `rec` in `dst`. Returns `false` if `dst` changed concurrently
160/// (the key already exists), signalling the caller to re-loop.
161async 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
168/// Conditionally overwrite `dst` with `rec`, expecting revision `expected`.
169/// Returns `false` if a concurrent mutation raced the merge window (`put`
170/// returned `Conflict`), signalling the caller to re-loop.
171async 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    /// A store that returns one spurious `Conflict` on the first conditional
277    /// (`expected.is_some()`) `put` per key — a concurrent writer that races
278    /// the first overwrite — then behaves normally. Forces the sync re-loop to
279    /// retry and still converge.
280    #[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    /// A store whose conditional `put` *always* races (a concurrent writer that
354    /// never stops). Initial creates (`expected == None`) commit; every
355    /// overwrite conflicts. Used to prove the re-loop is bounded and terminates.
356    #[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        // Both stores converge to the same revision.
492        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        // Regression for #133: a Session whose committed body legitimately holds
558        // a blank line and a repeated line must survive a divergent sync intact.
559        // `sync` merges against an empty base, so the old split+dedup path
560        // silently dropped the blank line and collapsed the repeat, corrupting
561        // BOTH stores (ADR 0005 violation). The merge must be side A verbatim
562        // plus side B's divergent tail.
563        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        // Blank line preserved, "yes\nyes" not collapsed, both appends present.
588        assert_eq!(text, "a\n\nyes\nyes\nfrom_a\nfrom_b\n");
589        // Both stores converge to the same merged revision.
590        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        // B races the first overwrite (non-quiescent during the merge window).
598        // A single pass would swallow that conflict and leave B un-synced; the
599        // re-loop must retry until both stores converge.
600        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        // Both stores converged despite B racing the first overwrite.
618        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    /// Build ours/theirs Structured records diverging from a shared base
625    /// revision (disjoint field edits), plus the base body to retain.
626    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        // Retain the shared base body under its revision hash.
644        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        // Disjoint field edits both applied against the real base.
654        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        // No ancestry → empty base → the Structured merge cannot tell a one-sided
669        // edit from a real conflict, so it surfaces a conflict.
670        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        // Both stores race *every* overwrite — a non-quiescent pair that never
678        // settles. The re-loop must be bounded: sync returns (does not hang)
679        // and reports the divergence as unresolved rather than spinning.
680        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        // No overwrite ever committed, so nothing converged.
694        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}