Development Documentation (main branch) - For stable release docs, see docs.rs/eidetica
Skip to main content

eidetica/
testing.rs

1//! In-process multi-instance test harness (Tier 0).
2//!
3//! Standing up several real `Instance`s that sync with one another takes a pile
4//! of identical boilerplate: create a user and key, enable sync, register a
5//! transport, start serving, resolve the bound address. [`Cluster`] does *that*
6//! plumbing and nothing else — it hands back wired peers and leaves every policy
7//! decision (auth, what to write, how to drive sync) to the test.
8//!
9//! That split is deliberate. A correctness harness for authenticated CRDT sync
10//! must keep the two things it exists to exercise — the **transport** and the
11//! **auth** — under the test's control, not baked into the setup:
12//!
13//! - **Transport is a seam.** [`ClusterBuilder::transport`] takes any
14//!   [`TestTransport`]; the default is [`HttpLoopback`]. A controllable in-memory
15//!   transport (deliver / reorder / drop / single-step — Tier 1) is a drop-in
16//!   here, which is the point: Tier 1 extends this, it doesn't replace it.
17//! - **Auth is the test's.** The harness never grants keys or permissions for
18//!   you. A peer exposes its `User`, key id, and key name; the test creates its
19//!   database with whatever auth posture it's exercising. [`add_auth_keys`] and
20//!   [`set_global_auth_key`] are policy-neutral *tools* the test composes — they
21//!   apply the keys you pass, they don't choose them.
22//!
23//! What the harness owns is plumbing only: wiring peers, marking a tree
24//! sync-enabled ([`Peer::serve`]), driving an exchange ([`Cluster::exchange`]),
25//! and observing convergence ([`Cluster::converged`]). It does not hold your
26//! databases — the test opens and keeps those itself.
27//!
28//! This is **topology A** (multi-peer sync): N independent `Instance`s, each
29//! owning an in-memory backend. Sync is driven explicitly (no background timers),
30//! so a test fully orders the exchange. The multi-client / single-service
31//! topology is a separate harness that lands with the `service` feature.
32//!
33//! Gated behind `cfg(any(test, feature = "testing"))` alongside [`FixedClock`]
34//! and [`Instance::create_backend_with_clock`]; never compiled into a release build.
35//!
36//! ```no_run
37//! # async fn ex() -> eidetica::Result<()> {
38//! use eidetica::{
39//!     auth::{Permission, types::AuthKey},
40//!     crdt::Doc,
41//!     testing::{Cluster, set_global_auth_key},
42//!     user::types::SyncSettings,
43//! };
44//!
45//! let mut net = Cluster::builder().peers(2).build().await?;
46//!
47//! // Peer 0 creates a database with auth the *test* chooses, then serves it.
48//! let key0 = net.peer(0).key_id().clone();
49//! let mut settings = Doc::new();
50//! settings.set("name", "chat");
51//! let db = net.peer_mut(0).user_mut().create_database(settings, &key0).await?;
52//! let room = db.root_id().clone();
53//! set_global_auth_key(&db, AuthKey::active(None, Permission::Write(10))).await?;
54//! net.peer_mut(0).serve(&room).await?;
55//!
56//! // Peer 1 bootstraps with its own key, then converges against peer 0.
57//! let key1 = net.peer(1).key_id().clone();
58//! let signing_key1 = net.peer(1).user().get_signing_key(&key1)?;
59//! let name1 = net.peer(1).key_name().to_string();
60//! let addr0 = net.peer(0).address().clone();
61//! net.peer(1)
62//!     .sync()
63//!     .sync_with_peer_for_bootstrap_with_key(
64//!         &addr0,
65//!         &room,
66//!         &signing_key1,
67//!         &name1,
68//!         Permission::Write(10),
69//!     )
70//!     .await?;
71//! net.peer_mut(1)
72//!     .user_mut()
73//!     .track_database(room.clone(), &key1, SyncSettings::disabled())
74//!     .await?;
75//!
76//! net.exchange(1, 0, &room).await?;
77//! assert!(net.converged(&[0, 1], &room).await?);
78//! # Ok(()) }
79//! ```
80
81use std::sync::{
82    Arc,
83    atomic::{AtomicBool, Ordering},
84};
85
86use async_trait::async_trait;
87
88use crate::{
89    Database, Entry, Instance, NewUser, Result, Snapshot,
90    auth::{Permission, crypto::PublicKey, types::AuthKey},
91    backend::{BackendImpl, VerificationStatus, database::InMemory},
92    clock::{Clock, FixedClock},
93    crdt::Doc,
94    entry::ID,
95    sync::{
96        Address, Sync,
97        error::SyncError,
98        handler::SyncHandler,
99        protocol::{RequestContext, SyncRequest, SyncResponse},
100        transports::{SyncTransport, TransportBuilder, http::HttpTransport},
101    },
102    user::{User, types::SyncSettings},
103};
104
105/// Display name given to every peer's signing key. Exposed per peer via
106/// [`Peer::key_name`] so a test can name it in a bootstrap request.
107const KEY_NAME: &str = "test-key";
108
109/// How a peer makes itself reachable to other peers. The one seam Tier 1 swaps:
110/// implement this over an in-memory, controllable network and the rest of the
111/// harness is unchanged.
112#[async_trait]
113pub trait TestTransport: Send + std::marker::Sync {
114    /// Register this transport on `sync`, start serving, and return the address
115    /// other peers use to reach it. Called once per peer at build time.
116    async fn serve(&self, sync: &Sync) -> Result<Address>;
117}
118
119/// Default [`TestTransport`]: HTTP over an OS-assigned loopback port.
120#[derive(Debug, Default, Clone)]
121pub struct HttpLoopback;
122
123#[async_trait]
124impl TestTransport for HttpLoopback {
125    async fn serve(&self, sync: &Sync) -> Result<Address> {
126        sync.register_transport("http", HttpTransport::builder().bind("127.0.0.1:0"))
127            .await?;
128        sync.accept_connections().await?;
129        Ok(Address::http(sync.get_server_address().await?))
130    }
131}
132
133/// Builder for a [`Cluster`]. Obtain via [`Cluster::builder`].
134pub struct ClusterBuilder {
135    peers: usize,
136    clock: Option<Arc<dyn Clock>>,
137    transport: Arc<dyn TestTransport>,
138}
139
140impl ClusterBuilder {
141    /// Number of peers (independent `Instance`s) to start. Defaults to 2.
142    pub fn peers(mut self, n: usize) -> Self {
143        self.peers = n;
144        self
145    }
146
147    /// Share a single clock across every peer (e.g. a [`FixedClock`] the test
148    /// drives by hand). When unset, each peer gets its own fresh `FixedClock`,
149    /// mirroring the standard `test_instance()` setup.
150    pub fn clock(mut self, clock: Arc<dyn Clock>) -> Self {
151        self.clock = Some(clock);
152        self
153    }
154
155    /// Use a custom [`TestTransport`] for every peer. Defaults to [`HttpLoopback`].
156    pub fn transport(mut self, transport: Arc<dyn TestTransport>) -> Self {
157        self.transport = transport;
158        self
159    }
160
161    /// Build the cluster: for each peer, open an instance, create a user and
162    /// signing key, enable sync, and serve over the configured transport.
163    pub async fn build(self) -> Result<Cluster> {
164        let mut peers = Vec::with_capacity(self.peers);
165        for i in 0..self.peers {
166            let clock: Arc<dyn Clock> = match &self.clock {
167                Some(shared) => shared.clone(),
168                None => Arc::new(FixedClock::default()),
169            };
170            // Create the backend and bootstrap its admin user in one step; each
171            // peer is a fresh in-memory instance, so its only user is this one.
172            let (instance, mut user) = Instance::create_backend_with_clock(
173                Box::new(InMemory::new()),
174                clock,
175                NewUser::passwordless(format!("peer{i}")),
176            )
177            .await?;
178            let key_id = user.add_private_key(Some(KEY_NAME)).await?;
179
180            instance.enable_sync().await?;
181            let sync = instance
182                .sync()
183                .expect("sync handle present immediately after enable_sync");
184            let address = self.transport.serve(&sync).await?;
185
186            peers.push(Peer {
187                instance,
188                user,
189                key_id,
190                sync,
191                address,
192            });
193        }
194        Ok(Cluster { peers })
195    }
196}
197
198/// A set of in-process eidetica peers wired for multi-peer sync. Each peer is a
199/// full `Instance` with its own backend. The cluster owns the wiring; the test
200/// owns the databases, the auth, and the order of operations.
201pub struct Cluster {
202    peers: Vec<Peer>,
203}
204
205impl Cluster {
206    /// Start building a cluster. Defaults: 2 peers, per-peer `FixedClock`,
207    /// [`HttpLoopback`] transport.
208    pub fn builder() -> ClusterBuilder {
209        ClusterBuilder {
210            peers: 2,
211            clock: None,
212            transport: Arc::new(HttpLoopback),
213        }
214    }
215
216    /// Number of peers in the cluster.
217    pub fn len(&self) -> usize {
218        self.peers.len()
219    }
220
221    /// Whether the cluster has no peers.
222    pub fn is_empty(&self) -> bool {
223        self.peers.is_empty()
224    }
225
226    /// Shared access to a peer (instance, user, key, address).
227    pub fn peer(&self, i: usize) -> &Peer {
228        &self.peers[i]
229    }
230
231    /// Mutable access to a peer, for `create_database` / `open_database` on its
232    /// `User`.
233    pub fn peer_mut(&mut self, i: usize) -> &mut Peer {
234        &mut self.peers[i]
235    }
236
237    /// Have peer `to` bootstrap `tree` from peer `from`: request the tree with
238    /// `to`'s own key (asking for `permission`), flush, and track it locally
239    /// (sync disabled — the test turns on [`serve`]/[`auto_sync`] if it wants
240    /// more). The joining-peer dance as one call.
241    ///
242    /// `permission` is the access level `to` requests; the harness does not pick
243    /// it — auth posture stays the test's. `from` must already be serving `tree`
244    /// (see [`Peer::serve`]) with a policy that admits this request.
245    ///
246    /// [`serve`]: Peer::serve
247    /// [`auto_sync`]: Cluster::auto_sync
248    pub async fn bootstrap(
249        &mut self,
250        from: usize,
251        to: usize,
252        tree: &ID,
253        permission: Permission,
254    ) -> Result<()> {
255        let from_addr = self.peers[from].address.clone();
256        let to_key = self.peers[to].key_id.clone();
257        let to_signing_key = self.peers[to]
258            .user
259            .get_signing_key(&to_key)
260            .expect("peer key must be unlocked");
261        self.peers[to]
262            .sync
263            .sync_with_peer_for_bootstrap_with_key(
264                &from_addr,
265                tree,
266                &to_signing_key,
267                KEY_NAME,
268                permission,
269            )
270            .await?;
271        self.peers[to].sync.flush().await?;
272        self.peers[to]
273            .user
274            .track_database(tree.clone(), &to_key, SyncSettings::disabled())
275            .await?;
276        Ok(())
277    }
278
279    /// Drive a sync exchange for `tree`, initiated by peer `from` against peer
280    /// `to`. eidetica's `sync_with_peer` exchanges in *both* directions, so after
281    /// this both peers hold each other's entries for the tree. Both sync queues
282    /// are flushed; flush errors propagate (this is a bug-finding tool — it does
283    /// not swallow them).
284    pub async fn exchange(&self, from: usize, to: usize, tree: &ID) -> Result<()> {
285        let to_addr = self.peers[to].address.clone();
286        self.peers[from]
287            .sync
288            .sync_with_peer(&to_addr, Some(tree))
289            .await?;
290        self.peers[from].sync.flush().await?;
291        self.peers[to].sync.flush().await?;
292        Ok(())
293    }
294
295    /// Turn on **background, automatic** sync of `tree` between peers `a` and `b`,
296    /// in both directions. After this a commit on either peer is queued for the
297    /// other automatically (sync-on-commit) — no per-write [`exchange`] call. Use
298    /// [`flush`] to push the queue immediately, or let the background interval
299    /// carry it.
300    ///
301    /// Both peers must already hold `tree` (e.g. one [`Peer::serve`]d it and the
302    /// other bootstrapped it). This wires the peer relationship both ways
303    /// (register peer + dial-back address + per-tree sync target) and re-tracks
304    /// the tree as `on_commit` on each side.
305    ///
306    /// [`exchange`]: Cluster::exchange
307    /// [`flush`]: Cluster::flush
308    pub async fn auto_sync(&mut self, a: usize, b: usize, tree: &ID) -> Result<()> {
309        let a_pub = self.peers[a].sync.get_device_pubkey()?;
310        let b_pub = self.peers[b].sync.get_device_pubkey()?;
311        let a_addr = self.peers[a].address.clone();
312        let b_addr = self.peers[b].address.clone();
313        let a_key = self.peers[a].key_id.clone();
314        let b_key = self.peers[b].key_id.clone();
315
316        // Each peer learns how to reach the other. A prior bootstrap may already
317        // have registered the peer, so registration is idempotent here.
318        register_peer_idempotent(&self.peers[a].sync, &b_pub).await?;
319        self.peers[a].sync.add_peer_address(&b_pub, b_addr).await?;
320        register_peer_idempotent(&self.peers[b].sync, &a_pub).await?;
321        self.peers[b].sync.add_peer_address(&a_pub, a_addr).await?;
322
323        // Each peer tracks the tree on-commit and targets the other for it.
324        self.peers[a]
325            .user
326            .track_database(tree.clone(), &a_key, SyncSettings::on_commit())
327            .await?;
328        self.peers[a].sync.add_tree_sync(&b_pub, tree).await?;
329        self.peers[b]
330            .user
331            .track_database(tree.clone(), &b_key, SyncSettings::on_commit())
332            .await?;
333        self.peers[b].sync.add_tree_sync(&a_pub, tree).await?;
334        Ok(())
335    }
336
337    /// [`auto_sync`] every peer pair in the cluster — a full mesh, so a commit on
338    /// any peer fans out to all the others. Both peers of every pair must already
339    /// hold `tree`. Does not flush; call [`flush_all`] to drain the setup pushes.
340    ///
341    /// [`auto_sync`]: Cluster::auto_sync
342    /// [`flush_all`]: Cluster::flush_all
343    pub async fn auto_sync_all(&mut self, tree: &ID) -> Result<()> {
344        let n = self.len();
345        for a in 0..n {
346            for b in (a + 1)..n {
347                self.auto_sync(a, b, tree).await?;
348            }
349        }
350        Ok(())
351    }
352
353    /// Push peer `peer`'s pending auto-sync queue to its targets now, instead of
354    /// waiting for the background interval. The deterministic barrier for
355    /// auto-sync tests: commit, `flush`, assert.
356    pub async fn flush(&self, peer: usize) -> Result<()> {
357        self.peers[peer].sync.flush().await
358    }
359
360    /// Drain the whole cluster: repeatedly [`flush`] every peer until all pending
361    /// auto-sync work has propagated everywhere, then settle.
362    ///
363    /// A single pass over the peers is **not** enough in general. `flush` visits
364    /// each peer once, in index order, so a pass advances an in-flight entry at
365    /// most one hop along its sync path (and only in the index direction — an
366    /// entry that must travel "backwards", from a higher-indexed peer to a lower
367    /// one, waits for the next pass). One pass suffices only when every peer
368    /// pushes directly to every other (a full mesh); a sparser topology — a relay
369    /// chain — needs up to one pass per hop. `len() + 1` passes covers the worst
370    /// case, since no propagation path through `len()` peers is longer than
371    /// `len() - 1` hops. This is the whole-cluster barrier: after it, every
372    /// deliverable entry has reached every peer.
373    ///
374    /// [`flush`]: Cluster::flush
375    pub async fn flush_all(&self) -> Result<()> {
376        for _ in 0..(self.len() + 1) {
377            for peer in 0..self.len() {
378                self.flush(peer).await?;
379            }
380        }
381        Ok(())
382    }
383
384    /// The [`Snapshot`] peer `peer` currently holds for `tree` — the canonical
385    /// (sorted, deduplicated) tip set identifying its state. [`Snapshot::EMPTY`]
386    /// if the peer has never seen the tree.
387    pub async fn snapshot(&self, peer: usize, tree: &ID) -> Result<Snapshot> {
388        self.peers[peer].instance.backend().snapshot(tree).await
389    }
390
391    /// True if the named `peers` all agree on `tree`'s [`Snapshot`] — the
392    /// convergence invariant. The caller names which peers should have converged;
393    /// a peer that never received the tree has an empty snapshot and will not
394    /// match. Comparison is `Snapshot` set-equality, so tip order never matters.
395    pub async fn converged(&self, peers: &[usize], tree: &ID) -> Result<bool> {
396        let mut reference: Option<Snapshot> = None;
397        for &i in peers {
398            let snapshot = self.snapshot(i, tree).await?;
399            match &reference {
400                None => reference = Some(snapshot),
401                Some(r) if *r != snapshot => return Ok(false),
402                Some(_) => {}
403            }
404        }
405        Ok(true)
406    }
407
408    /// Whether *every* peer agrees on `tree`'s tip set — the common convergence
409    /// check. Shorthand for [`converged`] over all peers; the explicit
410    /// `&[peers]` form stays for partition tests that expect only a subset to
411    /// agree.
412    ///
413    /// [`converged`]: Cluster::converged
414    pub async fn converged_all(&self, tree: &ID) -> Result<bool> {
415        let all: Vec<usize> = (0..self.peers.len()).collect();
416        self.converged(&all, tree).await
417    }
418
419    /// Drive bidirectional [`exchange`] across every peer pair, round after
420    /// round, until the whole cluster holds an identical tip set for `tree` —
421    /// then return `true`. Bounded to `peers` rounds (a complete graph converges
422    /// in one, the budget is slack for safety); returns the final convergence
423    /// status if the budget is spent without settling.
424    ///
425    /// Quiescent only: there must be no concurrent writes while this runs (it
426    /// has no way to observe them). Every peer must already hold and serve
427    /// `tree` so it can answer an exchange — `bootstrap` then [`Peer::serve`] on
428    /// each joiner. The fixpoint barrier the N-peer / partition-heal tests
429    /// assert against.
430    ///
431    /// [`exchange`]: Cluster::exchange
432    pub async fn converge(&self, tree: &ID) -> Result<bool> {
433        let n = self.peers.len();
434        for _ in 0..n.max(1) {
435            if self.converged_all(tree).await? {
436                return Ok(true);
437            }
438            for i in 0..n {
439                for j in (i + 1)..n {
440                    self.exchange(i, j, tree).await?;
441                }
442            }
443        }
444        self.converged_all(tree).await
445    }
446
447    // ===== invariant assertions =====
448    //
449    // Tip-set equality ([`converged`]) proves two peers *agree*, but it is a weak
450    // invariant: it says nothing about *what* they agreed on. Two peers can share
451    // a tip set yet differ below it, or converge onto a state that quietly dropped
452    // a signed entry, or store a received entry as `Failed`. These walk the full
453    // entry set behind the tips and assert the properties tip equality misses.
454    // They panic (not return `false`) with a diagnostic — invariant violation is a
455    // test failure, and the message should name the offending peer and entry.
456
457    /// The concrete local backend engine for peer `peer`. `Cluster` peers always
458    /// run on an in-memory backend, so the off-seam raw reads the invariant checks
459    /// need — the full entry dump ([`BackendImpl::get_tree`]) and per-entry
460    /// verification status — are always reachable through it.
461    fn local_engine(&self, peer: usize) -> Arc<dyn BackendImpl> {
462        self.peers[peer]
463            .instance
464            .backend()
465            .local_engine()
466            .expect("Cluster peers run on a local in-memory backend")
467    }
468
469    /// Every entry peer `peer` holds for `tree`, in id order. The full DAG of the
470    /// tree — settings, auth, and every store — not just the tips.
471    pub async fn entries(&self, peer: usize, tree: &ID) -> Result<Vec<Entry>> {
472        let mut entries = self.local_engine(peer).get_tree(tree).await?;
473        entries.sort_by_key(|e| e.id());
474        Ok(entries)
475    }
476
477    /// The id of every entry peer `peer` holds for `tree`, sorted.
478    pub async fn entry_ids(&self, peer: usize, tree: &ID) -> Result<Vec<ID>> {
479        Ok(self
480            .entries(peer, tree)
481            .await?
482            .into_iter()
483            .map(|e| e.id())
484            .collect())
485    }
486
487    /// Assert no peer in `peers` is missing an entry another holds for `tree` —
488    /// the merge converged onto the *union* of histories, never silently dropping
489    /// one peer's signed entry. Stronger than [`converged`], which only compares
490    /// tips.
491    ///
492    /// Limitation: if *every* peer dropped the same entry the union is also short
493    /// it, so this can't see that loss — use [`assert_all_present`] with an
494    /// externally-known id set for the absolute form.
495    ///
496    /// [`converged`]: Cluster::converged
497    /// [`assert_all_present`]: Cluster::assert_all_present
498    pub async fn assert_no_lost_entries(&self, peers: &[usize], tree: &ID) -> Result<()> {
499        use std::collections::BTreeSet;
500        let mut union: BTreeSet<ID> = BTreeSet::new();
501        let mut per_peer: Vec<(usize, BTreeSet<ID>)> = Vec::with_capacity(peers.len());
502        for &p in peers {
503            let ids: BTreeSet<ID> = self.entry_ids(p, tree).await?.into_iter().collect();
504            union.extend(ids.iter().cloned());
505            per_peer.push((p, ids));
506        }
507        for (p, ids) in &per_peer {
508            let missing: Vec<&ID> = union.difference(ids).collect();
509            assert!(
510                missing.is_empty(),
511                "peer {p} lost {} entr{} other peers hold for the tree: {missing:?}",
512                missing.len(),
513                if missing.len() == 1 { "y" } else { "ies" },
514            );
515        }
516        Ok(())
517    }
518
519    /// Assert every id in `expected` is present on every peer in `peers`. The
520    /// absolute form of [`assert_no_lost_entries`]: the test names entries it knows
521    /// were committed (e.g. ids captured from its own writes) and demands they
522    /// survive the merge everywhere.
523    ///
524    /// [`assert_no_lost_entries`]: Cluster::assert_no_lost_entries
525    pub async fn assert_all_present(
526        &self,
527        peers: &[usize],
528        tree: &ID,
529        expected: &[ID],
530    ) -> Result<()> {
531        for &p in peers {
532            let ids: std::collections::BTreeSet<ID> =
533                self.entry_ids(p, tree).await?.into_iter().collect();
534            let missing: Vec<&ID> = expected.iter().filter(|id| !ids.contains(id)).collect();
535            assert!(
536                missing.is_empty(),
537                "peer {p} is missing expected entries: {missing:?}",
538            );
539        }
540        Ok(())
541    }
542
543    /// Assert every entry peer `peer` holds for `tree` carries a well-formed
544    /// signature. A synced CRDT under global auth must never store an unsigned or
545    /// malformed-signature entry; this catches one that slipped through.
546    pub async fn assert_all_signed(&self, peer: usize, tree: &ID) -> Result<()> {
547        for e in self.entries(peer, tree).await? {
548            assert!(
549                !e.auth().is_unsigned(),
550                "peer {peer} holds an unsigned entry: {}",
551                e.id(),
552            );
553            if let Some(reason) = e.auth().malformed_reason() {
554                panic!(
555                    "peer {peer} holds a malformed-signature entry {}: {reason}",
556                    e.id(),
557                );
558            }
559        }
560        Ok(())
561    }
562
563    /// Assert no entry peer `peer` holds for `tree` is in the `Failed` verification
564    /// state — every entry, including those received over sync, verified against
565    /// the tree's auth. Stronger than tip equality: a peer can converge on the
566    /// right tips while having stored a received entry that does not verify.
567    ///
568    /// This is only a meaningful convergence invariant once sync runs a per-entry
569    /// verification pass that promotes received entries after their signing
570    /// context arrives. On a build where sync ingestion records a placeholder
571    /// status instead of a real signature check (see the TODO on
572    /// [`VerificationStatus`] and `docs/src/design/verification.md`), the stored
573    /// status does not reflect verification and this assertion should not be used
574    /// — a bootstrapped peer legitimately holds entries marked `Failed` that no
575    /// pass has yet promoted. Provided for the harness's forward path: exercise it
576    /// once verification-on-ingest is in place.
577    ///
578    /// [`VerificationStatus`]: crate::backend::VerificationStatus
579    pub async fn assert_all_verified(&self, peer: usize, tree: &ID) -> Result<()> {
580        let engine = self.local_engine(peer);
581        for e in self.entries(peer, tree).await? {
582            let status = engine.get_verification_status(&e.id()).await?;
583            assert!(
584                matches!(status, VerificationStatus::Verified),
585                "peer {peer} stored entry {} as {status:?}, expected Verified",
586                e.id(),
587            );
588        }
589        Ok(())
590    }
591}
592
593/// One peer in a [`Cluster`]: a full `Instance` plus the handles a test needs to
594/// act as that peer. It does **not** hold the peer's application databases — the
595/// test opens and keeps those.
596pub struct Peer {
597    instance: Instance,
598    user: User,
599    key_id: PublicKey,
600    sync: Arc<Sync>,
601    address: Address,
602}
603
604impl Peer {
605    /// This peer's `Instance`.
606    pub fn instance(&self) -> &Instance {
607        &self.instance
608    }
609
610    /// This peer's logged-in user session.
611    pub fn user(&self) -> &User {
612        &self.user
613    }
614
615    /// Mutable user session, for `create_database` / `open_database` directly.
616    pub fn user_mut(&mut self) -> &mut User {
617        &mut self.user
618    }
619
620    /// This peer's signing key id (the `SigKey` for its database operations).
621    pub fn key_id(&self) -> &PublicKey {
622        &self.key_id
623    }
624
625    /// The display name of this peer's signing key, for naming it in a bootstrap
626    /// request.
627    pub fn key_name(&self) -> &str {
628        KEY_NAME
629    }
630
631    /// This peer's sync handle.
632    pub fn sync(&self) -> &Arc<Sync> {
633        &self.sync
634    }
635
636    /// The address other peers use to reach this peer.
637    pub fn address(&self) -> &Address {
638        &self.address
639    }
640
641    /// Mark `tree` sync-enabled on this peer so its sync handler will serve it to
642    /// bootstrapping peers. Uses the user-opened [`Database`] handle to write the user's
643    /// preference; the instance callback reconciles the host's combined sync
644    /// state through the same path a real consumer takes. The database must
645    /// already be tracked (it is, on any
646    /// peer that created it via `create_database` or joined it via
647    /// [`Cluster::bootstrap`]). Pure plumbing: set whatever auth the test needs on
648    /// the database *before* calling this.
649    ///
650    pub async fn serve(&mut self, tree: &ID) -> Result<()> {
651        self.user.open_database(tree).await?.share().await
652    }
653}
654
655/// Register `pubkey` as a peer of `sync`, treating an already-registered peer as
656/// success — a prior bootstrap commonly registers it first.
657async fn register_peer_idempotent(sync: &Sync, pubkey: &PublicKey) -> Result<()> {
658    match sync.register_peer(pubkey, Some("peer")).await {
659        Ok(()) => Ok(()),
660        Err(crate::Error::Sync(e))
661            if matches!(*e, crate::sync::error::SyncError::PeerAlreadyExists(_)) =>
662        {
663            Ok(())
664        }
665        Err(e) => Err(e),
666    }
667}
668
669// ===== auth tools (policy-neutral: apply the keys the caller passes) =====
670
671/// Apply per-key auth to a database via a settings transaction. The caller
672/// chooses the keys and permissions; this just writes them.
673pub async fn add_auth_keys(db: &Database, keys: &[(&PublicKey, AuthKey)]) -> Result<()> {
674    let txn = db.new_transaction().await?;
675    let settings = txn.get_settings()?;
676    for (pubkey, key) in keys {
677        settings.set_auth_key(pubkey, key.clone()).await?;
678    }
679    txn.commit().await?;
680    Ok(())
681}
682
683/// Set the global (wildcard) auth key on a database via a settings transaction.
684/// The caller chooses the permission level.
685pub async fn set_global_auth_key(db: &Database, key: AuthKey) -> Result<()> {
686    let txn = db.new_transaction().await?;
687    let settings = txn.get_settings()?;
688    settings.set_global_auth_key(key).await?;
689    txn.commit().await?;
690    Ok(())
691}
692
693// ===== SimTransport: in-memory, controllable transport (Tier 1 seam) =====
694//
695// `HttpLoopback` is real HTTP over loopback: it delivers in wired order, so
696// "convergence is order-independent" is unprovable and a partition can only be
697// modelled coarsely (don't call `exchange`). `SimTransport` swaps the one seam
698// the harness left open — [`TestTransport`] — for an in-process fabric that
699// routes a [`SyncRequest`] straight to the target peer's [`SyncHandler`]: no
700// sockets, no ports, deterministic, and *controllable*. A test holds a
701// [`SimNetwork`] handle and partitions links mid-run.
702//
703// This is Tier 1 of the harness. It does not replace Tier 0 — it plugs into it:
704// `Cluster::builder().transport(Arc::new(SimLoopback::new(net.clone())))`.
705
706/// In-memory message fabric shared by every [`SimTransport`] in a cluster, and
707/// the control handle a test uses to inject faults. A drop-in for
708/// [`HttpLoopback`] via [`ClusterBuilder::transport`] that additionally lets a
709/// test [`partition`] links and [`heal`] them.
710///
711/// `Clone` is a shared handle (an `Arc` inside): the copy a test keeps and the
712/// copies inside each peer's transport all see the same fabric.
713///
714/// [`partition`]: SimNetwork::partition
715/// [`heal`]: SimNetwork::heal
716#[derive(Clone, Default)]
717pub struct SimNetwork {
718    inner: Arc<std::sync::Mutex<SimState>>,
719}
720
721#[derive(Default)]
722struct SimState {
723    /// Peer address -> that peer's serving handler, populated when it serves.
724    handlers: std::collections::HashMap<String, Arc<dyn SyncHandler>>,
725    /// Directed links currently dropping traffic: `(from_addr, to_addr)`.
726    blocked: std::collections::HashSet<(String, String)>,
727    /// Monotonic id source for peer addresses.
728    next_id: usize,
729    /// When set, `SendEntries` pushes are captured in `queue` instead of being
730    /// delivered to the receiver's handler inline. Request/response traffic
731    /// (handshake, tree-sync) always delivers inline regardless.
732    manual_delivery: bool,
733    /// Captured, not-yet-delivered messages, in send order. Only populated in
734    /// manual-delivery mode.
735    queue: Vec<InFlight>,
736    /// Monotonic id source for captured messages (delivery handles).
737    next_seq: usize,
738}
739
740/// A `SendEntries` push captured in manual-delivery mode, awaiting an explicit
741/// [`SimNetwork::deliver_one`] / [`deliver`](SimNetwork::deliver) / etc. The
742/// sender already received an optimistic `Ack`, so this models a message
743/// in-flight on the wire: the network decides when, in what order, and how many
744/// times the receiver actually sees it.
745struct InFlight {
746    /// Stable delivery handle, unique for the life of the fabric.
747    seq: usize,
748    /// Sender's sim address (becomes the receiver's `remote_address`).
749    from: String,
750    /// Receiver's sim address (whose handler will process the request).
751    to: String,
752    /// The captured request — always a `SyncRequest::SendEntries`.
753    request: SyncRequest,
754}
755
756impl SimNetwork {
757    /// A fresh, empty fabric.
758    pub fn new() -> Self {
759        Self::default()
760    }
761
762    fn lock(&self) -> std::sync::MutexGuard<'_, SimState> {
763        self.inner.lock().expect("SimNetwork mutex poisoned")
764    }
765
766    /// Hand out the next unique peer address (`sim-peer-N`, in serve order).
767    fn alloc_address(&self) -> String {
768        let mut s = self.lock();
769        let id = s.next_id;
770        s.next_id += 1;
771        format!("sim-peer-{id}")
772    }
773
774    fn register(&self, address: &str, handler: Arc<dyn SyncHandler>) {
775        self.lock().handlers.insert(address.to_string(), handler);
776    }
777
778    fn unregister(&self, address: &str) {
779        self.lock().handlers.remove(address);
780    }
781
782    /// Clone out the handler for `address` (drops the lock before any await).
783    fn handler_for(&self, address: &str) -> Option<Arc<dyn SyncHandler>> {
784        self.lock().handlers.get(address).cloned()
785    }
786
787    fn is_blocked(&self, from: &str, to: &str) -> bool {
788        self.lock()
789            .blocked
790            .contains(&(from.to_string(), to.to_string()))
791    }
792
793    /// If manual-delivery is on, capture `request` as an in-flight message from
794    /// `from` to `to` and return `true` (the caller answers the sender with an
795    /// optimistic `Ack`). Otherwise return `false` and let the caller deliver
796    /// inline. Called only for `SendEntries`.
797    fn capture(&self, from: &str, to: &Address, request: &SyncRequest) -> bool {
798        let mut s = self.lock();
799        if !s.manual_delivery {
800            return false;
801        }
802        let seq = s.next_seq;
803        s.next_seq += 1;
804        s.queue.push(InFlight {
805            seq,
806            from: from.to_string(),
807            to: to.address.clone(),
808            request: request.clone(),
809        });
810        true
811    }
812
813    /// Drop all traffic between `a` and `b` in *both* directions until [`heal`].
814    /// A send across a blocked link fails as a connection error, so an
815    /// auto-sync peer's queued entries stay pending and redeliver after heal —
816    /// a message-level partition, finer than withholding `exchange` calls.
817    ///
818    /// [`heal`]: SimNetwork::heal
819    pub fn partition(&self, a: &Address, b: &Address) {
820        let mut s = self.lock();
821        s.blocked.insert((a.address.clone(), b.address.clone()));
822        s.blocked.insert((b.address.clone(), a.address.clone()));
823    }
824
825    /// Restore traffic between `a` and `b` (both directions).
826    pub fn heal(&self, a: &Address, b: &Address) {
827        let mut s = self.lock();
828        s.blocked.remove(&(a.address.clone(), b.address.clone()));
829        s.blocked.remove(&(b.address.clone(), a.address.clone()));
830    }
831
832    /// Restore every link in the fabric.
833    pub fn heal_all(&self) {
834        self.lock().blocked.clear();
835    }
836
837    // ----- store-and-forward delivery control -----
838    //
839    // By default the fabric delivers every request inline (synchronous, in wired
840    // order) — same as `HttpLoopback`, just without sockets. Turn on
841    // *manual delivery* and `SendEntries` pushes are instead captured as
842    // [`InFlight`] messages, and the test decides when each one reaches its
843    // receiver. The sender still gets an immediate `Ack`, so an undelivered
844    // message looks delivered-from-the-sender's-side — a message that's left the
845    // sender but not yet arrived. This is what makes reorder, duplicate, and
846    // selective drop expressible; a real socket transport can't hold a message
847    // mid-flight under test control.
848    //
849    // Handshake and tree-sync are request/response and carry data the caller
850    // needs back, so they always deliver inline — only the fire-and-forget
851    // `SendEntries` push is deferrable. Bootstrap and `exchange` therefore work
852    // unchanged in manual mode; only auto-sync's entry pushes get captured.
853
854    /// Capture `SendEntries` pushes instead of delivering them inline (`true`),
855    /// or return to inline delivery (`false`). Flip this *after* setup
856    /// (bootstrap / `auto_sync`) so only the entry pushes a test cares about get
857    /// captured. Turning it back off does not flush the queue — already-captured
858    /// messages still need an explicit deliver.
859    ///
860    /// The two fault families compose with **partition taking precedence**: a
861    /// send across a [`partition`](Self::partition)ed link fails as a connection
862    /// error *before* capture is considered, so it parks in the retry queue rather
863    /// than the in-flight queue. Manual delivery only ever captures pushes on
864    /// links that are up. Tests generally use one family or the other, not both at
865    /// once.
866    pub fn set_manual_delivery(&self, manual: bool) {
867        self.lock().manual_delivery = manual;
868    }
869
870    /// The delivery handles of every captured-but-undelivered message, in the
871    /// order they were sent. Pass these to [`deliver`](Self::deliver),
872    /// [`duplicate`](Self::duplicate) or [`drop_message`](Self::drop_message) to
873    /// drive an out-of-order, duplicated, or lossy schedule.
874    pub fn pending(&self) -> Vec<usize> {
875        self.lock().queue.iter().map(|m| m.seq).collect()
876    }
877
878    /// Pop the captured message with handle `seq` (its request, sender, and
879    /// receiver). Returns `None` if no such message is queued.
880    fn take(&self, seq: usize) -> Option<InFlight> {
881        let mut s = self.lock();
882        let idx = s.queue.iter().position(|m| m.seq == seq)?;
883        Some(s.queue.remove(idx))
884    }
885
886    /// Deliver one captured message to its receiver: look up the receiver's
887    /// handler and run the request through it, exactly as an inline send would.
888    /// The response is discarded — the sender already got its optimistic `Ack`.
889    /// Returns `false` if the receiver is no longer serving (the message is
890    /// effectively lost).
891    async fn deliver_inflight(&self, msg: InFlight) -> bool {
892        let Some(handler) = self.handler_for(&msg.to) else {
893            return false;
894        };
895        let context = RequestContext {
896            remote_address: Some(Address::new(SimTransport::TRANSPORT_TYPE, msg.from)),
897            // Only tree-sync carries a pubkey, and that never enters the queue.
898            peer_pubkey: None,
899        };
900        handler.handle_request(&msg.request, &context).await;
901        true
902    }
903
904    /// Deliver the captured message `seq` to its receiver and remove it from the
905    /// queue. Delivering in an order other than [`pending`](Self::pending)
906    /// returned is how a test reorders the wire. Returns `false` if no such
907    /// message is queued (or its receiver has stopped).
908    pub async fn deliver(&self, seq: usize) -> bool {
909        match self.take(seq) {
910            Some(msg) => self.deliver_inflight(msg).await,
911            None => false,
912        }
913    }
914
915    /// Deliver the oldest captured message. Returns `false` when the queue is
916    /// empty.
917    pub async fn deliver_one(&self) -> bool {
918        let seq = match self.lock().queue.first() {
919            Some(m) => m.seq,
920            None => return false,
921        };
922        self.deliver(seq).await
923    }
924
925    /// Deliver every captured message in send order, draining the queue. Returns
926    /// the number delivered. (Delivery never enqueues more — the receiver's
927    /// handler processes entries, it doesn't push back through this fabric.)
928    pub async fn deliver_all(&self) -> usize {
929        let mut n = 0;
930        while self.deliver_one().await {
931            n += 1;
932        }
933        n
934    }
935
936    /// Clone captured message `seq` so it will be delivered a second time,
937    /// modelling a duplicate on the wire. The copy gets a fresh handle and lands
938    /// at the back of the queue; the original stays put. Returns the new handle,
939    /// or `None` if `seq` isn't queued. Idempotent sync must converge regardless.
940    pub fn duplicate(&self, seq: usize) -> Option<usize> {
941        let mut s = self.lock();
942        let original = s.queue.iter().find(|m| m.seq == seq)?;
943        let copy = InFlight {
944            seq: s.next_seq,
945            from: original.from.clone(),
946            to: original.to.clone(),
947            request: original.request.clone(),
948        };
949        let new_seq = copy.seq;
950        s.next_seq += 1;
951        s.queue.push(copy);
952        Some(new_seq)
953    }
954
955    /// Drop captured message `seq` without delivering it — a lost packet.
956    /// Returns `true` if a message was removed.
957    pub fn drop_message(&self, seq: usize) -> bool {
958        self.take(seq).is_some()
959    }
960
961    /// Drop every captured message without delivering. Returns the number lost.
962    pub fn drop_all(&self) -> usize {
963        let mut s = self.lock();
964        let n = s.queue.len();
965        s.queue.clear();
966        n
967    }
968}
969
970/// [`TestTransport`] backed by a [`SimNetwork`]: an in-memory drop-in for
971/// [`HttpLoopback`]. Build a cluster over it with
972/// `Cluster::builder().transport(Arc::new(SimLoopback::new(net.clone())))` and
973/// keep `net` to drive partitions.
974pub struct SimLoopback {
975    network: SimNetwork,
976}
977
978impl SimLoopback {
979    /// Wrap a [`SimNetwork`]. Share one network across the cluster (clone the
980    /// handle) so the test and every peer route through the same fabric.
981    pub fn new(network: SimNetwork) -> Self {
982        Self { network }
983    }
984}
985
986#[async_trait]
987impl TestTransport for SimLoopback {
988    async fn serve(&self, sync: &Sync) -> Result<Address> {
989        let address = self.network.alloc_address();
990        sync.register_transport(
991            "sim",
992            SimTransportBuilder {
993                address: address.clone(),
994                network: self.network.clone(),
995            },
996        )
997        .await?;
998        // accept_connections awaits StartServer, which calls start_server and
999        // registers our handler before returning — no post-serve race.
1000        sync.accept_connections().await?;
1001        Ok(Address::new(SimTransport::TRANSPORT_TYPE, address))
1002    }
1003}
1004
1005/// Builder that hands the peer's address + shared fabric to its [`SimTransport`].
1006struct SimTransportBuilder {
1007    address: String,
1008    network: SimNetwork,
1009}
1010
1011#[async_trait]
1012impl TransportBuilder for SimTransportBuilder {
1013    type Transport = SimTransport;
1014
1015    async fn build(self, _persisted: Doc) -> Result<(Self::Transport, Option<Doc>)> {
1016        Ok((
1017            SimTransport {
1018                address: self.address,
1019                network: self.network,
1020                running: AtomicBool::new(false),
1021            },
1022            None,
1023        ))
1024    }
1025}
1026
1027/// In-memory [`SyncTransport`]. Routes a [`SyncRequest`] straight to the target
1028/// peer's [`SyncHandler`] through the shared [`SimNetwork`] — no sockets, no
1029/// serialization — and honors the network's partition state.
1030pub struct SimTransport {
1031    /// This peer's own sim address (the key its handler is registered under).
1032    address: String,
1033    network: SimNetwork,
1034    running: AtomicBool,
1035}
1036
1037impl SimTransport {
1038    const TRANSPORT_TYPE: &'static str = "sim";
1039}
1040
1041#[async_trait]
1042impl SyncTransport for SimTransport {
1043    fn transport_type(&self) -> &'static str {
1044        Self::TRANSPORT_TYPE
1045    }
1046
1047    fn can_handle_address(&self, address: &Address) -> bool {
1048        address.transport_type == Self::TRANSPORT_TYPE
1049    }
1050
1051    async fn start_server(&self, handler: Arc<dyn SyncHandler>) -> Result<()> {
1052        self.network.register(&self.address, handler);
1053        self.running.store(true, Ordering::SeqCst);
1054        Ok(())
1055    }
1056
1057    async fn stop_server(&self) -> Result<()> {
1058        self.network.unregister(&self.address);
1059        self.running.store(false, Ordering::SeqCst);
1060        Ok(())
1061    }
1062
1063    async fn send_request(&self, address: &Address, request: &SyncRequest) -> Result<SyncResponse> {
1064        if !self.can_handle_address(address) {
1065            return Err(SyncError::UnsupportedTransport {
1066                transport_type: address.transport_type.clone(),
1067            }
1068            .into());
1069        }
1070        // A partitioned link looks like a connection failure to the sender; the
1071        // background sync layer keeps the entries queued for a later flush.
1072        if self.network.is_blocked(&self.address, &address.address) {
1073            return Err(SyncError::ConnectionFailed {
1074                address: address.address.clone(),
1075                reason: "sim link partitioned".to_string(),
1076            }
1077            .into());
1078        }
1079        // In manual-delivery mode, capture entry pushes as in-flight messages
1080        // and report success to the sender (an optimistic Ack — see
1081        // `SimNetwork::set_manual_delivery`). Request/response traffic falls
1082        // through to inline delivery so its result reaches the caller.
1083        if matches!(request, SyncRequest::SendEntries(_))
1084            && self.network.capture(&self.address, address, request)
1085        {
1086            return Ok(SyncResponse::Ack);
1087        }
1088
1089        let handler = self.network.handler_for(&address.address).ok_or_else(|| {
1090            SyncError::ConnectionFailed {
1091                address: address.address.clone(),
1092                reason: "no sim peer serving this address".to_string(),
1093            }
1094        })?;
1095
1096        // Mirror the HTTP transport's context: only SyncTree carries a pubkey.
1097        let peer_pubkey = match request {
1098            SyncRequest::SyncTree(r) => r.peer_pubkey.clone(),
1099            _ => None,
1100        };
1101        let context = RequestContext {
1102            remote_address: Some(Address::new(Self::TRANSPORT_TYPE, self.address.clone())),
1103            peer_pubkey,
1104        };
1105        Ok(handler.handle_request(request, &context).await)
1106    }
1107
1108    fn is_server_running(&self) -> bool {
1109        self.running.load(Ordering::SeqCst)
1110    }
1111
1112    fn get_server_address(&self) -> Result<String> {
1113        if self.running.load(Ordering::SeqCst) {
1114            Ok(self.address.clone())
1115        } else {
1116            Err(SyncError::ServerNotRunning.into())
1117        }
1118    }
1119}
1120
1121#[cfg(test)]
1122mod tests {
1123    use super::*;
1124    use crate::{crdt::Doc, store::DocStore};
1125
1126    async fn write(db: &Database, key: &str, value: &str) -> Result<()> {
1127        let tx = db.new_transaction().await?;
1128        tx.get_store::<DocStore>("data")
1129            .await?
1130            .set_string(key, value)
1131            .await?;
1132        tx.commit().await?;
1133        Ok(())
1134    }
1135
1136    async fn read(db: &Database, key: &str) -> Result<String> {
1137        let tx = db.new_transaction().await?;
1138        tx.get_store::<DocStore>("data")
1139            .await?
1140            .get_string(key)
1141            .await
1142    }
1143
1144    /// Peer 0 creates a database (auth chosen by the test) and serves it; it
1145    /// writes `a`; peer 1 bootstraps and opens it. Returns the room id and both
1146    /// peers' open handles — the shared starting point for the tests below.
1147    async fn shared_room(net: &mut Cluster) -> Result<(ID, Database, Database)> {
1148        let key0 = net.peer(0).key_id().clone();
1149        let device0 = net.peer(0).instance().id();
1150        let mut settings = Doc::new();
1151        settings.set("name", "chat");
1152        let db0 = net
1153            .peer_mut(0)
1154            .user_mut()
1155            .create_database(settings, &key0)
1156            .await?;
1157        let room = db0.root_id().clone();
1158
1159        add_auth_keys(
1160            &db0,
1161            &[
1162                (&key0, AuthKey::active(Some("admin"), Permission::Admin(10))),
1163                (
1164                    &device0,
1165                    AuthKey::active(Some("device"), Permission::Admin(10)),
1166                ),
1167            ],
1168        )
1169        .await?;
1170        set_global_auth_key(&db0, AuthKey::active(None, Permission::Admin(10))).await?;
1171        net.peer_mut(0).serve(&room).await?;
1172        write(&db0, "a", "from-peer-0").await?;
1173
1174        net.bootstrap(0, 1, &room, Permission::Write(10)).await?;
1175        let db1 = net.peer_mut(1).user_mut().open_database(&room).await?;
1176        assert_eq!(
1177            read(&db1, "a").await?,
1178            "from-peer-0",
1179            "bootstrap carries the write"
1180        );
1181        Ok((room, db0, db1))
1182    }
1183
1184    /// Manual mode: peer 1 writes, an explicit `exchange` brings peer 0 up to
1185    /// date, and both converge. Sync is fully ordered by the test.
1186    #[tokio::test]
1187    async fn exchange_round_trips_and_converges() -> Result<()> {
1188        let mut net = Cluster::builder().peers(2).build().await?;
1189        let (room, db0, db1) = shared_room(&mut net).await?;
1190
1191        write(&db1, "b", "from-peer-1").await?;
1192        net.exchange(1, 0, &room).await?;
1193
1194        assert_eq!(read(&db0, "a").await?, "from-peer-0");
1195        assert_eq!(read(&db0, "b").await?, "from-peer-1");
1196        assert!(net.converged(&[0, 1], &room).await?);
1197        Ok(())
1198    }
1199
1200    /// Background mode: after `auto_sync`, commits propagate on their own — no
1201    /// `exchange` per write. `flush` is the only barrier the test needs.
1202    #[tokio::test]
1203    async fn auto_sync_propagates_on_commit() -> Result<()> {
1204        let mut net = Cluster::builder().peers(2).build().await?;
1205        let (room, db0, db1) = shared_room(&mut net).await?;
1206
1207        net.auto_sync(0, 1, &room).await?;
1208
1209        // Peer 0 commits; no exchange call — auto-sync carries it.
1210        write(&db0, "c", "auto-from-0").await?;
1211        net.flush(0).await?;
1212        assert_eq!(read(&db1, "c").await?, "auto-from-0");
1213
1214        // And the reverse direction.
1215        write(&db1, "d", "auto-from-1").await?;
1216        net.flush(1).await?;
1217        assert_eq!(read(&db0, "d").await?, "auto-from-1");
1218
1219        assert!(net.converged(&[0, 1], &room).await?);
1220        Ok(())
1221    }
1222
1223    /// [`SimTransport`] is a drop-in for [`HttpLoopback`]: the same bootstrap +
1224    /// `exchange` flow converges over the in-memory fabric, no sockets involved.
1225    #[tokio::test]
1226    async fn sim_transport_is_a_drop_in() -> Result<()> {
1227        let mut net = Cluster::builder()
1228            .peers(2)
1229            .transport(Arc::new(SimLoopback::new(SimNetwork::new())))
1230            .build()
1231            .await?;
1232        let (room, db0, db1) = shared_room(&mut net).await?;
1233
1234        write(&db1, "b", "from-peer-1").await?;
1235        net.exchange(1, 0, &room).await?;
1236
1237        assert_eq!(read(&db0, "a").await?, "from-peer-0");
1238        assert_eq!(read(&db0, "b").await?, "from-peer-1");
1239        assert!(net.converged(&[0, 1], &room).await?);
1240        Ok(())
1241    }
1242
1243    /// A partition drops traffic at the message level: with `auto_sync` wired,
1244    /// a commit on peer 0 cannot reach peer 1 while the link is cut, then
1245    /// redelivers from the still-pending queue once the link heals. This is what
1246    /// `HttpLoopback` can't do — there, withholding `exchange` is the only
1247    /// partition, and it can't model "wired but not delivering".
1248    #[tokio::test]
1249    async fn sim_partition_blocks_then_heal_delivers() -> Result<()> {
1250        let fabric = SimNetwork::new();
1251        let mut net = Cluster::builder()
1252            .peers(2)
1253            .transport(Arc::new(SimLoopback::new(fabric.clone())))
1254            .build()
1255            .await?;
1256        let (room, db0, db1) = shared_room(&mut net).await?;
1257        let addr0 = net.peer(0).address().clone();
1258        let addr1 = net.peer(1).address().clone();
1259
1260        net.auto_sync(0, 1, &room).await?;
1261
1262        // Cut the link, then commit on peer 0. The flush attempt cannot reach
1263        // peer 1 (a connection error to the sender), so the entry stays queued.
1264        fabric.partition(&addr0, &addr1);
1265        write(&db0, "c", "during-partition").await?;
1266        let _ = net.flush(0).await; // send fails across the cut; entry remains queued
1267
1268        assert!(
1269            read(&db1, "c").await.is_err(),
1270            "peer 1 must not see the write while partitioned"
1271        );
1272        assert!(
1273            !net.converged(&[0, 1], &room).await?,
1274            "peers must diverge under partition"
1275        );
1276
1277        // Heal and flush again: the queued entry now reaches peer 1.
1278        fabric.heal(&addr0, &addr1);
1279        net.flush(0).await?;
1280
1281        assert_eq!(read(&db1, "c").await?, "during-partition");
1282        assert!(net.converged(&[0, 1], &room).await?);
1283        Ok(())
1284    }
1285
1286    /// Manual delivery captures a flushed push instead of delivering it: after
1287    /// `flush` the sender believes it sent (optimistic Ack), but the receiver
1288    /// sees nothing until the test releases the message. `deliver_all` then
1289    /// drains the wire and the peers converge.
1290    #[tokio::test]
1291    async fn sim_manual_delivery_holds_then_releases() -> Result<()> {
1292        let fabric = SimNetwork::new();
1293        let mut net = Cluster::builder()
1294            .peers(2)
1295            .transport(Arc::new(SimLoopback::new(fabric.clone())))
1296            .build()
1297            .await?;
1298        let (room, db0, db1) = shared_room(&mut net).await?;
1299        net.auto_sync(0, 1, &room).await?;
1300
1301        // Capture pushes from here on, then commit + flush: the entry leaves the
1302        // sender but is held on the wire.
1303        fabric.set_manual_delivery(true);
1304        write(&db0, "c", "held").await?;
1305        net.flush(0).await?;
1306
1307        assert_eq!(fabric.pending().len(), 1, "the push is captured, not lost");
1308        assert!(
1309            read(&db1, "c").await.is_err(),
1310            "receiver sees nothing until the message is delivered"
1311        );
1312
1313        // Release the wire: the held push reaches peer 1 and they converge.
1314        assert_eq!(fabric.deliver_all().await, 1);
1315        assert_eq!(read(&db1, "c").await?, "held");
1316        assert!(net.converged(&[0, 1], &room).await?);
1317        Ok(())
1318    }
1319
1320    /// Two causally-ordered pushes delivered *back to front*. A receiver that
1321    /// gets a child before its parent must still converge once both arrive —
1322    /// this is the order-independence that an in-wired transport can't probe.
1323    #[tokio::test]
1324    async fn sim_reordered_delivery_still_converges() -> Result<()> {
1325        let fabric = SimNetwork::new();
1326        let mut net = Cluster::builder()
1327            .peers(2)
1328            .transport(Arc::new(SimLoopback::new(fabric.clone())))
1329            .build()
1330            .await?;
1331        let (room, db0, db1) = shared_room(&mut net).await?;
1332        net.auto_sync(0, 1, &room).await?;
1333        fabric.set_manual_delivery(true);
1334
1335        // Two separate flushes => two distinct in-flight messages; the second
1336        // entry's parent is the first, so they're causally ordered.
1337        write(&db0, "first", "1").await?;
1338        net.flush(0).await?;
1339        write(&db0, "second", "2").await?;
1340        net.flush(0).await?;
1341
1342        let pending = fabric.pending();
1343        assert_eq!(pending.len(), 2, "two independent pushes on the wire");
1344
1345        // Deliver child-before-parent: reverse send order.
1346        for &seq in pending.iter().rev() {
1347            assert!(fabric.deliver(seq).await, "message {seq} should deliver");
1348        }
1349
1350        assert_eq!(read(&db1, "first").await?, "1");
1351        assert_eq!(read(&db1, "second").await?, "2");
1352        assert!(net.converged(&[0, 1], &room).await?);
1353        Ok(())
1354    }
1355
1356    /// A duplicated push: the same message delivered twice. Sync must be
1357    /// idempotent — the second copy is a no-op, not a corruption or a crash.
1358    #[tokio::test]
1359    async fn sim_duplicate_delivery_is_idempotent() -> Result<()> {
1360        let fabric = SimNetwork::new();
1361        let mut net = Cluster::builder()
1362            .peers(2)
1363            .transport(Arc::new(SimLoopback::new(fabric.clone())))
1364            .build()
1365            .await?;
1366        let (room, db0, db1) = shared_room(&mut net).await?;
1367        net.auto_sync(0, 1, &room).await?;
1368        fabric.set_manual_delivery(true);
1369
1370        write(&db0, "c", "once").await?;
1371        net.flush(0).await?;
1372
1373        let seq = fabric.pending()[0];
1374        fabric.duplicate(seq).expect("message is queued");
1375        assert_eq!(fabric.pending().len(), 2, "original plus its duplicate");
1376
1377        // Both copies delivered; the second is a redelivery of the same entries.
1378        assert_eq!(fabric.deliver_all().await, 2);
1379
1380        assert_eq!(read(&db1, "c").await?, "once");
1381        assert!(net.converged(&[0, 1], &room).await?);
1382        Ok(())
1383    }
1384
1385    /// A dropped push is a genuine loss — the sender got its Ack and won't
1386    /// retry, so the receiver stays behind. A later reconciling `exchange`
1387    /// (tree-sync, which delivers inline) repairs the divergence: lossy delivery
1388    /// doesn't strand the cluster as long as some full sync eventually runs.
1389    #[tokio::test]
1390    async fn sim_dropped_push_is_recovered_by_resync() -> Result<()> {
1391        let fabric = SimNetwork::new();
1392        let mut net = Cluster::builder()
1393            .peers(2)
1394            .transport(Arc::new(SimLoopback::new(fabric.clone())))
1395            .build()
1396            .await?;
1397        let (room, db0, db1) = shared_room(&mut net).await?;
1398        net.auto_sync(0, 1, &room).await?;
1399        fabric.set_manual_delivery(true);
1400
1401        write(&db0, "c", "dropped").await?;
1402        net.flush(0).await?;
1403
1404        let seq = fabric.pending()[0];
1405        assert!(fabric.drop_message(seq), "the push is dropped on the wire");
1406        assert!(
1407            read(&db1, "c").await.is_err(),
1408            "a dropped push never reaches the receiver"
1409        );
1410        assert!(
1411            !net.converged(&[0, 1], &room).await?,
1412            "the cluster diverges after a drop"
1413        );
1414
1415        // A reconciling tree-sync pulls what the dropped push lost.
1416        fabric.set_manual_delivery(false);
1417        net.exchange(1, 0, &room).await?;
1418
1419        assert_eq!(read(&db1, "c").await?, "dropped");
1420        assert!(net.converged(&[0, 1], &room).await?);
1421        Ok(())
1422    }
1423}