fix: deletable self-sent mail, honest refused-handshake reporting, registration remedy
This commit is contained in:
parent
4a780ed648
commit
39c4c4d6a5
3 changed files with 211 additions and 55 deletions
|
|
@ -174,7 +174,9 @@ fn hint(error: Option<&fumi::error::Error>) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Some(fumi::error::Error::NotRegistered) => {
|
Some(fumi::error::Error::NotRegistered) => {
|
||||||
eprintln!("run: fumi register <user@host>")
|
eprintln!(
|
||||||
|
"for a new mailbox:\n fumi register <user@host>\nfor an existing one, pin the server and restore:\n fumi trust <host> <key>\n fumi restore <user@host>"
|
||||||
|
)
|
||||||
}
|
}
|
||||||
_ => {}
|
_ => {}
|
||||||
}
|
}
|
||||||
|
|
@ -675,11 +677,17 @@ fn resolve_id_prefixes(store: &Store, ids: &[String]) -> Result<Vec<[u8; ID_LEN]
|
||||||
let mut full = Vec::with_capacity(ids.len());
|
let mut full = Vec::with_capacity(ids.len());
|
||||||
for id in ids {
|
for id in ids {
|
||||||
let prefix = id.to_ascii_lowercase();
|
let prefix = id.to_ascii_lowercase();
|
||||||
let matches: Vec<[u8; ID_LEN]> = stored
|
let mut matches: Vec<[u8; ID_LEN]> = stored
|
||||||
.iter()
|
.iter()
|
||||||
.filter(|m| hex(m).starts_with(&prefix))
|
.filter(|m| hex(m).starts_with(&prefix))
|
||||||
.copied()
|
.copied()
|
||||||
.collect();
|
.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() {
|
match matches.len() {
|
||||||
0 => return Err(anyhow!(format!("no message matching {id:?}"))),
|
0 => return Err(anyhow!(format!("no message matching {id:?}"))),
|
||||||
1 => full.push(matches[0]),
|
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(951_782_400), "2000-02-29 00:00"); // leap day
|
||||||
assert_eq!(format_utc(-1), "1969-12-31 23:59");
|
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]);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -55,6 +55,16 @@ pub enum Error {
|
||||||
port: u16,
|
port: u16,
|
||||||
source: std::io::Error,
|
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")]
|
#[cfg(feature = "store")]
|
||||||
Storage(rusqlite::Error),
|
Storage(rusqlite::Error),
|
||||||
Noise(snow::Error),
|
Noise(snow::Error),
|
||||||
|
|
@ -100,7 +110,10 @@ impl fmt::Display for Error {
|
||||||
b32(known),
|
b32(known),
|
||||||
b32(offered)
|
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!(
|
Self::ChainLimit { index } => write!(
|
||||||
f,
|
f,
|
||||||
"rotation index {index} exceeds the chain limit of {MAX_CHAIN}"
|
"rotation index {index} exceeds the chain limit of {MAX_CHAIN}"
|
||||||
|
|
@ -112,6 +125,16 @@ impl fmt::Display for Error {
|
||||||
Self::Unreachable { host, port, source } => {
|
Self::Unreachable { host, port, source } => {
|
||||||
write!(f, "cannot reach {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")]
|
#[cfg(feature = "store")]
|
||||||
Self::Storage(e) => write!(f, "storage error: {e}"),
|
Self::Storage(e) => write!(f, "storage error: {e}"),
|
||||||
Self::Noise(e) => write!(f, "noise 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 {
|
impl std::error::Error for Error {
|
||||||
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
|
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
|
||||||
match self {
|
match self {
|
||||||
Self::Unreachable { source, .. } => Some(source),
|
Self::Unreachable { source, .. } | Self::HandshakeRefused { source, .. } => {
|
||||||
|
Some(source)
|
||||||
|
}
|
||||||
#[cfg(feature = "store")]
|
#[cfg(feature = "store")]
|
||||||
Self::Storage(e) => Some(e),
|
Self::Storage(e) => Some(e),
|
||||||
Self::Noise(e) => Some(e),
|
Self::Noise(e) => Some(e),
|
||||||
|
|
|
||||||
199
core/src/tcp.rs
199
core/src/tcp.rs
|
|
@ -25,66 +25,54 @@ pub struct TcpTransport {
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TcpTransport {
|
impl TcpTransport {
|
||||||
/// Runs the Noise_NX handshake as the initiator. The NX pattern has the
|
/// Runs the Noise_NX handshake as the initiator, retrying a bounded
|
||||||
/// server transmit its static key during the handshake, so pinning is a
|
/// number of times when the connection is dropped mid-handshake: a rate
|
||||||
/// single code path: an existing pin must match, an absent one is
|
/// limiter drops connections without a reply, so a refusal there is
|
||||||
/// accepted but reported as unpinned by `pinned()`.
|
/// 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(
|
pub fn connect(
|
||||||
host: &str,
|
host: &str,
|
||||||
port: u16,
|
port: u16,
|
||||||
pinned: Option<[u8; KEY_LEN]>,
|
pinned: Option<[u8; KEY_LEN]>,
|
||||||
timeout: u64,
|
timeout: u64,
|
||||||
) -> Result<TcpTransport, Error> {
|
) -> Result<TcpTransport, Error> {
|
||||||
let stream = TcpStream::connect((host, port)).map_err(|source| Error::Unreachable {
|
let mut refused = None;
|
||||||
host: host.to_string(),
|
for attempt in 1..=HANDSHAKE_ATTEMPTS {
|
||||||
port,
|
match handshake(host, port, timeout) {
|
||||||
source,
|
Ok((stream, noise, server_static, bind)) => {
|
||||||
})?;
|
if let Some(pinned) = pinned {
|
||||||
stream.set_read_timeout(Some(Duration::from_secs(timeout)))?;
|
if !ct_eq(&pinned, &server_static) {
|
||||||
stream.set_write_timeout(Some(Duration::from_secs(timeout)))?;
|
return Err(Error::PinMismatch {
|
||||||
let mut stream = stream;
|
host: host.to_string(),
|
||||||
|
pinned,
|
||||||
let params: snow::params::NoiseParams = NOISE_PARAMS.parse()?;
|
presented: server_static,
|
||||||
let mut noise = Builder::new(params).prologue(PROLOGUE).build_initiator()?;
|
});
|
||||||
|
}
|
||||||
let mut buf = [0u8; 65535];
|
}
|
||||||
let len = noise.write_message(&[], &mut buf)?;
|
return Ok(TcpTransport {
|
||||||
write_u16_len(&mut stream, &buf[..len])?;
|
stream,
|
||||||
|
noise,
|
||||||
let len = read_u16_len(&mut stream)?;
|
buf: Vec::new(),
|
||||||
let mut message = vec![0u8; len];
|
bind,
|
||||||
stream.read_exact(&mut message)?;
|
pinned: pinned.is_some(),
|
||||||
noise.read_message(&message, &mut buf)?;
|
host: host.to_string(),
|
||||||
if !noise.is_handshake_finished() {
|
});
|
||||||
return Err(Error::Other("handshake did not complete".into()));
|
}
|
||||||
}
|
Err(Error::Io(source)) if dropped_mid_handshake(&source) => {
|
||||||
|
if let Some(backoff_ms) = HANDSHAKE_BACKOFF_MS.get(attempt - 1) {
|
||||||
let server_static: [u8; KEY_LEN] = noise
|
std::thread::sleep(Duration::from_millis(*backoff_ms));
|
||||||
.get_remote_static()
|
}
|
||||||
.ok_or_else(|| Error::Other("server sent no static key".into()))?
|
refused = Some(source);
|
||||||
.try_into()
|
}
|
||||||
.map_err(|_| Error::Other("server static key is not 32 bytes".into()))?;
|
Err(e) => return Err(e),
|
||||||
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,
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
Err(Error::HandshakeRefused {
|
||||||
Ok(TcpTransport {
|
|
||||||
stream,
|
|
||||||
noise: transport,
|
|
||||||
buf: Vec::new(),
|
|
||||||
bind,
|
|
||||||
pinned: pinned.is_some(),
|
|
||||||
host: host.to_string(),
|
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 {
|
impl Transport for TcpTransport {
|
||||||
/// Sends one application frame (u32 length || op || body, split across
|
/// Sends one application frame (u32 length || op || body, split across
|
||||||
/// as many Noise messages as it needs) and reads the response frame,
|
/// 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.len() as u16).to_be_bytes())?;
|
||||||
stream.write_all(packet)
|
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:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue