diff --git a/Cargo.lock b/Cargo.lock index 280e23c..7e7467b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1330,7 +1330,7 @@ dependencies = [ [[package]] name = "chat-proto" version = "0.1.0" -source = "git+https://github.com/logos-messaging/chat_proto?rev=37ec98a151f6d50aab2905802ac0a896477e62ea#37ec98a151f6d50aab2905802ac0a896477e62ea" +source = "git+https://github.com/logos-messaging/chat_proto?rev=948f4ad6dedb01b44e279204d8eb30a1eda5a330#948f4ad6dedb01b44e279204d8eb30a1eda5a330" dependencies = [ "prost", ] @@ -1463,14 +1463,15 @@ name = "components" version = "0.1.0" dependencies = [ "base64", + "chat-proto", "crossbeam-channel", "crypto", "hex", "libchat", "logos-account", + "prost", "reqwest 0.12.28", "serde", - "serde_json", "storage", "thiserror", "tracing", @@ -3629,8 +3630,6 @@ version = "0.1.0" dependencies = [ "base64", "crossbeam-channel", - "libchat", - "logos-generic-chat", "serde", "serde_json", "thiserror", diff --git a/Cargo.toml b/Cargo.toml index 1efb279..47d84ea 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -44,6 +44,7 @@ logos-delivery = { path = "extensions/logos-delivery-rust"} logos-generic-chat = { path = "crates/generic-chat" } shared-traits = { path = "core/shared-traits" } storage = { path = "core/storage" } +chat-proto = { git = "https://github.com/logos-messaging/chat_proto", rev = "948f4ad6dedb01b44e279204d8eb30a1eda5a330" } # External Workspace dependency declarations (sorted) blake2 = "0.10" diff --git a/bin/chat-cli/src/main.rs b/bin/chat-cli/src/main.rs index 4a51208..bf110a3 100644 --- a/bin/chat-cli/src/main.rs +++ b/bin/chat-cli/src/main.rs @@ -10,7 +10,7 @@ use clap::{Parser, ValueEnum}; use crossbeam_channel::Receiver; use logos_chat::{ AccountDirectory, ChatClient, ChatStore, Event, LogosConfig, P2pConfig, RegistrationService, - Transport, + RegistryPublishMode, Transport, }; use app::ChatApp; @@ -65,6 +65,28 @@ struct Cli { /// Example: `--registry-url http://127.0.0.1:18080`. #[arg(long)] registry_url: Option, + + /// How keypackage and account bundles are submitted to the store: over its + /// HTTP POST API, or published on the delivery network for the store to + /// pick up by subscription. Queries always use the HTTP API. + #[arg(long, value_enum, default_value_t = RegistryPublishKind::Http)] + registry_publish: RegistryPublishKind, +} + +#[derive(Copy, Clone, Debug, ValueEnum)] +#[value(rename_all = "kebab-case")] +enum RegistryPublishKind { + Http, + Delivery, +} + +impl From for RegistryPublishMode { + fn from(kind: RegistryPublishKind) -> Self { + match kind { + RegistryPublishKind::Http => RegistryPublishMode::Http, + RegistryPublishKind::Delivery => RegistryPublishMode::Delivery, + } + } } fn main() -> Result<()> { @@ -95,6 +117,7 @@ fn main() -> Result<()> { if let Some(registry_url) = cli.registry_url.as_deref() { config.set_registry_url(registry_url); } + config.set_registry_publish_mode(cli.registry_publish.into()); config.set_p2p_config(p2p_config); let (client, events) = logos_chat::open(config) .map_err(|e| anyhow::anyhow!("{e:?}")) @@ -115,6 +138,7 @@ fn main() -> Result<()> { if let Some(registry_url) = cli.registry_url.as_deref() { config.set_registry_url(registry_url); } + config.set_registry_publish_mode(cli.registry_publish.into()); let (client, events) = logos_chat::open_with_transport(config, transport) .map_err(|e| anyhow::anyhow!("{e:?}")) .context("failed to open chat client")?; diff --git a/bin/chat-cli/src/transport/file.rs b/bin/chat-cli/src/transport/file.rs index 75b1d31..eb17215 100644 --- a/bin/chat-cli/src/transport/file.rs +++ b/bin/chat-cli/src/transport/file.rs @@ -14,7 +14,7 @@ pub enum FileTransportError { Io(#[from] io::Error), } -#[derive(Debug)] +#[derive(Clone, Debug)] pub struct FileTransport { transport_dir: PathBuf, inbound_rx: Option>>, diff --git a/core/conversations/Cargo.toml b/core/conversations/Cargo.toml index 0b3950d..e44b1e4 100644 --- a/core/conversations/Cargo.toml +++ b/core/conversations/Cargo.toml @@ -17,7 +17,7 @@ storage = { workspace = true } # External dependencies (sorted) alloy = "2.0" base64 = "0.22" -chat-proto = { git = "https://github.com/logos-messaging/chat_proto", rev = "37ec98a151f6d50aab2905802ac0a896477e62ea" } +chat-proto = { workspace = true } de-mls = { git = "https://github.com/vacp2p/de-mls", rev = "2c7a8669c1492c749c02efd2c5ac45e93e4926a3"} # Expose Mls Extensions (#131) double-ratchets = { path = "../double-ratchets" } hashgraph-like-consensus = "0.6.0" diff --git a/crates/generic-chat/src/lib.rs b/crates/generic-chat/src/lib.rs index f3c42d3..91d1318 100644 --- a/crates/generic-chat/src/lib.rs +++ b/crates/generic-chat/src/lib.rs @@ -23,4 +23,6 @@ pub use logos_account::AccountDirectory; // Re-export bundled registry implementations so callers can pick one without // pulling in `components` directly. -pub use components::{EphemeralRegistry, HttpRegistry, HttpRegistryError}; +pub use components::{ + ContactRegistry, ContactRegistryError, EphemeralRegistry, RegistryPublishMode, +}; diff --git a/crates/logos-chat/src/logos.rs b/crates/logos-chat/src/logos.rs index a5d87ca..beebbdd 100644 --- a/crates/logos-chat/src/logos.rs +++ b/crates/logos-chat/src/logos.rs @@ -2,7 +2,8 @@ //! //! [`open`] commits to the Logos service stack so independently built clients //! share the same production services instead of each re-deriving them: a -//! delegate identity, the HTTP keypackage + account registry, and encrypted +//! delegate identity, the keypackage + account registry (queried over HTTP, +//! with submissions over HTTP or the delivery network), and encrypted //! on-disk storage. The stack is generic over the transport — any //! [`Transport`] can be injected via [`open_with_transport`] — and the //! concrete [`LogosChatClient`] commits to the embedded logos-delivery node, @@ -14,7 +15,7 @@ //! constructors off the alias; they are crate-level functions ([`open`], //! [`open_with_transport`]) taking the all-inclusive [`LogosConfig`] instead. -use components::HttpRegistry; +use components::{ContactRegistry, RegistryPublishMode}; use crossbeam_channel::Receiver; use embedded_logos_delivery::{EmbeddedLogosDelivery, P2pConfig}; use libchat::{ChatStorage, StorageConfig}; @@ -40,6 +41,7 @@ pub struct LogosConfig { db_path: String, db_key: String, registry_url: String, + registry_publish_mode: RegistryPublishMode, p2p_config: P2pConfig, group_v2_config: Option, } @@ -53,6 +55,7 @@ impl LogosConfig { db_path: db_path.into(), db_key: db_key.into(), registry_url: REGISTRY_ENDPOINT.to_string(), + registry_publish_mode: RegistryPublishMode::default(), p2p_config: P2pConfig::default(), group_v2_config: None, } @@ -64,6 +67,13 @@ impl LogosConfig { self.registry_url = registry_url.into(); } + /// Choose how keypackage and account bundles are submitted to the store: + /// HTTP POST (the default) or published over the delivery transport for the + /// store to pick up by subscription. Reads always use the HTTP query API. + pub fn set_registry_publish_mode(&mut self, mode: RegistryPublishMode) { + self.registry_publish_mode = mode; + } + /// Override the embedded node's p2p settings (defaults to /// [`P2pConfig::default`]). Only [`open`] starts an embedded node, so /// [`open_with_transport`] ignores this. @@ -106,18 +116,33 @@ pub fn open(config: LogosConfig) -> Result<(LogosChatClient, Receiver), C /// Open a client on the Logos stack per `config` with the injected transport, /// persisting to the encrypted database. +/// +/// The registry publishes per `config`'s +/// [`registry publish mode`](LogosConfig::set_registry_publish_mode): over +/// HTTP (the default), or over a clone of `transport` — sharing the client's +/// own delivery stack, which is why the transport must be `Clone`. #[allow(clippy::type_complexity)] -pub fn open_with_transport( +pub fn open_with_transport( config: LogosConfig, transport: T, -) -> Result<(ChatClient, Receiver), ClientError> { +) -> Result< + ( + ChatClient, ChatStorage>, + Receiver, + ), + ClientError, +> { // A fresh account endorsing a fresh delegate each open: the account // key is dropped after publishing the bundle, so devices cannot be // added later. A caller-supplied, custody-holding account replaces // this once the platform provides one. let account = TestLogosAccount::new(); let delegate = DelegateSigner::random(); - let mut registry = HttpRegistry::new(config.registry_url); + let mut registry = ContactRegistry::new( + transport.clone(), + config.registry_url, + config.registry_publish_mode, + ); account .add_delegate_signer(&mut registry, delegate.public_key()) .map_err(|e| ClientError::BundlePublish(e.to_string()))?; @@ -136,10 +161,12 @@ pub fn open_with_transport( } /// The Logos client: a [`ChatClient`] wired to the Logos service stack — a -/// [`DelegateSigner`] identity acting for a fresh dev account, the HTTP -/// keypackage + account registry ([`HttpRegistry`], which is both the -/// keypackage store and the account → device directory), and encrypted -/// [`ChatStorage`] — running an embedded logos-delivery node as its -/// transport. Open one with [`open`], or swap the transport via +/// [`DelegateSigner`] identity acting for a fresh dev account, the keypackage + +/// account registry ([`ContactRegistry`], which is both the keypackage store +/// and the account → device directory; it queries over HTTP and submits over +/// HTTP or the delivery network per [`LogosConfig::set_registry_publish_mode`]), +/// and encrypted [`ChatStorage`] — running an embedded logos-delivery node as +/// its transport. Open one with [`open`], or swap the transport via /// [`open_with_transport`]. -pub type LogosChatClient = ChatClient; +pub type LogosChatClient = + ChatClient, ChatStorage>; diff --git a/extensions/components/Cargo.toml b/extensions/components/Cargo.toml index 7e1a79f..c4d1dd5 100644 --- a/extensions/components/Cargo.toml +++ b/extensions/components/Cargo.toml @@ -12,10 +12,11 @@ storage = { workspace = true } # External dependencies (sorted) base64 = "0.22" +chat-proto = { workspace = true } crossbeam-channel = { workspace = true } hex = "0.4.3" +prost = "0.14.1" reqwest = { version = "0.12", default-features = false, features = ["blocking", "json", "rustls-tls"] } serde = { version = "1.0", features = ["derive"] } -serde_json = "1.0" thiserror = "2" tracing = "0.1" diff --git a/extensions/components/src/contact_registry.rs b/extensions/components/src/contact_registry.rs index 6d63e33..a9e9d92 100644 --- a/extensions/components/src/contact_registry.rs +++ b/extensions/components/src/contact_registry.rs @@ -1,2 +1,2 @@ pub mod ephemeral; -pub mod http; +pub mod store; diff --git a/extensions/components/src/contact_registry/http.rs b/extensions/components/src/contact_registry/http.rs deleted file mode 100644 index 0cf5987..0000000 --- a/extensions/components/src/contact_registry/http.rs +++ /dev/null @@ -1,431 +0,0 @@ -use std::fmt::Debug; -use std::time::{Duration, SystemTime, UNIX_EPOCH}; - -use base64::Engine; -use base64::engine::general_purpose::STANDARD as BASE64; -use crypto::{Ed25519Signature, Ed25519VerifyingKey}; -use libchat::{IdentityProvider, RegistrationService}; -use logos_account::{AccountDirectory, BundleError, DeviceSet, SignedDeviceBundle, verify_bundle}; -use serde::{Deserialize, Serialize}; - -/// HTTP client for the testnet KeyPackage Registry service. -/// -/// Throwaway transport for issue #110 — replaced by λLEZ in v0.3. -/// -/// The wire carries `device_id` (the hex device verifying key), an opaque -/// `payload` blob, and its `signature`. The signed bytes and the transmitted -/// `payload` bytes are identical, so every verifier checks the signature over -/// exactly what it received — no field-by-field reconstruction to keep in sync. -/// The `payload` is opaque to the server: it verifies `signature` over `payload` -/// with `device_id`'s key (proof-of-possession — only the holder of that key can -/// publish under `device_id`) without decoding the payload. -#[derive(Clone)] -pub struct HttpRegistry { - base_url: String, - http: reqwest::blocking::Client, -} - -#[derive(Debug, thiserror::Error)] -pub enum HttpRegistryError { - #[error("http: {0}")] - Http(#[from] reqwest::Error), - #[error("server returned status {0}: {1}")] - Server(u16, String), - #[error("decode: {0}")] - Decode(String), - #[error("clock before unix epoch")] - Clock, - #[error("signature verification failed")] - SignatureInvalid, - #[error("bundle: {0}")] - Bundle(#[from] BundleError), -} - -#[derive(Debug, Serialize)] -struct SubmitRequest { - /// hex of the 32-byte device verifying key — the verification + storage key. - device_id: String, - /// base64 of the canonical signed payload (see [`encode_payload`]). - payload: String, - /// base64 of the 64-byte Ed25519 signature over `payload`. - signature: String, -} - -#[derive(Debug, Deserialize)] -struct FetchResponse { - payload: String, - signature: String, -} - -#[derive(Debug, Serialize)] -struct SubmitAccountRequest { - /// hex of the 32-byte account verifying key — verification + storage key. - account_pub: String, - /// base64 of the canonical signed device-list payload. - payload: String, - /// base64 of the 64-byte account signature over `payload`. - signature: String, -} - -#[derive(Debug, Deserialize)] -struct FetchAccountResponse { - payload: String, - signature: String, - #[allow(dead_code)] // server's prune clock; freshness is taken from the bundle's lamport - updated_at: i64, -} - -impl HttpRegistry { - pub fn new(base_url: impl Into) -> Self { - Self::with_timeout(base_url, Duration::from_secs(10)) - } - - pub fn with_timeout(base_url: impl Into, timeout: Duration) -> Self { - let http = reqwest::blocking::Client::builder() - .timeout(timeout) - .build() - .expect("reqwest client builder is infallible with these options"); - Self { - base_url: base_url.into().trim_end_matches('/').to_string(), - http, - } - } -} - -impl Debug for HttpRegistry { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("HttpRegistry") - .field("base_url", &self.base_url) - .finish() - } -} - -impl RegistrationService for HttpRegistry { - type Error = HttpRegistryError; - - fn register( - &mut self, - identity: &dyn IdentityProvider, - key_bundle: Vec, - ) -> Result<(), HttpRegistryError> { - let device_id = hex::encode(identity.public_key().as_ref()); - let timestamp_ms = SystemTime::now() - .duration_since(UNIX_EPOCH) - .map_err(|_| HttpRegistryError::Clock)? - .as_millis() as u64; - - // Sign exactly the bytes that go on the wire. - let payload = encode_payload(timestamp_ms, &key_bundle); - let signature = identity.sign(&payload); - - let req = SubmitRequest { - device_id, - payload: BASE64.encode(&payload), - signature: BASE64.encode(signature.as_ref()), - }; - - let url = format!("{}/v0/keypackage", self.base_url); - let resp = send_retrying(|| self.http.post(&url).json(&req))?; - if !resp.status().is_success() { - let status = resp.status().as_u16(); - let body = resp.text().unwrap_or_default(); - return Err(HttpRegistryError::Server(status, body)); - } - Ok(()) - } - - fn retrieve(&self, device_id: &str) -> Result>, HttpRegistryError> { - let url = format!("{}/v0/keypackage/{}", self.base_url, device_id); - let resp = send_retrying(|| self.http.get(&url))?; - if resp.status().as_u16() == 404 { - return Ok(None); - } - if !resp.status().is_success() { - let status = resp.status().as_u16(); - let body = resp.text().unwrap_or_default(); - return Err(HttpRegistryError::Server(status, body)); - } - let body: FetchResponse = resp.json()?; - - let payload = BASE64 - .decode(&body.payload) - .map_err(|e| HttpRegistryError::Decode(e.to_string()))?; - let signature_arr: [u8; 64] = BASE64 - .decode(&body.signature) - .map_err(|e| HttpRegistryError::Decode(e.to_string()))? - .as_slice() - .try_into() - .map_err(|_| HttpRegistryError::Decode("signature not 64 bytes".into()))?; - - // Verify over the received payload bytes, using the key we asked for - // (`device_id`). A bundle the requested device didn't sign won't verify. - let device_pubkey: [u8; 32] = hex::decode(device_id) - .map_err(|e| HttpRegistryError::Decode(e.to_string()))? - .as_slice() - .try_into() - .map_err(|_| HttpRegistryError::Decode("device_id not a 32-byte key".into()))?; - let verifying_key = Ed25519VerifyingKey::from_bytes(&device_pubkey) - .map_err(|_| HttpRegistryError::Decode("device_id not a valid ed25519 vk".into()))?; - verifying_key - .verify(&payload, &Ed25519Signature::from(signature_arr)) - .map_err(|_| HttpRegistryError::SignatureInvalid)?; - - let (_timestamp_ms, key_package) = decode_payload(&payload) - .ok_or_else(|| HttpRegistryError::Decode("short payload".into()))?; - - Ok(Some(key_package.to_vec())) - } -} - -impl AccountDirectory for HttpRegistry { - type Error = HttpRegistryError; - - fn publish(&mut self, bundle: &SignedDeviceBundle) -> Result<(), Self::Error> { - let req = SubmitAccountRequest { - account_pub: hex::encode(bundle.account_pub.as_ref()), - payload: BASE64.encode(&bundle.payload), - signature: BASE64.encode(bundle.signature.as_ref()), - }; - - let url = format!("{}/v0/account", self.base_url); - let resp = send_retrying(|| self.http.post(&url).json(&req))?; - if !resp.status().is_success() { - let status = resp.status().as_u16(); - let body = resp.text().unwrap_or_default(); - return Err(HttpRegistryError::Server(status, body)); - } - Ok(()) - } - - fn fetch(&self, account: &Ed25519VerifyingKey) -> Result, Self::Error> { - let url = format!( - "{}/v0/account/{}", - self.base_url, - hex::encode(account.as_ref()) - ); - let resp = send_retrying(|| self.http.get(&url))?; - if resp.status().as_u16() == 404 { - return Ok(None); - } - if !resp.status().is_success() { - let status = resp.status().as_u16(); - let body = resp.text().unwrap_or_default(); - return Err(HttpRegistryError::Server(status, body)); - } - let body: FetchAccountResponse = resp.json()?; - - let payload = BASE64 - .decode(&body.payload) - .map_err(|e| HttpRegistryError::Decode(e.to_string()))?; - let signature_arr: [u8; 64] = BASE64 - .decode(&body.signature) - .map_err(|e| HttpRegistryError::Decode(e.to_string()))? - .as_slice() - .try_into() - .map_err(|_| HttpRegistryError::Decode("signature not 64 bytes".into()))?; - - // The directory service is untrusted: verify the account signature over - // the exact received bytes, and that the bundle is bound to the account - // we asked for, before handing back any device keys. - let bundle = SignedDeviceBundle { - account_pub: account.clone(), - payload, - signature: Ed25519Signature::from(signature_arr), - }; - let device_set = verify_bundle(account, &bundle)?; - Ok(Some(device_set)) - } -} - -/// Canonical binary payload — the bytes that are both signed and transmitted -/// verbatim. Opaque to the server; decoded only by consumers: -/// -/// ```text -/// timestamp_ms : u64 little-endian (8 bytes) -/// key_package : remaining bytes (variable, last → no length prefix needed) -/// ``` -/// -/// The fixed-width field first with the one variable field last makes every -/// byte string parse exactly one way — no delimiter, no ambiguity, even though -/// `key_package` is arbitrary bytes. The device verifying key is carried -/// alongside as `device_id`, not embedded here. -fn encode_payload(timestamp_ms: u64, key_package: &[u8]) -> Vec { - let mut out = Vec::with_capacity(8 + key_package.len()); - out.extend_from_slice(×tamp_ms.to_le_bytes()); - out.extend_from_slice(key_package); - out -} - -/// Inverse of [`encode_payload`]. Returns `None` if the payload is shorter than -/// the fixed header (`8`). -fn decode_payload(payload: &[u8]) -> Option<(u64, &[u8])> { - if payload.len() < 8 { - return None; - } - let timestamp_ms = u64::from_le_bytes(payload[..8].try_into().ok()?); - Some((timestamp_ms, &payload[8..])) -} - -/// Retry budget for the registry's transient, load-induced 5xx/429 responses. -/// The service is reliable request-by-request but sheds concurrent bursts, so a -/// few backed-off retries let a request land once the burst clears. On that path -/// each retry returns fast, so the added cost is the ~3s worst-case backoff sum, -/// well inside chat_module's ~20s init IPC budget. A fully unreachable registry -/// instead costs up to MAX_RETRIES times the reqwest timeout, which no retry -/// budget can rescue. -const MAX_RETRIES: u32 = 4; -const RETRY_BASE_MS: u64 = 200; -const RETRY_MAX_BACKOFF_MS: u64 = 2000; - -/// Send a request built by `build`, retrying transient failures — network errors -/// and 5xx/429 responses — with exponential backoff and full jitter. The -/// registry is reliable request-by-request but sheds concurrent bursts with a -/// 5xx, so a backed-off retry lands once the burst clears; a 4xx (and any other -/// final response) is returned to the caller unchanged. `build` is re-invoked per -/// attempt because sending consumes the builder. -fn send_retrying( - build: impl Fn() -> reqwest::blocking::RequestBuilder, -) -> Result { - let mut attempt = 0; - loop { - let outcome = build().send(); - let transient = match &outcome { - Err(_) => true, // network error / timeout: worth another try - Ok(resp) => is_transient_status(resp.status()), - }; - if !transient || attempt >= MAX_RETRIES { - return Ok(outcome?); - } - std::thread::sleep(backoff_with_jitter(attempt)); - attempt += 1; - } -} - -/// Whether a response status is worth retrying: 5xx (the registry sheds -/// concurrent load with these) or 429 (explicit backpressure). A 4xx is the -/// caller's fault and won't change on retry. -fn is_transient_status(status: reqwest::StatusCode) -> bool { - status.is_server_error() || status == reqwest::StatusCode::TOO_MANY_REQUESTS -} - -/// Full-jitter exponential backoff: a random delay in -/// `[0, min(RETRY_MAX_BACKOFF_MS, RETRY_BASE_MS * 2^attempt)]`. The jitter -/// decorrelates concurrent publishers so their retries don't collide into the -/// same burst that failed them. -fn backoff_with_jitter(attempt: u32) -> Duration { - let exp = RETRY_BASE_MS.saturating_mul(1u64 << attempt.min(16)); - Duration::from_millis(jitter_below(exp.min(RETRY_MAX_BACKOFF_MS))) -} - -/// A value in `[0, max]`, seeded from the wall clock's sub-second nanos — enough -/// entropy to spread retries across processes without pulling in an RNG crate. -fn jitter_below(max: u64) -> u64 { - if max == 0 { - return 0; - } - let nanos = SystemTime::now() - .duration_since(UNIX_EPOCH) - .map(|d| d.subsec_nanos() as u64) - .unwrap_or(0); - nanos % (max + 1) -} - -#[cfg(test)] -mod tests { - use super::*; - use crypto::Ed25519SigningKey; - - /// `encode_payload` / `decode_payload` round-trip, including a key_package - /// containing bytes that a delimiter scheme would choke on (`:`, `|`, NUL). - #[test] - fn payload_roundtrips_with_arbitrary_bytes() { - let ts = 1_700_000_000_000u64; - let key_package = b"mls:bytes|with\x00delimiters".to_vec(); - - let payload = encode_payload(ts, &key_package); - let (got_ts, got_kp) = decode_payload(&payload).unwrap(); - assert_eq!(got_ts, ts); - assert_eq!(got_kp, key_package.as_slice()); - } - - #[test] - fn decode_rejects_short_payload() { - assert!(decode_payload(&[0u8; 7]).is_none()); - } - - /// Tampering with any byte of the payload breaks verification. - #[test] - fn signature_binds_payload() { - let signing = Ed25519SigningKey::generate(); - let verifying = signing.verifying_key(); - - let payload = encode_payload(1_700_000_000_000, b"original-keypackage"); - let signature = signing.sign(&payload); - - let tampered = encode_payload(1_700_000_000_000, b"tampered-keypackage"); - verifying - .verify(&tampered, &signature) - .expect_err("signature must not verify against a different payload"); - } - - /// End-to-end of the wire crypto: verify over the received payload bytes - /// using the key recovered from device_id, exactly as `retrieve` does. - #[test] - fn sign_then_verify_over_payload() { - let signing = Ed25519SigningKey::generate(); - let pubkey: [u8; 32] = signing.verifying_key().as_ref().try_into().unwrap(); - let payload = encode_payload(1_700_000_000_000, b"fake-mls-keypackage-bytes"); - let signature = signing.sign(&payload); - - // retrieve side: recover key from device_id (hex of pubkey), verify payload. - let device_id = hex::encode(pubkey); - let recovered: [u8; 32] = hex::decode(&device_id) - .unwrap() - .as_slice() - .try_into() - .unwrap(); - Ed25519VerifyingKey::from_bytes(&recovered) - .unwrap() - .verify(&payload, &signature) - .expect("recovered key must verify the register-time signature"); - } - - /// Only 5xx and 429 are retried; 2xx/4xx are returned to the caller as-is. - #[test] - fn only_5xx_and_429_are_transient() { - use reqwest::StatusCode; - for s in [500u16, 502, 503, 504, 429] { - assert!( - is_transient_status(StatusCode::from_u16(s).unwrap()), - "{s} should be retried" - ); - } - for s in [200u16, 201, 400, 401, 404, 409] { - assert!( - !is_transient_status(StatusCode::from_u16(s).unwrap()), - "{s} should not be retried" - ); - } - } - - /// Backoff never exceeds the exponential ceiling for its attempt, nor the - /// absolute cap — and the exponent shift can't overflow at high attempts. - #[test] - fn backoff_stays_within_the_cap() { - for attempt in 0..40u32 { - let ceiling = RETRY_BASE_MS - .saturating_mul(1u64 << attempt.min(16)) - .min(RETRY_MAX_BACKOFF_MS); - let delay = backoff_with_jitter(attempt).as_millis() as u64; - assert!(delay <= ceiling, "attempt {attempt}: {delay} > {ceiling}"); - } - } - - #[test] - fn jitter_is_bounded() { - assert_eq!(jitter_below(0), 0); - for _ in 0..200 { - assert!(jitter_below(50) <= 50); - } - } -} diff --git a/extensions/components/src/contact_registry/store.rs b/extensions/components/src/contact_registry/store.rs new file mode 100644 index 0000000..251cf39 --- /dev/null +++ b/extensions/components/src/contact_registry/store.rs @@ -0,0 +1,636 @@ +use std::fmt::{self, Debug}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use base64::Engine; +use base64::engine::general_purpose::STANDARD as BASE64; +use chat_proto::logoschat::store::{AccountSubmissionV1, KeyPackageSubmissionV1}; +use crypto::{Ed25519Signature, Ed25519VerifyingKey}; +use libchat::{AddressedEnvelope, DeliveryService, IdentityProvider, RegistrationService}; +use logos_account::{AccountDirectory, BundleError, DeviceSet, SignedDeviceBundle, verify_bundle}; +use prost::Message; +use prost::bytes::Bytes; +use serde::{Deserialize, Serialize}; + +/// Delivery address the store listens on for keypackage submissions. The +/// transport maps it to its content topic (e.g. +/// `/logos-chat/1/store-keypackage-v0/proto` on logos-delivery); the store +/// subscribes to the same topic. +pub const KEYPACKAGE_SUBMIT_ADDRESS: &str = "store-keypackage-v0"; + +/// Delivery address the store listens on for account device-list bundles. +pub const ACCOUNT_SUBMIT_ADDRESS: &str = "store-account-v0"; + +/// Request timeout for the store's HTTP API (queries, and submissions in +/// [`RegistryPublishMode::Http`]). +const HTTP_TIMEOUT: Duration = Duration::from_secs(10); + +/// How a [`ContactRegistry`] submits bundles to the store. Reads always use +/// the store's HTTP query API; only the write half switches. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum RegistryPublishMode { + /// Submit via the store's HTTP POST endpoints (synchronous, acknowledged). + #[default] + Http, + /// Publish over the delivery network on the well-known store addresses; + /// the store subscribes and persists what verifies. Fire-and-forget — + /// there is no per-submission acknowledgement, which the registry can + /// afford because consumers verify every bundle on retrieval anyway. + Delivery, +} + +/// The keypackage store and account → device directory. +/// +/// Reads (keypackage retrieve, account fetch) always go over the store's HTTP +/// query API. Writes (register, publish) go over whichever wire +/// [`RegistryPublishMode`] selects: the store's HTTP POST endpoints (a JSON +/// body with hex + base64 fields) or a protobuf submission +/// ([`KeyPackageSubmissionV1`] / [`AccountSubmissionV1`]) published on the +/// well-known store addresses — matching the `/proto` content topics those +/// addresses map to, and carrying the keys, payload and signature as raw bytes. +/// +/// A single registry serves both wires so it can be used behind one +/// `ChatClient` registry type; the delivery transport `D` is unused in +/// [`RegistryPublishMode::Http`]. +#[derive(Clone)] +pub struct ContactRegistry { + base_url: String, + http: reqwest::blocking::Client, + delivery: D, + publish_mode: RegistryPublishMode, +} + +impl Debug for ContactRegistry { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("ContactRegistry") + .field("base_url", &self.base_url) + .field("publish_mode", &self.publish_mode) + .field("delivery", &self.delivery) + .finish() + } +} + +#[derive(Debug, thiserror::Error)] +pub enum ContactRegistryError { + #[error("http: {0}")] + Http(#[from] reqwest::Error), + #[error("server returned status {0}: {1}")] + Server(u16, String), + #[error("decode: {0}")] + Decode(String), + #[error("clock before unix epoch")] + Clock, + #[error("signature verification failed")] + SignatureInvalid, + #[error("bundle: {0}")] + Bundle(#[from] BundleError), + #[error("publish over delivery: {0}")] + Publish(String), +} + +impl ContactRegistry { + /// A registry that queries the store's HTTP API at `base_url` and submits + /// per `publish_mode` — over `delivery` or over the same HTTP API. + pub fn new( + delivery: D, + base_url: impl Into, + publish_mode: RegistryPublishMode, + ) -> Self { + let http = reqwest::blocking::Client::builder() + .timeout(HTTP_TIMEOUT) + .build() + .expect("reqwest client builder is infallible with these options"); + Self { + base_url: base_url.into().trim_end_matches('/').to_string(), + http, + delivery, + publish_mode, + } + } +} + +impl ContactRegistry { + /// POST `body` as JSON to `path` on the store, mapping a non-success status + /// to [`ContactRegistryError::Server`] with the server's own message. + fn http_post(&self, path: &str, body: &S) -> Result<(), ContactRegistryError> { + let url = format!("{}{}", self.base_url, path); + let resp = send_retrying(|| self.http.post(&url).json(body))?; + if !resp.status().is_success() { + let status = resp.status().as_u16(); + return Err(ContactRegistryError::Server( + status, + resp.text().unwrap_or_default(), + )); + } + Ok(()) + } + + /// GET `url`, returning `None` on 404 (never published) and the decoded, + /// still-unverified bundle on success. The caller verifies the signature. + fn http_fetch(&self, url: &str) -> Result, ContactRegistryError> { + let resp = send_retrying(|| self.http.get(url))?; + if resp.status().as_u16() == 404 { + return Ok(None); + } + if !resp.status().is_success() { + let status = resp.status().as_u16(); + return Err(ContactRegistryError::Server( + status, + resp.text().unwrap_or_default(), + )); + } + let body: FetchResponse = resp.json()?; + let payload = BASE64 + .decode(&body.payload) + .map_err(|e| ContactRegistryError::Decode(e.to_string()))?; + let signature: [u8; 64] = BASE64 + .decode(&body.signature) + .map_err(|e| ContactRegistryError::Decode(e.to_string()))? + .as_slice() + .try_into() + .map_err(|_| ContactRegistryError::Decode("signature not 64 bytes".into()))?; + Ok(Some(FetchedBundle { payload, signature })) + } +} + +impl ContactRegistry { + /// Encode `submission` as protobuf and publish it on `delivery_address`. + /// Protobuf encoding into a `Vec` cannot fail, so the only error here is the + /// transport's. + fn publish_submission( + &mut self, + delivery_address: &str, + submission: &M, + ) -> Result<(), ContactRegistryError> { + self.delivery + .publish(AddressedEnvelope { + delivery_address: delivery_address.to_string(), + data: submission.encode_to_vec(), + }) + .map_err(|e| ContactRegistryError::Publish(e.to_string())) + } +} + +impl RegistrationService for ContactRegistry { + type Error = ContactRegistryError; + + fn register( + &mut self, + identity: &dyn IdentityProvider, + key_bundle: Vec, + ) -> Result<(), Self::Error> { + let timestamp_ms = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_err(|_| ContactRegistryError::Clock)? + .as_millis() as u64; + + // The signed bytes are the same on both wires; only the encoding of the + // submission around them differs. Sign once, then branch on transport. + let payload = encode_payload(timestamp_ms, &key_bundle); + let signature = identity.sign(&payload); + let device_id = identity.public_key().as_ref(); + + match self.publish_mode { + RegistryPublishMode::Http => self.http_post( + "/v0/keypackage", + &SubmitRequest { + device_id: hex::encode(device_id), + payload: BASE64.encode(&payload), + signature: BASE64.encode(signature.as_ref()), + }, + ), + RegistryPublishMode::Delivery => { + let req = KeyPackageSubmissionV1 { + device_id: Bytes::copy_from_slice(device_id), + payload: Bytes::from(payload), + signature: Bytes::copy_from_slice(signature.as_ref()), + }; + self.publish_submission(KEYPACKAGE_SUBMIT_ADDRESS, &req) + } + } + } + + fn retrieve(&self, device_id: &str) -> Result>, Self::Error> { + let url = format!("{}/v0/keypackage/{}", self.base_url, device_id); + let Some(FetchedBundle { payload, signature }) = self.http_fetch(&url)? else { + return Ok(None); + }; + + // Verify over the received payload bytes, using the key we asked for + // (`device_id`). A bundle the requested device didn't sign won't verify. + let device_pubkey: [u8; 32] = hex::decode(device_id) + .map_err(|e| ContactRegistryError::Decode(e.to_string()))? + .as_slice() + .try_into() + .map_err(|_| ContactRegistryError::Decode("device_id not a 32-byte key".into()))?; + let verifying_key = Ed25519VerifyingKey::from_bytes(&device_pubkey) + .map_err(|_| ContactRegistryError::Decode("device_id not a valid ed25519 vk".into()))?; + verifying_key + .verify(&payload, &Ed25519Signature::from(signature)) + .map_err(|_| ContactRegistryError::SignatureInvalid)?; + + let (_timestamp_ms, key_package) = decode_payload(&payload) + .ok_or_else(|| ContactRegistryError::Decode("short payload".into()))?; + Ok(Some(key_package.to_vec())) + } +} + +impl AccountDirectory for ContactRegistry { + type Error = ContactRegistryError; + + fn publish(&mut self, bundle: &SignedDeviceBundle) -> Result<(), Self::Error> { + // The bundle is already signed; both wires carry its exact bytes. + match self.publish_mode { + RegistryPublishMode::Http => self.http_post( + "/v0/account", + &SubmitAccountRequest { + account_pub: hex::encode(bundle.account_pub.as_ref()), + payload: BASE64.encode(&bundle.payload), + signature: BASE64.encode(bundle.signature.as_ref()), + }, + ), + RegistryPublishMode::Delivery => { + let req = AccountSubmissionV1 { + account_pub: Bytes::copy_from_slice(bundle.account_pub.as_ref()), + payload: Bytes::copy_from_slice(&bundle.payload), + signature: Bytes::copy_from_slice(bundle.signature.as_ref()), + }; + self.publish_submission(ACCOUNT_SUBMIT_ADDRESS, &req) + } + } + } + + fn fetch(&self, account: &Ed25519VerifyingKey) -> Result, Self::Error> { + let url = format!( + "{}/v0/account/{}", + self.base_url, + hex::encode(account.as_ref()) + ); + let Some(FetchedBundle { payload, signature }) = self.http_fetch(&url)? else { + return Ok(None); + }; + + // The directory service is untrusted: verify the account signature over + // the exact received bytes, and that the bundle is bound to the account + // we asked for, before handing back any device keys. + let bundle = SignedDeviceBundle { + account_pub: account.clone(), + payload, + signature: Ed25519Signature::from(signature), + }; + let device_set = verify_bundle(account, &bundle)?; + Ok(Some(device_set)) + } +} + +/// Keypackage submission as the HTTP POST body. The delivery path carries the +/// same fields as protobuf; this JSON shape is the store's HTTP endpoint only. +#[derive(Debug, Serialize)] +struct SubmitRequest { + /// hex of the 32-byte device verifying key — the verification + storage key. + device_id: String, + /// base64 of the canonical signed payload (see [`encode_payload`]). + payload: String, + /// base64 of the 64-byte Ed25519 signature over `payload`. + signature: String, +} + +/// Account device-list submission as the HTTP POST body, like [`SubmitRequest`]. +#[derive(Debug, Serialize)] +struct SubmitAccountRequest { + /// hex of the 32-byte account verifying key — verification + storage key. + account_pub: String, + /// base64 of the canonical signed device-list payload. + payload: String, + /// base64 of the 64-byte account signature over `payload`. + signature: String, +} + +/// The `payload` + `signature` of a store fetch response; both keypackage and +/// account queries return this shape. +#[derive(Debug, Deserialize)] +struct FetchResponse { + payload: String, + signature: String, +} + +/// A fetch response with its base64 fields decoded but not yet verified. +struct FetchedBundle { + payload: Vec, + signature: [u8; 64], +} + +/// Canonical binary payload — the bytes that are both signed and transmitted +/// verbatim. Opaque to the server; decoded only by consumers: +/// +/// ```text +/// timestamp_ms : u64 little-endian (8 bytes) +/// key_package : remaining bytes (variable, last → no length prefix needed) +/// ``` +/// +/// The fixed-width field first with the one variable field last makes every +/// byte string parse exactly one way — no delimiter, no ambiguity, even though +/// `key_package` is arbitrary bytes. The device verifying key is carried +/// alongside as `device_id`, not embedded here. +fn encode_payload(timestamp_ms: u64, key_package: &[u8]) -> Vec { + let mut out = Vec::with_capacity(8 + key_package.len()); + out.extend_from_slice(×tamp_ms.to_le_bytes()); + out.extend_from_slice(key_package); + out +} + +/// Inverse of [`encode_payload`]. Returns `None` if the payload is shorter than +/// the fixed header (`8`). +fn decode_payload(payload: &[u8]) -> Option<(u64, &[u8])> { + if payload.len() < 8 { + return None; + } + let timestamp_ms = u64::from_le_bytes(payload[..8].try_into().ok()?); + Some((timestamp_ms, &payload[8..])) +} + +/// Retry budget for the registry's transient, load-induced 5xx/429 responses. +/// The service is reliable request-by-request but sheds concurrent bursts, so a +/// few backed-off retries let a request land once the burst clears. On that path +/// each retry returns fast, so the added cost is the ~3s worst-case backoff sum, +/// well inside chat_module's ~20s init IPC budget. A fully unreachable registry +/// instead costs up to MAX_RETRIES times the reqwest timeout, which no retry +/// budget can rescue. +const MAX_RETRIES: u32 = 4; +const RETRY_BASE_MS: u64 = 200; +const RETRY_MAX_BACKOFF_MS: u64 = 2000; + +/// Send a request built by `build`, retrying transient failures — network errors +/// and 5xx/429 responses — with exponential backoff and full jitter. The +/// registry is reliable request-by-request but sheds concurrent bursts with a +/// 5xx, so a backed-off retry lands once the burst clears; a 4xx (and any other +/// final response) is returned to the caller unchanged. `build` is re-invoked per +/// attempt because sending consumes the builder. +fn send_retrying( + build: impl Fn() -> reqwest::blocking::RequestBuilder, +) -> Result { + let mut attempt = 0; + loop { + let outcome = build().send(); + let transient = match &outcome { + Err(_) => true, // network error / timeout: worth another try + Ok(resp) => is_transient_status(resp.status()), + }; + if !transient || attempt >= MAX_RETRIES { + return Ok(outcome?); + } + std::thread::sleep(backoff_with_jitter(attempt)); + attempt += 1; + } +} + +/// Whether a response status is worth retrying: 5xx (the registry sheds +/// concurrent load with these) or 429 (explicit backpressure). A 4xx is the +/// caller's fault and won't change on retry. +fn is_transient_status(status: reqwest::StatusCode) -> bool { + status.is_server_error() || status == reqwest::StatusCode::TOO_MANY_REQUESTS +} + +/// Full-jitter exponential backoff: a random delay in +/// `[0, min(RETRY_MAX_BACKOFF_MS, RETRY_BASE_MS * 2^attempt)]`. The jitter +/// decorrelates concurrent publishers so their retries don't collide into the +/// same burst that failed them. +fn backoff_with_jitter(attempt: u32) -> Duration { + let exp = RETRY_BASE_MS.saturating_mul(1u64 << attempt.min(16)); + Duration::from_millis(jitter_below(exp.min(RETRY_MAX_BACKOFF_MS))) +} + +/// A value in `[0, max]`, seeded from the wall clock's sub-second nanos — enough +/// entropy to spread retries across processes without pulling in an RNG crate. +fn jitter_below(max: u64) -> u64 { + if max == 0 { + return 0; + } + let nanos = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.subsec_nanos() as u64) + .unwrap_or(0); + nanos % (max + 1) +} + +#[cfg(test)] +mod tests { + use super::*; + use crypto::Ed25519SigningKey; + use libchat::{IdentId, IdentIdRef}; + + #[derive(Debug, Default)] + struct CapturingDelivery { + published: Vec, + } + + impl DeliveryService for CapturingDelivery { + type Error = std::convert::Infallible; + fn publish(&mut self, envelope: AddressedEnvelope) -> Result<(), Self::Error> { + self.published.push(envelope); + Ok(()) + } + fn subscribe(&mut self, _delivery_address: &str) -> Result<(), Self::Error> { + Ok(()) + } + } + + struct TestIdent { + id: IdentId, + key: Ed25519SigningKey, + verifying: Ed25519VerifyingKey, + } + + impl TestIdent { + fn new() -> Self { + let key = Ed25519SigningKey::generate(); + let verifying = key.verifying_key(); + Self { + id: IdentId::new("test"), + key, + verifying, + } + } + } + + impl IdentityProvider for TestIdent { + fn id(&self) -> IdentIdRef<'_> { + &self.id + } + fn display_name(&self) -> String { + self.id.to_string() + } + fn sign(&self, payload: &[u8]) -> Ed25519Signature { + self.key.sign(payload) + } + fn public_key(&self) -> &Ed25519VerifyingKey { + &self.verifying + } + } + + #[test] + fn register_publishes_store_submission_on_keypackage_address() { + let mut registry = ContactRegistry::new( + CapturingDelivery::default(), + "http://unused.invalid", + RegistryPublishMode::Delivery, + ); + let ident = TestIdent::new(); + let key_bundle = b"kp-bytes".to_vec(); + registry.register(&ident, key_bundle.clone()).unwrap(); + + let [envelope] = ®istry.delivery.published[..] else { + panic!("expected exactly one published envelope"); + }; + assert_eq!(envelope.delivery_address, KEYPACKAGE_SUBMIT_ADDRESS); + + // Decode as the store does: the bytes on the wire are a protobuf + // submission, so a field-number or type change here breaks ingestion. + let wire = KeyPackageSubmissionV1::decode(&envelope.data[..]).unwrap(); + assert_eq!(wire.device_id.as_ref(), ident.verifying.as_ref()); + // The store verifies the signature over the payload bytes under the + // device key before persisting — the submission must pass that check. + assert!(wire.payload.ends_with(&key_bundle)); + let signature: [u8; 64] = wire.signature.as_ref().try_into().unwrap(); + ident + .verifying + .verify(&wire.payload, &Ed25519Signature::from(signature)) + .expect("store-side verification must succeed"); + } + + #[test] + fn account_publish_targets_account_address_verbatim() { + let mut registry = ContactRegistry::new( + CapturingDelivery::default(), + "http://unused.invalid", + RegistryPublishMode::Delivery, + ); + let account = Ed25519SigningKey::generate(); + let payload = b"signed-device-list".to_vec(); + let bundle = SignedDeviceBundle { + account_pub: account.verifying_key(), + signature: account.sign(&payload), + payload: payload.clone(), + }; + registry.publish(&bundle).unwrap(); + + let [envelope] = ®istry.delivery.published[..] else { + panic!("expected exactly one published envelope"); + }; + assert_eq!(envelope.delivery_address, ACCOUNT_SUBMIT_ADDRESS); + + let wire = AccountSubmissionV1::decode(&envelope.data[..]).unwrap(); + assert_eq!(wire.account_pub.as_ref(), bundle.account_pub.as_ref()); + // Payload travels verbatim so the store and consumers verify the exact + // signed bytes. + assert_eq!(wire.payload.as_ref(), payload.as_slice()); + assert_eq!(wire.signature.as_ref(), bundle.signature.as_ref()); + } + + #[test] + fn http_mode_never_touches_the_delivery_service() { + // Port 9 (discard) refuses immediately; the point is only that the + // submission goes down the HTTP path, not over delivery. + let mut registry = ContactRegistry::new( + CapturingDelivery::default(), + "http://127.0.0.1:9", + RegistryPublishMode::Http, + ); + let err = registry.register(&TestIdent::new(), vec![1]).unwrap_err(); + assert!(matches!(err, ContactRegistryError::Http(_))); + assert!(registry.delivery.published.is_empty()); + } + + /// `encode_payload` / `decode_payload` round-trip, including a key_package + /// containing bytes that a delimiter scheme would choke on (`:`, `|`, NUL). + #[test] + fn payload_roundtrips_with_arbitrary_bytes() { + let ts = 1_700_000_000_000u64; + let key_package = b"mls:bytes|with\x00delimiters".to_vec(); + + let payload = encode_payload(ts, &key_package); + let (got_ts, got_kp) = decode_payload(&payload).unwrap(); + assert_eq!(got_ts, ts); + assert_eq!(got_kp, key_package.as_slice()); + } + + #[test] + fn decode_rejects_short_payload() { + assert!(decode_payload(&[0u8; 7]).is_none()); + } + + /// Tampering with any byte of the payload breaks verification. + #[test] + fn signature_binds_payload() { + let signing = Ed25519SigningKey::generate(); + let verifying = signing.verifying_key(); + + let payload = encode_payload(1_700_000_000_000, b"original-keypackage"); + let signature = signing.sign(&payload); + + let tampered = encode_payload(1_700_000_000_000, b"tampered-keypackage"); + verifying + .verify(&tampered, &signature) + .expect_err("signature must not verify against a different payload"); + } + + /// End-to-end of the wire crypto: verify over the received payload bytes + /// using the key recovered from device_id, exactly as `retrieve` does. + #[test] + fn sign_then_verify_over_payload() { + let signing = Ed25519SigningKey::generate(); + let pubkey: [u8; 32] = signing.verifying_key().as_ref().try_into().unwrap(); + let payload = encode_payload(1_700_000_000_000, b"fake-mls-keypackage-bytes"); + let signature = signing.sign(&payload); + + // retrieve side: recover key from device_id (hex of pubkey), verify payload. + let device_id = hex::encode(pubkey); + let recovered: [u8; 32] = hex::decode(&device_id) + .unwrap() + .as_slice() + .try_into() + .unwrap(); + Ed25519VerifyingKey::from_bytes(&recovered) + .unwrap() + .verify(&payload, &signature) + .expect("recovered key must verify the register-time signature"); + } + + /// Only 5xx and 429 are retried; 2xx/4xx are returned to the caller as-is. + #[test] + fn only_5xx_and_429_are_transient() { + use reqwest::StatusCode; + for s in [500u16, 502, 503, 504, 429] { + assert!( + is_transient_status(StatusCode::from_u16(s).unwrap()), + "{s} should be retried" + ); + } + for s in [200u16, 201, 400, 401, 404, 409] { + assert!( + !is_transient_status(StatusCode::from_u16(s).unwrap()), + "{s} should not be retried" + ); + } + } + + /// Backoff never exceeds the exponential ceiling for its attempt, nor the + /// absolute cap — and the exponent shift can't overflow at high attempts. + #[test] + fn backoff_stays_within_the_cap() { + for attempt in 0..40u32 { + let ceiling = RETRY_BASE_MS + .saturating_mul(1u64 << attempt.min(16)) + .min(RETRY_MAX_BACKOFF_MS); + let delay = backoff_with_jitter(attempt).as_millis() as u64; + assert!(delay <= ceiling, "attempt {attempt}: {delay} > {ceiling}"); + } + } + + #[test] + fn jitter_is_bounded() { + assert_eq!(jitter_below(0), 0); + for _ in 0..200 { + assert!(jitter_below(50) <= 50); + } + } +} diff --git a/extensions/components/src/lib.rs b/extensions/components/src/lib.rs index 92ad4a4..70e05f3 100644 --- a/extensions/components/src/lib.rs +++ b/extensions/components/src/lib.rs @@ -4,7 +4,10 @@ mod storage; mod wakeup; pub use contact_registry::ephemeral::EphemeralRegistry; -pub use contact_registry::http::{HttpRegistry, HttpRegistryError}; +pub use contact_registry::store::{ + ACCOUNT_SUBMIT_ADDRESS, ContactRegistry, ContactRegistryError, KEYPACKAGE_SUBMIT_ADDRESS, + RegistryPublishMode, +}; pub use delivery::*; pub use storage::*; pub use wakeup::*; diff --git a/extensions/logos-delivery-rust/Cargo.toml b/extensions/logos-delivery-rust/Cargo.toml index e9dc546..bae3f20 100644 --- a/extensions/logos-delivery-rust/Cargo.toml +++ b/extensions/logos-delivery-rust/Cargo.toml @@ -7,8 +7,6 @@ links = "logosdelivery" [dependencies] # Workspace dependencies (sorted) crossbeam-channel = { workspace = true } -libchat = { workspace = true } -logos-generic-chat = { workspace = true } # External dependencies (sorted) base64 = "0.22"