diff --git a/Cargo.lock b/Cargo.lock index 7832a8a1..2a974160 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -457,12 +457,13 @@ checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06" [[package]] name = "before" version = "0.1.0" -source = "git+https://github.com/oxidecomputer/rumors?rev=90e4509dc66330affb38afa74813a8dd500edd93#90e4509dc66330affb38afa74813a8dd500edd93" +source = "git+https://github.com/oxidecomputer/rumors?rev=ab94cd77f36d32a08acd912d8cb7dc7f9017fa82#ab94cd77f36d32a08acd912d8cb7dc7f9017fa82" dependencies = [ "bytes", - "dashu-int", "dsi-bitstream", + "num-bigint", "serde", + "serde_bytes", "serde_json", "static_assertions", "suanpan", @@ -1040,25 +1041,6 @@ dependencies = [ "syn 2.0.119", ] -[[package]] -name = "dashu-base" -version = "0.5.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a64b04cdfc4c8533100fe00304eb9687173bda47b1f1dac8af12ba13712ed49d" - -[[package]] -name = "dashu-int" -version = "0.5.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b6ee98721d5d223e5b64b642dd9588b79d9ef415554b13720308b77d628c3be6" -dependencies = [ - "cfg-if", - "dashu-base", - "num-modular", - "rustversion", - "static_assertions", -] - [[package]] name = "data-encoding" version = "2.11.1" @@ -1135,17 +1117,6 @@ dependencies = [ "serde_core", ] -[[package]] -name = "derive-where" -version = "1.6.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d08b3a0bcc0d079199cd476b2cae8435016ec11d1c0986c6901c5ac223041534" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.119", -] - [[package]] name = "dice-mfg-msgs" version = "0.3.0" @@ -1418,18 +1389,6 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "869b0adbda23651a9c5c0c3d270aac9fcb52e8622a8f2b17e57802d7791962f2" -[[package]] -name = "enum-as-inner" -version = "0.6.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a1e6a265c649f3f5979b601d26f1d05ada116434c87741c9493cb56218f76cbc" -dependencies = [ - "heck 0.5.0", - "proc-macro2", - "quote", - "syn 2.0.119", -] - [[package]] name = "env_filter" version = "2.0.0" @@ -1775,6 +1734,12 @@ version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" +[[package]] +name = "hashbrown" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "43a3c133739dddd0d2990f9a4bdf8eb4b21ef50e4851ca85ab661199821d510e" + [[package]] name = "hashbrown" version = "0.14.5" @@ -1848,51 +1813,6 @@ version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" -[[package]] -name = "hickory-proto" -version = "0.24.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "92652067c9ce6f66ce53cc38d1169daa36e6e7eb7dd3b63b5103bd9d97117248" -dependencies = [ - "async-trait", - "cfg-if", - "data-encoding", - "enum-as-inner", - "futures-channel", - "futures-io", - "futures-util", - "idna", - "ipnet", - "once_cell", - "rand 0.8.7", - "thiserror 1.0.69", - "tinyvec", - "tokio", - "tracing", - "url", -] - -[[package]] -name = "hickory-resolver" -version = "0.24.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cbb117a1ca520e111743ab2f6688eddee69db4e0ea242545a604dce8a66fd22e" -dependencies = [ - "cfg-if", - "futures-util", - "hickory-proto", - "ipconfig", - "lru-cache", - "once_cell", - "parking_lot", - "rand 0.8.7", - "resolv-conf", - "smallvec 1.15.2", - "thiserror 1.0.69", - "tokio", - "tracing", -] - [[package]] name = "hkdf" version = "0.12.4" @@ -2292,19 +2212,6 @@ dependencies = [ "generic-array", ] -[[package]] -name = "ipconfig" -version = "0.3.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4d40460c0ce33d6ce4b0630ad68ff63d6661961c48b6dba35e5a4d81cfb48222" -dependencies = [ - "socket2", - "widestring", - "windows-registry", - "windows-result", - "windows-sys 0.61.2", -] - [[package]] name = "ipnet" version = "2.12.1" @@ -2572,12 +2479,6 @@ version = "0.2.16" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981" -[[package]] -name = "linked-hash-map" -version = "0.5.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0717cef1bc8b636c6e1c1bbdefc09e6322da8a9321966e8928ef80d20f7f770f" - [[package]] name = "linux-raw-sys" version = "0.4.15" @@ -2656,15 +2557,6 @@ dependencies = [ "hashbrown 0.17.1", ] -[[package]] -name = "lru-cache" -version = "0.1.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "31e24f1ad8321ca0e8a1e0ac13f23cb668e6f5466c2c57319f6a5cf1cc8e3b1c" -dependencies = [ - "linked-hash-map", -] - [[package]] name = "lru-slab" version = "0.1.2" @@ -2846,6 +2738,16 @@ version = "0.1.14" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72ef4a56884ca558e5ddb05a1d1e7e1bfd9a68d9ed024c21704cc98872dae1bb" +[[package]] +name = "num-bigint" +version = "0.4.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c89e69e7e0f03bea5ef08013795c25018e101932225a656383bd384495ecc367" +dependencies = [ + "num-integer", + "num-traits", +] + [[package]] name = "num-bigint-dig" version = "0.8.6" @@ -2888,12 +2790,6 @@ dependencies = [ "num-traits", ] -[[package]] -name = "num-modular" -version = "0.6.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bd8e500409e6cd603b03e477c26a6caecdc27ac58979a53e881c75eafc079f44" - [[package]] name = "num-primitive" version = "0.2.1" @@ -3469,25 +3365,6 @@ dependencies = [ "thiserror 1.0.69", ] -[[package]] -name = "qorb" -version = "0.4.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2f223a85f489c9548719fb367c6733e8bc9d4eef76d6ea59a18299acd9497c60" -dependencies = [ - "anyhow", - "async-trait", - "debug-ignore", - "derive-where", - "futures", - "hickory-resolver", - "rand 0.9.5", - "thiserror 2.0.20", - "tokio", - "tokio-stream", - "tracing", -] - [[package]] name = "quinn" version = "0.11.11" @@ -3794,12 +3671,6 @@ dependencies = [ "web-sys", ] -[[package]] -name = "resolv-conf" -version = "0.7.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1e061d1b48cb8d38042de4ae0a7a6401009d6143dc80d2e2d6f31f0bdd6470c7" - [[package]] name = "rfc6979" version = "0.4.0" @@ -3858,17 +3729,17 @@ dependencies = [ [[package]] name = "rumors" version = "0.1.0" -source = "git+https://github.com/oxidecomputer/rumors?rev=90e4509dc66330affb38afa74813a8dd500edd93#90e4509dc66330affb38afa74813a8dd500edd93" +source = "git+https://github.com/oxidecomputer/rumors?rev=ab94cd77f36d32a08acd912d8cb7dc7f9017fa82#ab94cd77f36d32a08acd912d8cb7dc7f9017fa82" dependencies = [ "async-stream", "before", "bytes", "ciborium", "futures", - "futures-util", "hex", "itertools", - "rand 0.8.7", + "rand 0.9.5", + "schnellru", "seq-macro", "serde", "sha3 0.12.0", @@ -3876,7 +3747,6 @@ dependencies = [ "thiserror 2.0.20", "tinyvec", "tokio", - "tokio-stream", ] [[package]] @@ -4122,6 +3992,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "schnellru" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "356285bbf17bea63d9e52e96bd18f039672ac92b55b8cb997d6162a2a37d1649" +dependencies = [ + "ahash", + "cfg-if", + "hashbrown 0.13.2", +] + [[package]] name = "scopeguard" version = "1.2.0" @@ -4220,6 +4101,16 @@ dependencies = [ "smallvec 0.6.14", ] +[[package]] +name = "serde_bytes" +version = "0.11.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a5d440709e79d88e51ac01c4b72fc6cb7314017bb7da9eeff678aa94c10e3ea8" +dependencies = [ + "serde", + "serde_core", +] + [[package]] name = "serde_core" version = "1.0.229" @@ -4783,10 +4674,7 @@ dependencies = [ [[package]] name = "suanpan" version = "0.1.0" -source = "git+https://github.com/oxidecomputer/rumors?rev=90e4509dc66330affb38afa74813a8dd500edd93#90e4509dc66330affb38afa74813a8dd500edd93" -dependencies = [ - "dashu-int", -] +source = "git+https://github.com/oxidecomputer/rumors?rev=ab94cd77f36d32a08acd912d8cb7dc7f9017fa82#ab94cd77f36d32a08acd912d8cb7dc7f9017fa82" [[package]] name = "subtle" @@ -4975,7 +4863,6 @@ dependencies = [ "p256", "percent-encoding", "pwd", - "qorb", "rand_core 0.6.4", "rumors", "rustix 1.1.4", @@ -5378,7 +5265,6 @@ dependencies = [ "futures-core", "pin-project-lite", "tokio", - "tokio-util", ] [[package]] @@ -5882,12 +5768,6 @@ dependencies = [ "rustls-pki-types", ] -[[package]] -name = "widestring" -version = "1.2.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "72069c3113ab32ab29e5584db3c6ec55d416895e60715417b5b883a357c3e471" - [[package]] name = "winapi" version = "0.3.9" diff --git a/Cargo.toml b/Cargo.toml index 597351bc..ad20af41 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -44,12 +44,11 @@ permission-slip-common = { git = "https://github.com/oxidecomputer/permission-sl progenitor = "0.14" progenitor-client = "0.14" pwd = "1" -qorb = { version = "0.4", default-features = false } rand = "0.8" rand_core = "0.6" reqwest = { version = "0.13", features = ["json", "rustls"] } rlimit = "0.10" -rumors = { git = "https://github.com/oxidecomputer/rumors", rev = "90e4509dc66330affb38afa74813a8dd500edd93" } +rumors = { git = "https://github.com/oxidecomputer/rumors", rev = "ab94cd77f36d32a08acd912d8cb7dc7f9017fa82" } rustix = { version = "1", features = ["fs", "process", "pty", "termios"] } rustls = "0.23" rustyline = { git = "https://github.com/kkawakam/rustyline", rev = "5ab22f2f8d65e85f0b925360437ece8eeeb6db92", features = ["signal-hook", "termios"] } @@ -77,22 +76,29 @@ tokio-util = "0.7" x509-cert = { version = "0.2", features = ["pem", "std"] } xdg = "3" -# Every sprockets handshake verifies P-384 attestation cert chains, and every -# gossip link stream is a handshake. The RustCrypto stack is ~100ms per verify -# unoptimized, which makes the sprockets-backed tests crawl, so optimize the -# crates that do the math; everything else keeps fast debug builds. +# Optimize the TLS, signature, and hashing primitives used by fresh attested +# connections. The test profile inherits these overrides; Sush itself stays +# unoptimized for debugging. +[profile.dev.package.aws-lc-sys] +opt-level = 3 [profile.dev.package.crypto-bigint] opt-level = 3 +[profile.dev.package.curve25519-dalek] +opt-level = 3 [profile.dev.package.ecdsa] opt-level = 3 [profile.dev.package.ed25519-dalek] opt-level = 3 [profile.dev.package.elliptic-curve] opt-level = 3 +[profile.dev.package.keccak] +opt-level = 3 [profile.dev.package.p384] opt-level = 3 [profile.dev.package.primeorder] opt-level = 3 +[profile.dev.package.sha2] +opt-level = 3 [profile.dev.package.sha3] opt-level = 3 diff --git a/server/Cargo.toml b/server/Cargo.toml index a4b1bd44..06b5acb5 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -38,7 +38,6 @@ memmap2.workspace = true p256.workspace = true percent-encoding.workspace = true pwd.workspace = true -qorb.workspace = true rumors.workspace = true rand_core.workspace = true rustix.workspace = true @@ -64,6 +63,7 @@ tokio-tungstenite.workspace = true x509-cert.workspace = true [dev-dependencies] +tokio = { workspace = true, features = ["test-util"] } attest-mock.workspace = true ciborium.workspace = true function_name.workspace = true diff --git a/server/src/bookmark.rs b/server/src/bookmark.rs index eeecd7e4..b606cd1f 100644 --- a/server/src/bookmark.rs +++ b/server/src/bookmark.rs @@ -2,72 +2,103 @@ // License, v. 2.0. If a copy of the MPL was not distributed with this // file, You can obtain one at https://mozilla.org/MPL/2.0/. -//! Durable gossip peer identity across restarts. +//! Store Rumors restart bookkeeping in this server's locker. //! -//! A rumors [`Bookmark`] records a peer's identity and how far it has -//! advanced, so that a restarted sled may reclaim its previous identity -//! instead of stranding it. The invariant is (as usual) that we must -//! never adopt stale data, because in this case it could lead to causality -//! violations (which are bad). -//! -//! The record format and when to load & store are dictated by rumors. -//! We use a [`Tenant`] of a [`Locker`] to store it on disk(s); -//! if a load fails or the slots disagree, we assume a new identity -//! rather than risk resuming with a stale one. The record keeps every -//! universe's identities, so a lost write costs at most a stranded -//! identity, never a stale one. - -use std::io::{self, Cursor}; +//! Rumors supplies an opaque record that lets it reclaim a departed peer's +//! identity after catching up. The locker stores the record across local disks +//! and rejects stale or conflicting copies. An unusable record starts fresh +//! bookkeeping; it must never cause us to reuse an uncertain identity. +//! The wrapper records the Rumors format version so an incompatible record is +//! discarded before it can repeatedly abort gossip at session startup. + +use std::io::Cursor; use std::sync::Arc; -use rumors::{Bookmark, BookmarkError, Serialized}; +use rumors::Bookmark; use slog::{Discard, Logger, o, warn}; -use thiserror::Error; -use tokio::io::AsyncWrite; use crate::format::{self, NoFormat, Record, Versioned}; use crate::locker::{Locker, StoreError, Tenant, TenantSpec, Verdict}; +/// Locker namespace and format marker for Rumors bookmark records. pub const BOOKMARK: TenantSpec = TenantSpec { file: "bookmark", magic: b"SUSHBOOKMARK", }; -/// Wrap the opaque bytes rumors writes. +/// The Rumors format version and its opaque restart bookkeeping. #[derive(serde::Deserialize, serde::Serialize)] -struct BookmarkRecord(#[serde(with = "format::cbor_bytes")] Vec); - +struct BookmarkRecord( + /// Identifies which Rumors codec can read these bytes. + u64, + /// Complete record supplied by Rumors. + #[serde(with = "format::cbor_bytes")] + Vec, +); + +/// Identify the Sush wrapper format, independently of the enclosed Rumors format. impl Versioned for BookmarkRecord { - const VERSION: u16 = 0; + /// The wrapper stores an explicit Rumors format version. + const VERSION: u16 = 1; } +/// Read the versioned wrapper through the locker's format chain. impl Record for BookmarkRecord { - type Previous = NoFormat; + /// The wrapper without an explicit Rumors version. + type Previous = v0::BookmarkRecord; } -impl TryFrom for BookmarkRecord { +/// An unlabelled record cannot establish which Rumors codec owns its bytes. +impl TryFrom for BookmarkRecord { + /// The reason recovery must start fresh. type Error = &'static str; - fn try_from(none: NoFormat) -> Result { - match none {} + /// Decline recovery rights whose format is unknown; replicated data is unaffected. + fn try_from(_: v0::BookmarkRecord) -> Result { + Err("bookmark has no Rumors format version") } } -/// What a bookmark load or store failed at. -#[derive(Debug, Error)] -pub enum BookmarkIoError { - #[error("serializing the bookmark record failed: {0}")] - Serialize(#[source] io::Error), - #[error("storing the bookmark failed: {0}")] - Store(#[source] StoreError), +/// Frozen wrapper format without a Rumors version marker. +mod v0 { + use super::*; + + /// Opaque bookmark bytes in wrapper format zero. + #[derive(serde::Deserialize, serde::Serialize)] + pub(super) struct BookmarkRecord(#[serde(with = "format::cbor_bytes")] pub(super) Vec); + + /// Identify this member of the durable format chain. + impl Versioned for BookmarkRecord { + /// The initial Sush wrapper format. + const VERSION: u16 = 0; + } + + /// Terminate the wrapper's history at its initial format. + impl Record for BookmarkRecord { + /// No preceding format exists. + type Previous = NoFormat; + } + + /// Complete the migration interface for the initial format. + impl TryFrom for BookmarkRecord { + /// No source value can exist. + type Error = &'static str; + /// The source type has no inhabitants. + fn try_from(none: NoFormat) -> Result { + match none {} + } + } } /// This server's bookmark storage. Every handle shares the one record. #[derive(Clone, Debug)] pub struct BookmarkSource { + /// Logger carrying this bookmark’s component context. log: Logger, + /// Shared locker record for this server across peer replacements. tenant: Arc, } +/// Create storage handles for the gossip manager. impl BookmarkSource { /// A source persisting to `locker`. /// [`Seed::grow`](crate::gossip::Seed::grow) makes the one source @@ -84,11 +115,8 @@ impl BookmarkSource { Self::new(&Logger::root(Discard, o!()), &Locker::null()) } - /// A persisting handle for a peer. Rumors persists a bookmark only - /// when a gossip session starts, and the gossip manager stops - /// every session before it hands a new peer its handle, so no two - /// peers persist concurrently; see the migration notes in - /// [`gossip`](crate::gossip). + /// A persisting handle for one peer. Before replacing it during migration, + /// the gossip manager waits for all sessions using the old handle to stop. pub fn handle(&self) -> SushBookmark { SushBookmark { log: self.log.clone(), @@ -111,18 +139,22 @@ impl BookmarkSource { /// One peer's handle on the [`BookmarkSource`]. #[derive(Debug)] pub struct SushBookmark { + /// Logger carrying this bookmark’s component context. log: Logger, + /// Shared locker record for this server across peer replacements. tenant: Arc, + /// Whether this handle bypasses storage after a persistence failure. shed: bool, } -impl BookmarkError for SushBookmark { - type Error = BookmarkIoError; -} - +/// Read locker records and store the complete bytes supplied by Rumors. impl Bookmark for SushBookmark { + /// Failure to durably store the locker record. + type Error = StoreError; + /// An owned snapshot of the stored Rumors record. type Reader = Cursor>; + /// Load a usable record, starting fresh if locker recovery rejects it. async fn load(&self) -> Result, Self::Error> { if self.shed { return Ok(None); @@ -131,7 +163,17 @@ impl Bookmark for SushBookmark { match guard.load().await { Verdict::Adopt(record) | Verdict::Restore(record) => { match format::decode::(&record) { - Ok(BookmarkRecord(bytes)) => Ok(Some(Cursor::new(bytes))), + Ok(BookmarkRecord(version, bytes)) + if version == rumors::BOOKMARK_FORMAT_VERSION => + { + Ok(Some(Cursor::new(bytes))) + } + Ok(BookmarkRecord(version, _)) => { + warn!(self.log, "starting fresh Rumors bookkeeping"; + "stored_format" => version, + "current_format" => rumors::BOOKMARK_FORMAT_VERSION); + Ok(None) + } // Stranding the old identity is harmless; // resuming from a misread record is not. Err(error) => { @@ -148,22 +190,14 @@ impl Bookmark for SushBookmark { } } - async fn store(&self, write: F) -> Result<(), Self::Error> - where - F: for<'a> FnOnce(&'a mut (dyn AsyncWrite + Unpin + Send)) -> Serialized<'a> + Send, - { + /// Store the encoded Rumors record unless this handle has shed persistence. + async fn store(&self, bytes: Vec) -> Result<(), Self::Error> { if self.shed { return Ok(()); } - let mut buf = Cursor::new(Vec::new()); - write(&mut buf).await.map_err(BookmarkIoError::Serialize)?; - let record = BookmarkRecord(buf.into_inner()); - + let record = BookmarkRecord(rumors::BOOKMARK_FORMAT_VERSION, bytes); let mut guard = self.tenant.lock().await; - guard - .store(&format::encode(&record)) - .await - .map_err(BookmarkIoError::Store) + guard.store(&format::encode(&record)).await } } @@ -175,14 +209,7 @@ mod test { use camino::Utf8PathBuf; use tempfile::TempDir; - use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _}; - - /// A serializer closure writing fixed bytes, shaped like rumors'. - fn record( - bytes: &'static [u8], - ) -> impl for<'a> FnOnce(&'a mut (dyn AsyncWrite + Unpin + Send)) -> Serialized<'a> + Send { - move |w| Box::pin(async move { w.write_all(bytes).await }) - } + use tokio::io::AsyncReadExt as _; /// Two slot directories, like two M.2s. fn slots(dir: &TempDir) -> Vec { @@ -196,14 +223,17 @@ mod test { .collect() } + /// Silence routine logging in storage tests. fn test_log() -> Logger { Logger::root(Discard, o!()) } + /// Create a bookmark source backed by the requested disk slots. fn source(slots: Vec) -> BookmarkSource { BookmarkSource::new(&test_log(), &Locker::new(&test_log(), slots).unwrap()) } + /// Read the exact bytes visible through the public Bookmark interface. async fn read_back(handle: &SushBookmark) -> Option> { let mut reader = handle.load().await.unwrap()?; let mut bytes = Vec::new(); @@ -211,6 +241,23 @@ mod test { Some(bytes) } + /// The locker-backed adapter preserves records across replacement and cancellation. + #[tokio::test] + async fn conforms_to_bookmark_contract() { + let mut directories = Vec::new(); + rumors::conformance::bookmark::check( + async || { + let dir = TempDir::with_prefix("sush-bookmark-conformance-").unwrap(); + let handle = source(slots(&dir)).handle(); + // Keep each fresh fixture's files until its check has finished. + directories.push(dir); + handle + }, + || tokio::time::sleep(std::time::Duration::from_secs(30)), + ) + .await; + } + /// A stored record loads back verbatim. #[tokio::test] async fn round_trip() { @@ -219,10 +266,43 @@ mod test { let handle = source.handle(); assert!(read_back(&handle).await.is_none()); - handle.store(record(b"who we are")).await.unwrap(); + handle.store(b"who we are".to_vec()).await.unwrap(); assert_eq!(read_back(&handle).await.unwrap(), b"who we are"); } + /// Unsupported bookmark formats start fresh and allow subsequent persistence. + #[tokio::test] + async fn incompatible_formats_allow_fresh_checkpoints() { + let dir = TempDir::with_prefix("sush-bookmark-").unwrap(); + let source = source(slots(&dir)); + let old = format::encode(&v0::BookmarkRecord(b"old".to_vec())); + let future = format::encode(&BookmarkRecord( + rumors::BOOKMARK_FORMAT_VERSION + 1, + b"future".to_vec(), + )); + for encoded in [old, future] { + source.tenant.lock().await.store(&encoded).await.unwrap(); + let handle = source.handle(); + assert!(read_back(&handle).await.is_none()); + handle.store(b"fresh".to_vec()).await.unwrap(); + assert_eq!(read_back(&source.handle()).await.unwrap(), b"fresh"); + } + } + + /// Pin both Sush wrapper layouts independently of Rumors' opaque payload. + #[test] + fn wrapper_formats_pin_their_bytes() { + assert_eq!( + format::encode(&v0::BookmarkRecord(b"old".to_vec())), + [0x82, 0x00, 0x44, 0x43, b'o', b'l', b'd'] + ); + // Version 6 is an example marker here, not a requirement on the linked codec. + assert_eq!( + format::encode(&BookmarkRecord(6, b"new".to_vec())), + [0x82, 0x01, 0x46, 0x82, 0x06, 0x43, b'n', b'e', b'w'] + ); + } + /// A discarded verdict is a fresh start, not an error. #[tokio::test] async fn discard_assumes_fresh_identity() { @@ -230,7 +310,7 @@ mod test { let slots = slots(&dir); for (slot, bytes) in slots.iter().zip([b"one", b"two"]) { let lone = source(vec![slot.clone()]); - lone.handle().store(record(bytes)).await.unwrap(); + lone.handle().store(bytes.to_vec()).await.unwrap(); } assert!(read_back(&source(slots).handle()).await.is_none()); } @@ -241,7 +321,7 @@ mod test { let dir = TempDir::with_prefix("sush-bookmark-").unwrap(); let source = source(slots(&dir)); - source.handle().store(record(b"shared")).await.unwrap(); + source.handle().store(b"shared".to_vec()).await.unwrap(); assert_eq!(read_back(&source.handle()).await.unwrap(), b"shared"); } @@ -251,15 +331,15 @@ mod test { async fn null_and_shed_touch_nothing() { let null = BookmarkSource::null(); let handle = null.handle(); - handle.store(record(b"lost")).await.unwrap(); + handle.store(b"lost".to_vec()).await.unwrap(); assert!(read_back(&handle).await.is_none()); let dir = TempDir::with_prefix("sush-bookmark-").unwrap(); let source = source(slots(&dir)); - source.handle().store(record(b"kept")).await.unwrap(); + source.handle().store(b"kept".to_vec()).await.unwrap(); let shed = source.shed_handle(); assert!(read_back(&shed).await.is_none()); - shed.store(record(b"dropped")).await.unwrap(); + shed.store(b"dropped".to_vec()).await.unwrap(); assert_eq!(read_back(&source.handle()).await.unwrap(), b"kept"); } } diff --git a/server/src/boundary.rs b/server/src/boundary.rs index 61ccfeda..d3f71912 100644 --- a/server/src/boundary.rs +++ b/server/src/boundary.rs @@ -366,11 +366,26 @@ mod test { serde_json::from_str(&format!("[{seed:?}{}]", ", 0".repeat(15))).unwrap() } + /// A nested history with event counts 2, 1, and 3 in successive regions. + /// + /// Keep these counts and forks fixed: the boundary format pin contains it. + fn executed_version() -> Version { + let mut left = rumors::before::Party::seed(); + let mut version = Version::new(); + left.tick(&mut version); + let mut right = left.fork(); + left.tick(&mut version); + let far_right = right.fork(); + far_right.ticks(&mut version, 2u64); + version + } + + /// A committed job and its nonempty execution history. fn boundary() -> Boundary { Boundary { network: network(1), burned: Bloom::new(), - executed: "(1, 1, (0, 0, 2))".parse().unwrap(), + executed: executed_version(), job: Some(Committed { session: SessionId::random(), job: JobId::random(), @@ -397,7 +412,7 @@ mod test { burned.insert(&network_key(network(2))); burned }, - executed: "(1, 1, (0, 0, 2))".parse().unwrap(), + executed: executed_version(), job: Some(Committed { session: "abandon-ability".parse().unwrap(), job: "zoo-zero".parse().unwrap(), diff --git a/server/src/format.rs b/server/src/format.rs index a6b0e94e..d586760d 100644 --- a/server/src/format.rs +++ b/server/src/format.rs @@ -174,22 +174,35 @@ pub(crate) mod cbor_bytes { use serde::{Deserializer, Serializer}; use std::fmt; + /// Encode the buffer as one byte string. pub fn serialize(bytes: &[u8], serializer: S) -> Result { serializer.serialize_bytes(bytes) } + /// Decode an owned buffer without imposing the decoder's scratch-buffer limit. pub fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result, D::Error> { + /// Accept borrowed or owned byte strings as an owned record. struct Bytes; + /// Preserve the complete byte string regardless of how the decoder supplies it. impl<'de> Visitor<'de> for Bytes { + /// The record's opaque bytes. type Value = Vec; + /// Describe the required CBOR value on a type mismatch. fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.write_str("bytes") } + /// Copy a decoder's temporary buffer. fn visit_bytes(self, v: &[u8]) -> Result { Ok(v.to_vec()) } + /// Take the decoder's allocation without another copy. + fn visit_byte_buf(self, v: Vec) -> Result { + Ok(v) + } } - deserializer.deserialize_bytes(Bytes) + // The result is owned. Requesting borrowed bytes would restrict + // Ciborium to records that fit in its fixed scratch buffer. + deserializer.deserialize_byte_buf(Bytes) } } @@ -289,6 +302,16 @@ mod test { assert_eq!(decoded, value); } + /// Records larger than the decoder's scratch buffer round-trip completely. + #[test] + fn large_record_round_trip() { + let value = TestV1 { + count: 7, + label: "x".repeat(16 * 1024), + }; + assert_eq!(decode::(&encode(&value)).unwrap(), value); + } + #[test] fn old_version_converts_up() { let old = encode(&TestV0 { count: 7 }); diff --git a/server/src/gossip.rs b/server/src/gossip.rs index f947763a..16b22fcb 100644 --- a/server/src/gossip.rs +++ b/server/src/gossip.rs @@ -6,7 +6,7 @@ //! //! Every peer seeds its own universe at startup, so peers meet across //! unrelated universes and their sessions fail with -//! [`Error::NetworkMismatch`]. That error carries everything both sides need +//! [`Mismatch::Network`]. That error carries everything both sides need //! to agree on which universe survives, without coordination: the greater //! minimum event count wins, and the lesser network id breaks ties. rumors //! documents the rule under `Peer`, "Bootstrapping without consensus". The @@ -28,8 +28,9 @@ use std::io; use std::net::{SocketAddr, SocketAddrV6}; use std::time::Duration; -use futures::StreamExt as _; -use rumors::{Error, Joined, Network, Peer, Rumors, Ticks}; +use futures::{Stream, StreamExt as _, stream}; +use rumors::error::Mismatch; +use rumors::{Changes, Error, Gossip, Joined, Network, Peer, Rumors, Ticks}; use serde::Serialize; use serde::de::DeserializeOwned; use sled_hardware_types::BaseboardId; @@ -37,8 +38,9 @@ use slog::{Logger, debug, info, o, warn}; use sprockets_tls::keys::SprocketsConfig; use tokio::sync::watch; use tokio::task::{AbortHandle, JoinSet}; -use tokio::time::{MissedTickBehavior, interval, timeout}; +use tokio::time::{Instant, MissedTickBehavior, Sleep, interval, interval_at, sleep, timeout}; use tokio::{select, spawn}; +use tokio_stream::wrappers::IntervalStream; use tokio_util::sync::CancellationToken; use rumors::link::routed::Endpoint; @@ -63,15 +65,36 @@ impl LinkedBaseboards { } } +/// Push changes promptly and probe each connection every ten seconds. +fn gossip_policy( + changes: Changes, +) -> impl Stream + Send { + let period = Duration::from_secs(10); + // Changes already requests the opening session. Delay the first heartbeat, + // and avoid a burst of overdue probes after a scheduling pause. + let mut heartbeat = interval_at(Instant::now() + period, period); + heartbeat.set_missed_tick_behavior(MissedTickBehavior::Delay); + stream::select( + changes.map(Gossip::from), + IntervalStream::new(heartbeat).map(|_| Gossip::Unconditionally), + ) +} + +/// Allow one second for each gossip, bootstrap, or retirement exchange. +fn session_deadline() -> Sleep { + sleep(Duration::from_secs(1)) +} + /// Manager timing. The defaults suit a rack, tests shrink them. #[derive(Clone, Debug)] pub struct GossipConfig { /// How often absent links are re-established. pub reconnect: Duration, - /// Ceiling on establishing one link, and on each data-stream dial - /// inside a live one. + /// Timeout for establishing a link, dialing a data stream, or initially + /// routing an incoming connection. pub connect_timeout: Duration, - /// Ceiling on one bootstrap join. + /// Ceiling on joining and attaching its bookmark, including local storage. + /// The wire exchange also has its own session deadline. pub join_timeout: Duration, } @@ -87,12 +110,13 @@ impl Default for GossipConfig { /// A gossip universe. #[derive(Clone, Debug)] -pub struct Universe { +pub struct Universe { /// The gossiped set. pub rumors: Rumors, } -impl Universe { +impl Universe { + /// Wrap the initial gossip network for publication. pub fn genesis(rumors: Rumors) -> Self { Self { rumors } } @@ -111,12 +135,12 @@ impl Universe { /// the pair whole, so a seed can never gossip against a source other /// than its own. #[derive(Debug)] -pub struct Seed { +pub struct Seed { rumors: Rumors, bookmarks: BookmarkSource, } -impl Seed { +impl Seed { /// Seed a fresh universe with this server as its only peer, over /// `locker`'s storage, making the locker's one [`BookmarkSource`]. /// @@ -129,14 +153,19 @@ impl Seed { /// is harmless. pub async fn grow(log: &Logger, locker: &Locker) -> Self where - T: DeserializeOwned + Serialize + Eq + Send + Sync + 'static, + T: DeserializeOwned + Serialize + Eq, { let bookmarks = BookmarkSource::new(log, locker); let handle = match locker.probe().await { Ok(()) => bookmarks.handle(), Err(_) => bookmarks.shed_handle(), }; - let rumors = match Peer::seed().bookmark(handle).await { + let rumors = match Peer::seed() + .gossip_when(gossip_policy) + .session_deadline(session_deadline) + .bookmark(handle) + .await + { Ok(peer) => peer.into_rumors(), Err(unbookmarked) => match unbookmarked.peer.bookmark(bookmarks.shed_handle()).await { Ok(peer) => peer.into_rumors(), @@ -146,6 +175,7 @@ impl Seed { Self { rumors, bookmarks } } + /// Borrow the seeded network without separating it from its bookmark source. pub fn rumors(&self) -> &Rumors { &self.rumors } @@ -260,7 +290,8 @@ enum Stopped { Failed, } -struct Manager { +/// Maintains peer links and publishes the current gossip network. +struct Manager { log: Logger, config: GossipConfig, endpoint: Endpoint, @@ -394,7 +425,7 @@ where let endpoint = self.endpoint.clone(); let handle = self.dials.spawn(async move { let result = match timeout(deadline, endpoint.link(peer)).await { - Ok(Ok(link)) => Ok(link), + Ok(Ok((_, link))) => Ok(link), Ok(Err(err)) => Err(err.to_string()), Err(_) => Err("link establishment timed out".to_string()), }; @@ -436,7 +467,7 @@ where } /// Spawn a session driver owning `link`: push our changes, and serve - /// whatever the peer initiates, until the link fails. + /// whatever the peer initiates, with periodic probes until the link fails. fn drive(&mut self, peer: SocketAddr, link: SprocketsLink) { debug!(self.log, "driving link"; "peer" => %peer); let rumors = self.rumors.clone(); @@ -449,31 +480,24 @@ where } } - /// Abandon our universe for the one `peer` belongs to, joining over the - /// fresh link to it. Every driver stops first: none may gossip across - /// the swap. On failure our universe is intact and the debt stands, so - /// the next link retries; either way all links are rebuilt, since the - /// old ones belong to the universe we are leaving. + /// Stop old sessions, then join the network reached through this fresh link. + /// On failure retain our network and retry later; all old links are rebuilt. /// - /// The new peer gets its own handle on the same bookmark storage. - /// That is safe because rumors persists a bookmark only when a - /// session starts, and aborting the drivers above ends every - /// session before the handle exists: the abandoned peer can never - /// store again. A store it already had in flight either loses to - /// the locker's sequence guard, or records a session that was - /// aborted before it sent anything, so nothing on the wire - /// outruns the record. If the received identity cannot be - /// persisted, we keep gossiping with a shed handle rather than - /// take the sled out of gossip; a stranded identity is harmless, - /// unlike a support shell that cannot reach a degraded rack. + /// Wait for the driver futures to drop before another peer uses the shared + /// bookmark storage. Locker orders any disk writes that outlive cancellation. + /// If attaching the bookmark fails, use a shed handle so storage failure + /// does not prevent gossip. async fn migrate(&mut self, peer: SocketAddr, mut link: SprocketsLink) { - self.drivers.abort_all(); + self.drivers.shutdown().await; self.live.clear(); info!( self.log, "joining the universe that beat ours"; "peer" => %peer, "ours" => %self.rumors.network(), ); - let bootstrap = Peer::bootstrap().bookmark(self.bookmarks.handle()); + let bootstrap = Peer::bootstrap() + .gossip_when(gossip_policy) + .session_deadline(session_deadline) + .bookmark(self.bookmarks.handle()); match timeout(self.config.join_timeout, bootstrap.join(&mut link)).await { Ok(Joined::Joined { peer }) => self.adopt(peer), Ok(Joined::Unbookmarked(unbookmarked)) => { @@ -519,21 +543,26 @@ async fn sessions( where T: DeserializeOwned + Serialize + Eq + Send + Sync + 'static, { - let ours = rumors.network(); - let mut driver = rumors.gossip_when(rumors.changes(), &mut link); + let mut driver = rumors.gossip(&mut link); while let Some(session) = driver.next().await { match session { Ok(_) => {} - Err(Error::NetworkMismatch { + Err(Error::Mismatch(Mismatch::Network { + local_network, + local_min_events, remote_network, remote_min_events, - local_min_events, - }) => { - let dominated = - remote_dominates(&local_min_events, &remote_min_events, ours, remote_network); + .. + })) => { + let dominated = remote_dominates( + &local_min_events, + &remote_min_events, + local_network, + remote_network, + ); debug!( log, "universe mismatch"; - "ours" => %ours, "theirs" => %remote_network, + "ours" => %local_network, "theirs" => %remote_network, "our_events" => %local_min_events, "their_events" => %remote_min_events, "we_lose" => dominated, @@ -544,13 +573,15 @@ where Stopped::Failed }; } - // A bookmark failure also stops every later session at the - // persist gate, so it deserves a warning where routine - // link churn does not. + // Storage failures need attention; reconnecting alone cannot repair them. Err(Error::Bookmark(error)) => { warn!(log, "bookmark failure stops gossip"; "error" => %error); return Stopped::Failed; } + Err(Error::Protocol(error)) => { + warn!(log, "gossip protocol violation"; "diagnostic" => ?error); + return Stopped::Failed; + } Err(err) => { debug!(log, "session failed"; "error" => %err); return Stopped::Failed; diff --git a/server/src/link.rs b/server/src/link.rs index 5549f81d..ee183306 100644 --- a/server/src/link.rs +++ b/server/src/link.rs @@ -2,45 +2,27 @@ // License, v. 2.0. If a copy of the MPL was not distributed with this // file, You can obtain one at https://mozilla.org/MPL/2.0/. -//! Sprockets as a rumors [`routed`] transport. +//! Sprockets transport for routed Rumors links. //! -//! [`rumors::link::routed`] adapts an accept/connect transport to the link -//! contract, mapping every link stream to its own connection so that flow -//! control and half-close are the transport's own. This module supplies the -//! sprockets end of that adapter: a dialer, a listener, and the connection -//! type they exchange. Peers are authenticated and attested by sprockets, -//! which is what rumors asks of a transport it trusts. -//! -//! A fresh connection costs a full attested handshake, most of a second -//! with the RoT in the loop, and the RoT serializes them. The dialer -//! therefore pools. rumors hands a completed stream's connection back -//! through [`Dial::recycle`], a qorb pool per peer holds it, and the next -//! stream to that peer draws it instead of dialing. Only clean returns -//! are reused: a connection dropped mid-stream is discarded at the -//! pool's door. - -use std::collections::{BTreeMap, BTreeSet, HashMap}; +//! Sprockets authenticates and attests each connection. Rumors owns reuse +//! within each link, avoiding repeated attestation for completed streams. +//! This adapter owns handshake timeouts and clean TLS shutdown. + +use std::collections::{BTreeMap, BTreeSet}; use std::io; use std::net::{Ipv6Addr, SocketAddr, SocketAddrV6}; use std::pin::Pin; -use std::sync::{Arc, Mutex}; +use std::sync::Arc; use std::task::{Context, Poll}; use std::time::Duration; -use async_trait::async_trait; use camino::Utf8PathBuf; -use qorb::backend::{self, Backend}; -use qorb::claim; -use qorb::policy::{Policy, SetConfig}; -use qorb::pool::Pool; -use qorb::resolvers::fixed::FixedResolver; -use rumors::link::STREAM_COUNT; use rumors::link::routed::{Config, Dial, Endpoint, Incoming, Listen, RoutedLink}; use sled_hardware_types::BaseboardId; use slog::{Logger, debug, o, warn}; use sprockets_tls::keys::SprocketsConfig; use sprockets_tls::{Client, Server}; -use tokio::io::{AsyncRead, AsyncReadExt as _, AsyncWrite, AsyncWriteExt as _, ReadBuf}; +use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt as _, ReadBuf}; use tokio::net::TcpStream; use tokio::runtime::Handle; use tokio::sync::{mpsc, watch}; @@ -57,92 +39,21 @@ pub type CorpusSource = Arc Vec + Send + Sync>; /// Completed handshakes the listener holds while the router catches up. const HANDSHAKE_QUEUE_DEPTH: usize = 64; -/// Inbound connections the router holds mid-header. Recycled idle -/// connections park here between streams, so the bound admits a whole -/// rack's stream complements at once. -const PENDING_HEADERS: usize = STREAM_COUNT * 32; - /// How long to wait after a failed accept before accepting again. const ACCEPT_RETRY: Duration = Duration::from_millis(100); -/// Connections per peer. The link contract's worst case is one control -/// stream plus a full complement for each of a pair's two links, and -/// claims parked on the ready byte can briefly overlap the next -/// session's opens. Slots are created on demand, so the cap is free -/// when idle. -const MAX_SLOTS: usize = 3 * STREAM_COUNT + 1; - -/// Idle connections kept warm per peer. Taking one stream leaves a -/// spare, so qorb fires no refill and reuses the recycled connection -/// rather than culling it. -const SPARES_WANTED: usize = 2; - -/// Floor between one slot's reconnect attempts after a failure, -/// matched to the RoT's roughly one-per-second handshake rate. -const MIN_CONNECTION_BACKOFF: Duration = Duration::from_secs(1); - -/// How often qorb re-checks an idle connection. The check is a no-op -/// (no ping), so the tick only bounces slot state. It must be finite: -/// qorb arms the timer by adding this to the current instant, and -/// `Duration::MAX`, which its docs suggest for disabling checks, -/// overflows the addition and kills the slot task. -const HEALTH_INTERVAL: Duration = Duration::from_secs(24 * 60 * 60); - -/// One attested connection as a peer's pool holds it. -struct PooledConn { - stream: Option>, - /// qorb returns a dropped claim to the pool no matter how its - /// stream ended, so this bit tells a clean return from an abort: - /// set as the connection is claimed, cleared only by - /// [`Dial::recycle`], and checked at the pool's door, where a dirty - /// connection is discarded. - dirty: bool, -} - -impl Drop for PooledConn { - fn drop(&mut self) { - let Some(mut stream) = self.stream.take() else { - return; - }; - // A dropped TLS connection sends a bare FIN, which rustls - // reports to the peer as an unexpected end of file, so the drop - // hands the connection to a task that shuts it down properly. - // Outside a runtime there is no one left to read the closing - // alert either, so let the connection drop abruptly. - if let Ok(handle) = Handle::try_current() { - handle.spawn(async move { - let _ = stream.shutdown().await; - }); - } - } -} - -/// One sprockets connection, carrying one link stream at a time. -/// -/// A dialed connection is a claim on its peer's pool: dropping it -/// returns it, and the pool reuses it only if [`Dial::recycle`] marked -/// it clean first. An accepted connection belongs to the router, and -/// dropping it closes it. -pub struct SprocketsConn(Conn); - -enum Conn { - /// Accepted by the listener. - Direct(Option>), - /// Claimed from a peer's pool. - Pooled(claim::Handle), -} +/// An attested connection, closed with a TLS shutdown on drop. +pub struct SprocketsConn(Option>); impl SprocketsConn { + /// Borrow the live TLS stream for I/O. fn stream(&mut self) -> Pin<&mut sprockets_tls::Stream> { - let stream = match &mut self.0 { - Conn::Direct(stream) => stream.as_mut(), - Conn::Pooled(handle) => handle.stream.as_mut(), - }; - Pin::new(stream.expect("stream present until drop")) + Pin::new(self.0.as_mut().expect("stream present until drop")) } } impl AsyncRead for SprocketsConn { + /// Read decrypted bytes from the TLS stream. fn poll_read( mut self: Pin<&mut Self>, cx: &mut Context<'_>, @@ -153,6 +64,7 @@ impl AsyncRead for SprocketsConn { } impl AsyncWrite for SprocketsConn { + /// Write bytes through the TLS stream. fn poll_write( mut self: Pin<&mut Self>, cx: &mut Context<'_>, @@ -161,26 +73,24 @@ impl AsyncWrite for SprocketsConn { self.stream().poll_write(cx, buf) } + /// Flush pending TLS output. fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { self.stream().poll_flush(cx) } + /// Send TLS shutdown after pending output. fn poll_shutdown(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { self.stream().poll_shutdown(cx) } } impl Drop for SprocketsConn { + /// Finish TLS shutdown in a task when a runtime is available. fn drop(&mut self) { - // A pooled connection's drop is its return to the pool, which - // discards it (shutting it down) unless it was recycled clean. - let Conn::Direct(stream) = &mut self.0 else { + let Some(mut stream) = self.0.take() else { return; }; - let Some(mut stream) = stream.take() else { - return; - }; - // As for `PooledConn`: shut down properly inside a runtime. + // Send close_notify so the peer sees EOF after accepted bytes. if let Ok(handle) = Handle::try_current() { handle.spawn(async move { let _ = stream.shutdown().await; @@ -293,74 +203,7 @@ fn baseboard(log: &Logger, platform_id: &str) -> Option { } } -/// Establishes a peer pool's connections, one attested handshake each. -struct PoolConnector { - log: Logger, - config: SprocketsConfig, - corpus: CorpusSource, - baseboards: Baseboards, - timeout: Duration, -} - -#[async_trait] -impl backend::Connector for PoolConnector { - type Connection = PooledConn; - - async fn connect(&self, backend: &Backend) -> Result { - // The resolver only ever names IPv6 backends; see `Dial::dial`. - let SocketAddr::V6(addr) = backend.address else { - return Err(backend::Error::from(io::Error::other(format!( - "sush gossip needs IPv6: {}", - backend.address - )))); - }; - let config = self.config.clone(); - let corpus = self.corpus.clone(); - let log = self.log.clone(); - // Sprockets connects are not cancel safe, and qorb may cancel - // this future, so the handshake runs in its own task and this - // future only awaits the result. The deadline lives inside the - // task, so an abandoned handshake still terminates at it - // instead of holding the RoT unbounded. The corpus source is - // consulted inside the task: if it panics, this connect fails, - // not the pool. - let deadline = self.timeout; - let dial = spawn(async move { - match timeout(deadline, Client::connect(config, addr, (corpus)(), log)).await { - Ok(connected) => connected.map_err(io::Error::other), - Err(_) => Err(io::Error::new( - io::ErrorKind::TimedOut, - format!("dialing {addr} timed out"), - )), - } - }); - let stream = dial.await.map_err(io::Error::other)??; - self.baseboards - .dialed(&self.log, addr, stream.peer_platform_id().as_str()); - Ok(PooledConn { - stream: Some(stream), - dirty: false, - }) - } - - async fn on_acquire(&self, conn: &mut PooledConn) -> Result<(), backend::Error> { - // Pessimistic: only a clean recycle clears it. - conn.dirty = true; - Ok(()) - } - - async fn on_recycle(&self, conn: &mut PooledConn) -> Result<(), backend::Error> { - if conn.dirty { - return Err(backend::Error::from(io::Error::new( - io::ErrorKind::ConnectionAborted, - "connection was dropped mid-stream", - ))); - } - Ok(()) - } -} - -/// Dials attested sprockets connections, pooled per peer. +/// Opens fresh attested connections with a handshake timeout. #[derive(Clone)] pub struct SprocketsDial { log: Logger, @@ -368,13 +211,10 @@ pub struct SprocketsDial { corpus: CorpusSource, baseboards: Baseboards, timeout: Duration, - pools: Arc>>>>, } impl SprocketsDial { - /// Dial with `config`, appraising peers against the corpus of the - /// moment, and giving up on any one handshake or claim after - /// `timeout`. + /// Appraise each peer against the current corpus and bound its handshake. pub fn new( log: &Logger, config: SprocketsConfig, @@ -382,90 +222,46 @@ impl SprocketsDial { baseboards: Baseboards, timeout: Duration, ) -> Self { - SprocketsDial { + Self { log: log.new(o!("component" => "sprockets dial")), config, corpus, baseboards, timeout, - pools: Arc::default(), - } - } - - /// The pool dialing `addr`, created on first use. - fn pool(&self, addr: SocketAddrV6) -> Arc> { - self.pools - .lock() - .expect("pool table lock") - .entry(addr) - .or_insert_with(|| Arc::new(self.build(addr))) - .clone() - } - - fn build(&self, addr: SocketAddrV6) -> Pool { - let connector = Arc::new(PoolConnector { - log: self.log.clone(), - config: self.config.clone(), - corpus: self.corpus.clone(), - baseboards: self.baseboards.clone(), - timeout: self.timeout, - }); - let resolver = Box::new(FixedResolver::new([SocketAddr::V6(addr)])); - let policy = Policy { - spares_wanted: SPARES_WANTED, - max_slots: MAX_SLOTS, - claim_timeout: self.timeout, - set_config: SetConfig { - max_count: MAX_SLOTS, - min_connection_backoff: MIN_CONNECTION_BACKOFF, - health_interval: HEALTH_INTERVAL, - ..SetConfig::default() - }, - ..Policy::default() - }; - match Pool::new(format!("gossip {addr}"), resolver, connector, policy) { - Ok(pool) => pool, - Err(err) => err.into_inner(), } } - - /// Drop the pools of peers outside `peers`, closing their idle - /// connections. - fn retain(&self, peers: &BTreeSet) { - self.pools - .lock() - .expect("pool table lock") - .retain(|addr, _| peers.contains(addr)); - } } impl Dial for SprocketsDial { + /// The peer's advertised listen address. type Addr = SocketAddr; + /// An attested TLS connection. type Conn = SprocketsConn; + /// Open and attest a fresh connection within the handshake timeout. async fn dial(&self, addr: &SocketAddr) -> io::Result { let SocketAddr::V6(addr) = *addr else { return Err(io::Error::other(format!("sush gossip needs IPv6: {addr}"))); }; - let handle = - self.pool(addr).claim().await.map_err(|err| { - io::Error::other(format!("claiming a connection to {addr}: {err}")) - })?; - Ok(SprocketsConn(Conn::Pooled(handle))) - } - - fn recycle(&self, _peer: &SocketAddr, mut conn: SprocketsConn) { - // Reusable only once the peer's router says it is ready: read - // that byte off the session's task, then let the drop return - // the claim clean. - spawn(async move { - let mut ready = [0u8; 1]; - if conn.read_exact(&mut ready).await.is_ok() - && let Conn::Pooled(handle) = &mut conn.0 - { - handle.dirty = false; + let config = self.config.clone(); + let corpus = self.corpus.clone(); + let log = self.log.clone(); + let deadline = self.timeout; + // Sprockets handshakes are not cancellation-safe. Keep the timeout + // inside the spawned task so an abandoned dial still terminates. + let dial = spawn(async move { + match timeout(deadline, Client::connect(config, addr, (corpus)(), log)).await { + Ok(connected) => connected.map_err(io::Error::other), + Err(_) => Err(io::Error::new( + io::ErrorKind::TimedOut, + format!("dialing {addr} timed out"), + )), } }); + let stream = dial.await.map_err(io::Error::other)??; + self.baseboards + .dialed(&self.log, addr, stream.peer_platform_id().as_str()); + Ok(SprocketsConn(Some(stream))) } } @@ -474,7 +270,6 @@ impl Dial for SprocketsDial { pub struct Transport { endpoint: Endpoint, incoming: Incoming, - dial: SprocketsDial, baseboards: Baseboards, bound: SocketAddrV6, } @@ -482,6 +277,8 @@ pub struct Transport { impl Transport { /// Listen on `listen_addr` and stand up the routing endpoint. Its /// router runs until `shutdown`, or until the listener fails. + /// `dial_timeout` bounds outgoing handshakes and incoming initial routing; + /// it does not limit idle gossip links. pub async fn new( log: &Logger, config: SprocketsConfig, @@ -497,16 +294,13 @@ impl Transport { corpus.clone(), baseboards.clone(), listen_addr, + dial_timeout, shutdown.clone(), ) .await?; let dial = SprocketsDial::new(log, config, corpus, baseboards.clone(), dial_timeout); - let router_config = Config { - pending_headers: PENDING_HEADERS, - ..Config::default() - }; let (endpoint, incoming, router) = - Endpoint::new(listen, SocketAddr::V6(bound), dial.clone(), router_config) + Endpoint::new(listen, SocketAddr::V6(bound), dial, Config::default()) .map_err(io::Error::other)?; let log = log.new(o!("component" => "link router")); spawn(async move { @@ -520,7 +314,6 @@ impl Transport { Ok(Transport { endpoint, incoming, - dial, baseboards, bound, }) @@ -541,10 +334,8 @@ impl Transport { &self.baseboards } - /// Drop the connection pools and recorded baseboards of peers - /// outside `peers`. + /// Forget recorded baseboards of peers outside `peers`. pub fn retain_peers(&self, peers: &BTreeSet) { - self.dial.retain(peers); self.baseboards.retain(peers); } @@ -564,7 +355,10 @@ impl Transport { /// concurrently, and [`Listen::accept`] is a queue receive, which the /// router may cancel freely. pub struct SprocketsListen { + /// Attested connections waiting for the router. connections: mpsc::Receiver, + /// Initial routing uses the same timeout as outgoing dials. + routing_timeout: Duration, } impl SprocketsListen { @@ -576,6 +370,7 @@ impl SprocketsListen { corpus: CorpusSource, baseboards: Baseboards, listen_addr: SocketAddrV6, + routing_timeout: Duration, shutdown: CancellationToken, ) -> io::Result<(Self, SocketAddrV6)> { let log = log.new(o!("component" => "sprockets listen")); @@ -590,13 +385,26 @@ impl SprocketsListen { }; let (tx, connections) = mpsc::channel(HANDSHAKE_QUEUE_DEPTH); spawn(pump(server, corpus, baseboards, tx, log, shutdown)); - Ok((SprocketsListen { connections }, bound)) + Ok(( + SprocketsListen { + connections, + routing_timeout, + }, + bound, + )) } } impl Listen for SprocketsListen { + /// An attested TLS connection. type Conn = SprocketsConn; + /// Bound initial routing without timing out established gossip. + fn routing_deadline(&self) -> impl Future + Send + 'static { + sleep(self.routing_timeout) + } + + /// Receive the next completed handshake or report listener shutdown. async fn accept(&mut self) -> io::Result { self.connections .recv() @@ -630,7 +438,7 @@ async fn pump( let id = stream.peer_platform_id().as_str(); baseboards.accepted(&log, *peer.ip(), id); } - let conn = SprocketsConn(Conn::Direct(Some(stream))); + let conn = SprocketsConn(Some(stream)); let _ = connections.send(conn).await; } Err(err) => { diff --git a/server/src/messages.rs b/server/src/messages.rs index b6885ee1..e3c68305 100644 --- a/server/src/messages.rs +++ b/server/src/messages.rs @@ -269,8 +269,8 @@ pub mod v0 { pub enum Error { #[error( "Concurrent sessions detected: \ - ours is {own_session}@{own_version}, \ - incoming is {incoming_session}@{incoming_version}" + ours is {own_session}@{own_version:?}, \ + incoming is {incoming_session}@{incoming_version:?}" )] ConcurrentSessions { own_session: SessionId, @@ -317,6 +317,11 @@ mod wire_format { } let decoded: VersionedMessage = from_cbor(bytes.as_slice()).unwrap(); assert_eq!(decoded, message, "wire format should round-trip"); + // Exercise Rumors' admission checks as well as the wire snapshot. + rumors::Peer::::seed() + .into_rumors() + .send(message) + .expect("locally authored Sush messages must encode faithfully"); } fn sid(name: &str) -> SessionId { diff --git a/server/src/state.rs b/server/src/state.rs index 18ebe8b8..1003253d 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -921,8 +921,8 @@ impl State { log, "stale session start"; "session_id" => %session_id, - "frontier" => %frontier, - "incoming_version" => %incoming_version, + "frontier" => ?frontier, + "incoming_version" => ?incoming_version, ); *frontier |= incoming_version.clone(); } @@ -2152,16 +2152,25 @@ mod test { ) } + /// Versions before, at, and after one event, plus a concurrent event. + fn session_versions() -> [Version; 4] { + let mut local = rumors::before::Clock::seed(); + let before = local.tick().clone(); + let mut remote = local.fork(); + let at = local.tick().clone(); + let after = local.tick().clone(); + // Both forks share the first event; their later events are independent. + let concurrent = remote.tick().clone(); + [before, at, after, concurrent] + } + /// The admission rules: a sled with no record admits, a record /// from a foreign universe admits, the committed session admits /// outright, a session started strictly above the executed join /// admits, and everything else hops. #[test] fn admission_rules() { - let started: Version = "(1, 1, (0, 0, 2))".parse().unwrap(); - let older: Version = "(1, 0, (0, 0, 2))".parse().unwrap(); - let newer: Version = "(2, 1, (0, 0, 3))".parse().unwrap(); - let concurrent: Version = "(1, 2, (0, 0, 1))".parse().unwrap(); + let [older, started, newer, concurrent] = session_versions(); let session = SessionId::random(); let job = JobId::random(); let committed = Boundary { @@ -2211,11 +2220,8 @@ mod test { /// record. #[test] fn floors_refuse_below() { - let floor: Version = "(1, 1, (0, 0, 2))".parse().unwrap(); - let at: Version = floor.clone(); - let below: Version = "(1, 0, (0, 0, 2))".parse().unwrap(); - let above: Version = "(2, 1, (0, 0, 3))".parse().unwrap(); - let concurrent: Version = "(1, 2, (0, 0, 1))".parse().unwrap(); + let [below, floor, above, concurrent] = session_versions(); + let at = floor.clone(); let past = past(None, Some(floor.clone())); for started in [&at, &below, &concurrent] { diff --git a/server/tests/distributed.rs b/server/tests/distributed.rs index da133f9b..11eeb458 100644 --- a/server/tests/distributed.rs +++ b/server/tests/distributed.rs @@ -2,7 +2,10 @@ // License, v. 2.0. If a copy of the MPL was not distributed with this // file, You can obtain one at https://mozilla.org/MPL/2.0/. -//! Two job managers on two gossiping sleds. What one accepts, both know. +//! Job managers on gossiping sleds. What one accepts, the others learn. +//! +//! These tests use a small multithreaded runtime, like the server, so +//! synchronous attestation work can progress in parallel within session deadlines. mod common; @@ -117,7 +120,7 @@ impl Sled { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn jobs_gossip_between_sleds() { let (_tmp, dir) = pki("sush-distributed-", 2); let mut root = common::ephemeral_root(); @@ -229,7 +232,7 @@ async fn jobs_gossip_between_sleds() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn rejoining_replays_without_reexecuting() { let (_tmp, dir) = pki("sush-replay-", 2); let mut root = common::ephemeral_root(); @@ -346,7 +349,7 @@ async fn rejoining_replays_without_reexecuting() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn interrupted_jobs_get_stopped() { let (_tmp, dir) = pki("sush-interrupted-", 2); let mut root = common::ephemeral_root(); @@ -436,7 +439,7 @@ async fn interrupted_jobs_get_stopped() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn stragglers_do_not_interrupt_live_jobs() { let (_tmp, dir) = pki("sush-straggler-", 3); let mut root = common::ephemeral_root(); @@ -583,7 +586,7 @@ async fn sign_job_for( } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn lost_suffix_never_reruns() { let (_tmp, dir) = pki("sush-lost-", 2); let mut root = common::ephemeral_root(); @@ -765,7 +768,7 @@ async fn lost_suffix_never_reruns() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn session_resumes_at_stored_successor() { let (_tmp, dir) = pki("sush-resume-", 3); let mut root = common::ephemeral_root(); @@ -922,7 +925,7 @@ async fn session_resumes_at_stored_successor() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn universe_flip_flop_raises_floor() { let (_tmp, dir) = pki("sush-flipflop-", 3); let mut root = common::ephemeral_root(); @@ -1173,7 +1176,7 @@ async fn universe_flip_flop_raises_floor() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn witnessed_session_survives_restart() { let (_tmp, dir) = pki("sush-witness-", 2); let mut root = common::ephemeral_root(); @@ -1264,7 +1267,7 @@ async fn witnessed_session_survives_restart() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn bookmarks_survive_restart() { let (_tmp, dir) = pki("sush-bookmark-", 2); let mut root = common::ephemeral_root(); @@ -1349,7 +1352,7 @@ async fn bookmarks_survive_restart() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn gossip_survives_bookmark_failure() { let (_tmp, dir) = pki("sush-nobookmark-", 2); let mut root = common::ephemeral_root(); diff --git a/server/tests/gossip.rs b/server/tests/gossip.rs index 2ea2c3d4..e0822b6f 100644 --- a/server/tests/gossip.rs +++ b/server/tests/gossip.rs @@ -3,7 +3,9 @@ // file, You can obtain one at https://mozilla.org/MPL/2.0/. //! The gossip manager converging localhost meshes, with real attested -//! handshakes throughout. +//! handshakes throughout. Network tests use a small multithreaded runtime, +//! like the server, so synchronous attestation work does not serialize every +//! handshake onto one runtime thread. mod common; @@ -106,7 +108,7 @@ fn converged(nodes: &[&Node]) -> Option { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn cold_start_converges() { let (_tmp, dir) = pki("sush-gossip-", 3); let log = test_logger(function_name!()); @@ -130,7 +132,7 @@ async fn cold_start_converges() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn staggered_start_converges() { let (_tmp, dir) = pki("sush-gossip-", 3); let log = test_logger(function_name!()); @@ -159,7 +161,7 @@ async fn staggered_start_converges() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn node_replacement_reconverges() { let (_tmp, dir) = pki("sush-gossip-", 4); let log = test_logger(function_name!()); @@ -189,7 +191,7 @@ async fn node_replacement_reconverges() { } #[named] -#[tokio::test] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn linked_follows_live_links() { let (_tmp, dir) = pki("sush-gossip-", 2); let log = test_logger(function_name!()); @@ -208,3 +210,70 @@ async fn linked_follows_live_links() { drop(b); eventually("dead peer unlinked", 120, async || a.linked().is_empty()).await; } + +/// Idle links probe every ten seconds without delaying changes or bursting +/// after a pause; a silent peer then trips the one-second session deadline. +#[tokio::test(start_paused = true)] +async fn idle_heartbeats_detect_a_silent_peer() { + use std::time::Duration; + + use futures::{FutureExt as _, StreamExt as _}; + use rumors::{Error, Joined, Led, Peer}; + use tokio::time::advance; + + let a: Rumors = Seed::grow( + &test_logger("idle_heartbeats_detect_a_silent_peer"), + &Locker::null(), + ) + .await + .into_rumors(); + let (mut near, mut far) = rumors::link::memory(); + let (served, joined) = futures::join!( + a.gossip_once(&mut near), + Peer::::bootstrap().join(&mut far), + ); + served.unwrap(); + let Joined::Joined { peer } = joined else { + panic!("join succeeds") + }; + let b = peer.into_rumors(); + let mut sessions = a.gossip(&mut near); + let (sent, received) = futures::join!(sessions.next(), b.gossip_once(&mut far)); + assert_eq!(sent.unwrap().unwrap().led, Led::Local); + received.unwrap(); + + advance(Duration::from_secs(9)).await; + assert!(sessions.next().now_or_never().is_none()); + a.send("a change before the heartbeat".into()).unwrap(); + let (sent, received) = futures::join!(sessions.next(), b.gossip_once(&mut far)); + sent.unwrap().unwrap(); + received.unwrap(); + assert_eq!(a.snapshot(), b.snapshot()); + + // The heartbeat still fires at ten seconds, even though the set is current. + advance(Duration::from_secs(1)).await; + assert!(sessions.next().now_or_never().is_none()); + let (sent, received) = futures::join!(sessions.next(), b.gossip_once(&mut far)); + assert_eq!(sent.unwrap().unwrap().led, Led::Local); + received.unwrap(); + + // Several missed intervals produce one probe, followed by a full interval. + advance(Duration::from_secs(35)).await; + assert!(sessions.next().now_or_never().is_none()); + let (sent, received) = futures::join!(sessions.next(), b.gossip_once(&mut far)); + assert_eq!(sent.unwrap().unwrap().led, Led::Local); + received.unwrap(); + assert!(sessions.next().now_or_never().is_none()); + advance(Duration::from_secs(9)).await; + assert!(sessions.next().now_or_never().is_none()); + + // Keep the remote link open but stop serving it: no EOF can expose failure. + advance(Duration::from_secs(1)).await; + assert!(sessions.next().now_or_never().is_none()); + advance(Duration::from_secs(1)).await; + assert!(matches!( + sessions.next().await, + Some(Err(Error::DeadlineExceeded)) + )); + assert!(sessions.next().await.is_none()); +} diff --git a/server/tests/link.rs b/server/tests/link.rs index 9cbdfb54..d4c274b6 100644 --- a/server/tests/link.rs +++ b/server/tests/link.rs @@ -2,8 +2,10 @@ // License, v. 2.0. If a copy of the MPL was not distributed with this // file, You can obtain one at https://mozilla.org/MPL/2.0/. -//! The sprockets transport against the rumors link contract, and a -//! two-peer gossip smoke test over it. +//! Attested transport conformance and two-peer gossip. +//! +//! A small multithreaded runtime, like the server, lets synchronous +//! attestation work progress in parallel within the production deadlines. mod common; @@ -14,7 +16,7 @@ use futures::StreamExt as _; use futures::stream; use slog::Logger; use tempfile::TempDir; -use tokio::time::timeout; +use tokio::time::{sleep, timeout}; use tokio::{join, spawn}; use tokio_util::sync::CancellationToken; @@ -24,13 +26,19 @@ use sush_server::link::{SprocketsLink, Transport}; use common::{corpus, dial_timeout, localhost, pki, sprockets_config, test_logger}; +/// Two attested endpoints and the PKI used by their handshakes. struct TestNet { + /// Endpoint for identity 1. a: Transport, + /// Endpoint for identity 2. b: Transport, + /// Cancels both endpoints' transport tasks. shutdown: CancellationToken, + /// Keeps certificate and attestation files alive during the test. _dir: TempDir, } +/// Bind one identity's authenticated transport for the test network. async fn transport( log: &Logger, dir: &Utf8PathBuf, @@ -50,6 +58,7 @@ async fn transport( } impl TestNet { + /// Create test PKI and start both attested endpoints. async fn new(test_name: &'static str) -> TestNet { let (tmp, dir) = pki("sush-link-", 2); let log = test_logger(test_name); @@ -69,7 +78,8 @@ impl TestNet { let endpoint = self.a.endpoint(); let peer = *self.b.endpoint().local_addr(); let (linked, accepted) = join!(endpoint.link(peer), self.b.accept()); - let link_a = linked.expect("peer router accepts the link"); + let (info, link_a) = linked.expect("peer router accepts the link"); + assert_eq!(info.peer, peer); let (from, link_b) = accepted.expect("router is live"); assert_eq!(from, *endpoint.local_addr()); (link_a, link_b) @@ -77,29 +87,31 @@ impl TestNet { } impl Drop for TestNet { + /// Cancel the network's transport tasks. fn drop(&mut self) { self.shutdown.cancel(); } } -// Slow (~35s each): CI always runs these, and local runs should -// whenever `link.rs` or the gossip configuration changes: +// These tests perform real attestation. CI runs them; run them locally +// when changing the transport or gossip configuration: // // cargo test --package sush-server --test link -- --include-ignored -#[tokio::test] +/// The attested transport satisfies the Rumors link contract. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[ignore] async fn conformance() { let mut net = TestNet::new("conformance").await; - timeout( - Duration::from_secs(600), - check(async || net.link_pair().await), + check( + async || net.link_pair().await, + || sleep(Duration::from_secs(600)), ) - .await - .expect("conformance suite timed out"); + .await; } -#[tokio::test] +/// Peers bootstrap and exchange messages over an attested link. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[ignore] async fn gossip_convergence() { let mut net = TestNet::new("gossip_convergence").await; @@ -107,13 +119,16 @@ async fn gossip_convergence() { // Alice seeds a universe with one message and serves sessions on her // end of the link. - let alice: Rumors = Peer::seed().into_rumors(); + let alice: Rumors = Peer::seed() + .gossip_when(|_| stream::pending::<()>()) + .session_deadline(|| sleep(Duration::from_secs(1))) + .into_rumors(); alice.send("from alice".to_string()).unwrap(); let server = spawn({ let alice = alice.clone(); async move { let mut link_a = link_a; - let mut driver = alice.gossip_when(stream::pending::<()>(), &mut link_a); + let mut driver = alice.gossip(&mut link_a); while let Some(session) = driver.next().await { session.expect("serving gossip session"); } @@ -121,21 +136,23 @@ async fn gossip_convergence() { }); // Bob joins Alice's universe through the link and hears her message. - let bob = timeout( + let rumors::Joined::Joined { peer: bob } = timeout( Duration::from_secs(60), - Peer::::bootstrap().join(&mut link_b), + Peer::::bootstrap() + .session_deadline(|| sleep(Duration::from_secs(1))) + .join(&mut link_b), ) .await - .expect("bootstrap timed out") - .expect("bootstrap failed") - .expect("mutual bootstrap bail") - .into_rumors(); + .expect("bootstrap timed out") else { + panic!("Alice must serve Bob's bootstrap"); + }; + let bob = bob.into_rumors(); assert_eq!(bob.network(), alice.network()); assert_eq!(bob.snapshot().len(), 1); // Bob's own message reaches Alice within one gossip session. bob.send("from bob".to_string()).unwrap(); - timeout(Duration::from_secs(60), bob.gossip(&mut link_b)) + timeout(Duration::from_secs(60), bob.gossip_once(&mut link_b)) .await .expect("gossip timed out") .expect("gossip failed");