From 39c4c4d6a507e129ae102d5b039fa563b2f4a3a5 Mon Sep 17 00:00:00 2001 From: randogoth Date: Tue, 29 Sep 2026 09:35:46 +0300 Subject: [PATCH] fix: deletable self-sent mail, honest refused-handshake reporting, registration remedy --- cli/src/main.rs | 38 ++++++++- core/src/error.rs | 29 ++++++- core/src/tcp.rs | 199 ++++++++++++++++++++++++++++++++++------------ 3 files changed, 211 insertions(+), 55 deletions(-) diff --git a/cli/src/main.rs b/cli/src/main.rs index 807c8eb..92cb8de 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -174,7 +174,9 @@ fn hint(error: Option<&fumi::error::Error>) { } } Some(fumi::error::Error::NotRegistered) => { - eprintln!("run: fumi register ") + eprintln!( + "for a new mailbox:\n fumi register \nfor an existing one, pin the server and restore:\n fumi trust \n fumi restore " + ) } _ => {} } @@ -675,11 +677,17 @@ fn resolve_id_prefixes(store: &Store, ids: &[String]) -> Result = stored + let mut matches: Vec<[u8; ID_LEN]> = stored .iter() .filter(|m| hex(m).starts_with(&prefix)) .copied() .collect(); + // A self-addressed message is stored once as received and once as + // sent with the same id; identical ids across folders are one + // message by construction, so dedupe before requiring more + // specificity. + matches.sort_unstable(); + matches.dedup(); match matches.len() { 0 => return Err(anyhow!(format!("no message matching {id:?}"))), 1 => full.push(matches[0]), @@ -887,4 +895,30 @@ mod tests { assert_eq!(format_utc(951_782_400), "2000-02-29 00:00"); // leap day assert_eq!(format_utc(-1), "1969-12-31 23:59"); } + + /// A self-addressed message sits in the inbox and the sent folder with + /// the same id; resolving it must yield one message, not an ambiguity. + #[test] + fn id_resolution_treats_a_self_addressed_message_as_one() { + let store = Store::open_in_memory().unwrap(); + let id = [7u8; ID_LEN]; + store.store_inbox(&id, b"envelope", 1, false, false).unwrap(); + store.store_sent(&id, "fumi@localhost", b"envelope", 1).unwrap(); + let full = resolve_id_prefixes(&store, &[hex(&id)]).unwrap(); + assert_eq!(full, vec![id]); + } + + /// Genuinely distinct messages sharing a prefix must still demand a + /// more specific id. + #[test] + fn id_resolution_rejects_an_ambiguous_prefix() { + let store = Store::open_in_memory().unwrap(); + let a = [7u8; ID_LEN]; + let mut b = [7u8; ID_LEN]; + b[31] = 8; + store.store_inbox(&a, b"a", 1, false, false).unwrap(); + store.store_inbox(&b, b"b", 2, false, false).unwrap(); + assert!(resolve_id_prefixes(&store, &["0707".into()]).is_err()); + assert_eq!(resolve_id_prefixes(&store, &[hex(&a)]).unwrap(), vec![a]); + } } diff --git a/core/src/error.rs b/core/src/error.rs index c2f1347..4270e37 100644 --- a/core/src/error.rs +++ b/core/src/error.rs @@ -55,6 +55,16 @@ pub enum Error { port: u16, source: std::io::Error, }, + /// The server dropped the connection mid-handshake after bounded + /// retries. A rate limiter drops connections without a reply, so this + /// is indistinguishable from a dying server; the decision it supports + /// is "retry later", not "treat the server as dead". + HandshakeRefused { + host: String, + port: u16, + source: std::io::Error, + attempts: usize, + }, #[cfg(feature = "store")] Storage(rusqlite::Error), Noise(snow::Error), @@ -100,7 +110,10 @@ impl fmt::Display for Error { b32(known), b32(offered) ), - Self::NotRegistered => write!(f, "not registered"), + Self::NotRegistered => write!( + f, + "not registered; this store has no account (registering or restoring creates one)" + ), Self::ChainLimit { index } => write!( f, "rotation index {index} exceeds the chain limit of {MAX_CHAIN}" @@ -112,6 +125,16 @@ impl fmt::Display for Error { Self::Unreachable { host, port, source } => { write!(f, "cannot reach {host}:{port}: {source}") } + Self::HandshakeRefused { + host, + port, + source, + attempts, + } => write!( + f, + "{host}:{port} refused the connection mid-handshake after {attempts} attempt(s) ({source})\n \ + possibly rate limited; a limiter drops connections without a reply, so a dying server looks the same" + ), #[cfg(feature = "store")] Self::Storage(e) => write!(f, "storage error: {e}"), Self::Noise(e) => write!(f, "noise error: {e}"), @@ -124,7 +147,9 @@ impl fmt::Display for Error { impl std::error::Error for Error { fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { match self { - Self::Unreachable { source, .. } => Some(source), + Self::Unreachable { source, .. } | Self::HandshakeRefused { source, .. } => { + Some(source) + } #[cfg(feature = "store")] Self::Storage(e) => Some(e), Self::Noise(e) => Some(e), diff --git a/core/src/tcp.rs b/core/src/tcp.rs index 18c5483..3a88846 100644 --- a/core/src/tcp.rs +++ b/core/src/tcp.rs @@ -25,66 +25,54 @@ pub struct TcpTransport { } impl TcpTransport { - /// Runs the Noise_NX handshake as the initiator. The NX pattern has the - /// server transmit its static key during the handshake, so pinning is a - /// single code path: an existing pin must match, an absent one is - /// accepted but reported as unpinned by `pinned()`. + /// Runs the Noise_NX handshake as the initiator, retrying a bounded + /// number of times when the connection is dropped mid-handshake: a rate + /// limiter drops connections without a reply, so a refusal there is + /// indistinguishable from a dying server and earns both a retry and an + /// honest report. Drops during the session proper are not retried; only + /// the handshake is idempotent. pub fn connect( host: &str, port: u16, pinned: Option<[u8; KEY_LEN]>, timeout: u64, ) -> Result { - let stream = TcpStream::connect((host, port)).map_err(|source| Error::Unreachable { - host: host.to_string(), - port, - source, - })?; - stream.set_read_timeout(Some(Duration::from_secs(timeout)))?; - stream.set_write_timeout(Some(Duration::from_secs(timeout)))?; - let mut stream = stream; - - let params: snow::params::NoiseParams = NOISE_PARAMS.parse()?; - let mut noise = Builder::new(params).prologue(PROLOGUE).build_initiator()?; - - let mut buf = [0u8; 65535]; - let len = noise.write_message(&[], &mut buf)?; - write_u16_len(&mut stream, &buf[..len])?; - - let len = read_u16_len(&mut stream)?; - let mut message = vec![0u8; len]; - stream.read_exact(&mut message)?; - noise.read_message(&message, &mut buf)?; - if !noise.is_handshake_finished() { - return Err(Error::Other("handshake did not complete".into())); - } - - let server_static: [u8; KEY_LEN] = noise - .get_remote_static() - .ok_or_else(|| Error::Other("server sent no static key".into()))? - .try_into() - .map_err(|_| Error::Other("server static key is not 32 bytes".into()))?; - let handshake_hash = noise.get_handshake_hash().to_vec(); - let bind = TransportBindValues::tcp(&handshake_hash, &server_static)?; - let transport = noise.into_transport_mode()?; - - if let Some(pinned) = pinned { - if !ct_eq(&pinned, &server_static) { - return Err(Error::PinMismatch { - host: host.to_string(), - pinned, - presented: server_static, - }); + let mut refused = None; + for attempt in 1..=HANDSHAKE_ATTEMPTS { + match handshake(host, port, timeout) { + Ok((stream, noise, server_static, bind)) => { + if let Some(pinned) = pinned { + if !ct_eq(&pinned, &server_static) { + return Err(Error::PinMismatch { + host: host.to_string(), + pinned, + presented: server_static, + }); + } + } + return Ok(TcpTransport { + stream, + noise, + buf: Vec::new(), + bind, + pinned: pinned.is_some(), + host: host.to_string(), + }); + } + Err(Error::Io(source)) if dropped_mid_handshake(&source) => { + if let Some(backoff_ms) = HANDSHAKE_BACKOFF_MS.get(attempt - 1) { + std::thread::sleep(Duration::from_millis(*backoff_ms)); + } + refused = Some(source); + } + Err(e) => return Err(e), } } - - Ok(TcpTransport { - stream, - noise: transport, - buf: Vec::new(), - bind, - pinned: pinned.is_some(), + Err(Error::HandshakeRefused { host: host.to_string(), + port, + source: refused.expect("a drop was seen on every failed attempt"), + attempts: HANDSHAKE_ATTEMPTS, }) } @@ -117,6 +105,70 @@ impl TcpTransport { } } +/// Attempts a dropped handshake gets before giving up, and the backoff +/// between them (none after the last). +const HANDSHAKE_ATTEMPTS: usize = 3; +const HANDSHAKE_BACKOFF_MS: [u64; 2] = [500, 1500]; + +/// Whether an io failure during the handshake means the server dropped the +/// connection rather than spoke something wrong. Only these kinds are +/// retried and reported as a refusal; a malformed reply is a different +/// failure and surfaces immediately. +fn dropped_mid_handshake(e: &std::io::Error) -> bool { + matches!( + e.kind(), + std::io::ErrorKind::UnexpectedEof + | std::io::ErrorKind::ConnectionReset + | std::io::ErrorKind::ConnectionAborted + | std::io::ErrorKind::BrokenPipe + | std::io::ErrorKind::NotConnected + ) +} + +/// One Noise_NX handshake attempt: dial, exchange the two handshake +/// messages, extract the server's static key and the bind values. The NX +/// pattern has the server transmit its static key during the handshake, so +/// pinning is a single code path in the caller: an existing pin must match, +/// an absent one is accepted but reported as unpinned. +fn handshake( + host: &str, + port: u16, + timeout: u64, +) -> Result<(TcpStream, TransportState, [u8; KEY_LEN], TransportBindValues), Error> { + let mut stream = TcpStream::connect((host, port)).map_err(|source| Error::Unreachable { + host: host.to_string(), + port, + source, + })?; + stream.set_read_timeout(Some(Duration::from_secs(timeout)))?; + stream.set_write_timeout(Some(Duration::from_secs(timeout)))?; + + let params: snow::params::NoiseParams = NOISE_PARAMS.parse()?; + let mut noise = Builder::new(params).prologue(PROLOGUE).build_initiator()?; + + let mut buf = [0u8; 65535]; + let len = noise.write_message(&[], &mut buf)?; + write_u16_len(&mut stream, &buf[..len])?; + + let len = read_u16_len(&mut stream)?; + let mut message = vec![0u8; len]; + stream.read_exact(&mut message)?; + noise.read_message(&message, &mut buf)?; + if !noise.is_handshake_finished() { + return Err(Error::Other("handshake did not complete".into())); + } + + let server_static: [u8; KEY_LEN] = noise + .get_remote_static() + .ok_or_else(|| Error::Other("server sent no static key".into()))? + .try_into() + .map_err(|_| Error::Other("server static key is not 32 bytes".into()))?; + let handshake_hash = noise.get_handshake_hash().to_vec(); + let bind = TransportBindValues::tcp(&handshake_hash, &server_static)?; + let noise = noise.into_transport_mode()?; + Ok((stream, noise, server_static, bind)) +} + impl Transport for TcpTransport { /// Sends one application frame (u32 length || op || body, split across /// as many Noise messages as it needs) and reads the response frame, @@ -177,3 +229,48 @@ fn write_u16_len(stream: &mut TcpStream, packet: &[u8]) -> std::io::Result<()> { stream.write_all(&(packet.len() as u16).to_be_bytes())?; stream.write_all(packet) } + +#[cfg(test)] +mod tests { + use super::*; + + /// A server that accepts and drops every connection mid-handshake is + /// what a rate limiter looks like from the client side: the handshake + /// is retried, then reported honestly as a refusal. + #[test] + fn a_dropped_handshake_is_retried_then_reported() { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let port = listener.local_addr().unwrap().port(); + let dropper = std::thread::spawn(move || { + for _ in 0..HANDSHAKE_ATTEMPTS { + if let Ok((stream, _)) = listener.accept() { + drop(stream); + } + } + }); + let error = TcpTransport::connect("127.0.0.1", port, None, 5) + .err() + .expect("expected HandshakeRefused, got a connection"); + match error { + Error::HandshakeRefused { attempts, .. } => assert_eq!(attempts, HANDSHAKE_ATTEMPTS), + other => panic!("expected HandshakeRefused, got {other:?}"), + } + dropper.join().unwrap(); + } + + /// A TCP-level refusal is not a mid-handshake drop: it surfaces + /// immediately as unreachable, without burning retries. + #[test] + fn a_closed_port_is_unreachable_without_retry() { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let port = listener.local_addr().unwrap().port(); + drop(listener); + let error = TcpTransport::connect("127.0.0.1", port, None, 1) + .err() + .expect("expected Unreachable, got a connection"); + match error { + Error::Unreachable { .. } => {} + other => panic!("expected Unreachable, got {other:?}"), + } + } +}