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}