990 lines
34 KiB
Rust
990 lines
34 KiB
Rust
//! The six operations, AUTH session setup, chain walking, token sync and the
|
|
//! send/fetch pipelines. Everything here is carrier-neutral: it speaks
|
|
//! through the `Transport` trait and the bodies are byte-identical on both
|
|
//! carriers (RNS.md sec 13.5, sec 15). The operations that take a `Store`
|
|
//! need the `store` feature; `resolve`, `walk_chain` and `accept_mac` are
|
|
//! always built.
|
|
|
|
use crate::account::RotationCert;
|
|
use crate::crypto::{ct_eq, hmac_sha256, KEY_LEN};
|
|
use crate::error::{expect_ok, Error};
|
|
use crate::transport::{
|
|
Reader, Transport, CERT_LEN, ID_LEN, LABEL_MAC, MAX_CHAIN, OP_RESOLVE, TOKEN_LEN,
|
|
};
|
|
|
|
#[cfg(feature = "store")]
|
|
use crate::account::{Account, Identity};
|
|
#[cfg(feature = "store")]
|
|
use std::collections::BTreeMap;
|
|
|
|
#[cfg(feature = "store")]
|
|
use crate::address::{Address, Scheme};
|
|
#[cfg(feature = "store")]
|
|
use crate::crypto::b32;
|
|
#[cfg(feature = "store")]
|
|
use crate::message::{
|
|
accept_field, build_frontmatter, message_id, parse_frontmatter, seal, unseal, Opened,
|
|
};
|
|
#[cfg(feature = "store")]
|
|
use crate::store::{now, Store, Stored};
|
|
#[cfg(feature = "store")]
|
|
use crate::tcp::TcpTransport;
|
|
#[cfg(feature = "store")]
|
|
use crate::transport::{
|
|
FLAG_REQUESTS, LABEL_AUTH, LABEL_REGISTER, OP_AUTH, OP_DELETE, OP_FETCH, OP_REGISTER, OP_SEND,
|
|
};
|
|
|
|
/// A connection plus what the caller needs to know about how it was made:
|
|
/// the server's static key when no pin was checked, for the sec 4 warning.
|
|
#[cfg(feature = "store")]
|
|
pub struct Session {
|
|
transport: Box<dyn Transport>,
|
|
pub unpinned_static: Option<[u8; KEY_LEN]>,
|
|
}
|
|
|
|
#[cfg(feature = "store")]
|
|
impl Session {
|
|
pub fn transport(&mut self) -> &mut dyn Transport {
|
|
self.transport.as_mut()
|
|
}
|
|
|
|
pub fn close(&mut self) {
|
|
self.transport.close();
|
|
}
|
|
}
|
|
|
|
/// The account address with the host's dial hint applied, when the caller
|
|
/// resolved DNS itself: account-bound operations re-read the account, so
|
|
/// the hint rides along this way rather than in the stored row.
|
|
#[cfg(feature = "store")]
|
|
fn account_with_dial(store: &Store, dial: Option<&str>) -> Result<Address, Error> {
|
|
let addr = store
|
|
.account()?
|
|
.ok_or(Error::NotRegistered)?;
|
|
Ok(match dial {
|
|
Some(dial) => addr.with_dial(dial),
|
|
None => addr,
|
|
})
|
|
}
|
|
|
|
/// Opens a session, enforcing sec 4's rule about unpinned servers. TCP pins
|
|
/// the server's Noise static key; on RNS the destination hash is the pin
|
|
/// (RNS.md sec 13.4).
|
|
#[cfg(feature = "store")]
|
|
pub fn connect(
|
|
store: &Store,
|
|
addr: &Address,
|
|
require_pin: bool,
|
|
timeout: u64,
|
|
) -> Result<Session, Error> {
|
|
match addr.scheme {
|
|
Scheme::Tcp => {
|
|
let pinned = store.server_pin(&addr.host, addr.port)?;
|
|
if pinned.is_none() && require_pin {
|
|
return Err(Error::NotPinned {
|
|
host: addr.host.clone(),
|
|
port: addr.port,
|
|
});
|
|
}
|
|
// The dial hint lets a host that resolves DNS itself (an
|
|
// Android app) reach the server by IP while `host` keeps its
|
|
// identity roles.
|
|
let transport = TcpTransport::connect(
|
|
addr.dial.as_deref().unwrap_or(&addr.host),
|
|
addr.port,
|
|
pinned,
|
|
timeout,
|
|
)?;
|
|
// Pin enforcement lives here, one layer above the transport: the
|
|
// identity host names the trust object even when a dial hint
|
|
// routed the packets elsewhere (sec 4).
|
|
if let Some(pinned) = pinned {
|
|
if !ct_eq(&pinned, transport.server_static()) {
|
|
return Err(Error::PinMismatch {
|
|
host: addr.host.clone(),
|
|
port: addr.port,
|
|
pinned,
|
|
presented: *transport.server_static(),
|
|
});
|
|
}
|
|
}
|
|
let unpinned_static = if pinned.is_none() {
|
|
Some(*transport.server_static())
|
|
} else {
|
|
None
|
|
};
|
|
Ok(Session {
|
|
transport: Box::new(transport),
|
|
unpinned_static,
|
|
})
|
|
}
|
|
Scheme::Rns => {
|
|
// There is nothing to pin on this carrier (upstream spec 13.4):
|
|
// the destination hash is the pin, so `require_pin` decides
|
|
// nothing here and every session is as authenticated as a
|
|
// pinned TCP one.
|
|
#[cfg(feature = "rns")]
|
|
{
|
|
let _ = require_pin;
|
|
let destination = addr.destination()?;
|
|
let transport = crate::rns::transport::RnsTransport::connect(
|
|
&destination,
|
|
std::time::Duration::from_secs(timeout),
|
|
)?;
|
|
Ok(Session {
|
|
transport: Box::new(transport),
|
|
unpinned_static: None,
|
|
})
|
|
}
|
|
#[cfg(not(feature = "rns"))]
|
|
{
|
|
let _ = (addr, require_pin);
|
|
Err(Error::Other(
|
|
"smol+rns:// addresses need the rns feature; rebuild with --features rns".into(),
|
|
))
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
fn string_field(out: &mut Vec<u8>, text: &[u8]) -> Result<(), Error> {
|
|
if text.len() > 255 {
|
|
return Err(Error::Other("field longer than one length byte".into()));
|
|
}
|
|
out.push(text.len() as u8);
|
|
out.extend_from_slice(text);
|
|
Ok(())
|
|
}
|
|
|
|
/// RESOLVE: the username's current key and its rotation chain (sec 6.1).
|
|
pub fn resolve(
|
|
transport: &mut dyn Transport,
|
|
username: &str,
|
|
) -> Result<([u8; KEY_LEN], Vec<[u8; CERT_LEN]>), Error> {
|
|
let mut body = Vec::new();
|
|
string_field(&mut body, username.as_bytes())?;
|
|
let response = transport.request(OP_RESOLVE, &body)?;
|
|
expect_ok(response.status, &format!("resolving {username}"))?;
|
|
let mut r = Reader::new(&response.body);
|
|
let identity: [u8; KEY_LEN] = r.take(KEY_LEN)?.try_into().unwrap();
|
|
let chain_len = r.u8()? as usize;
|
|
let mut chain = Vec::with_capacity(chain_len);
|
|
for _ in 0..chain_len {
|
|
chain.push(r.take(CERT_LEN)?.try_into().unwrap());
|
|
}
|
|
r.done()?;
|
|
Ok((identity, chain))
|
|
}
|
|
|
|
/// Accepts a key change only when a signed chain leads from the key we hold
|
|
/// to the one the server now returns (sec 7). Both keys must sign each link;
|
|
/// a link predating the key we hold is skipped, a gap is not.
|
|
pub fn walk_chain(
|
|
username: &str,
|
|
pinned: &[u8; KEY_LEN],
|
|
current: &[u8; KEY_LEN],
|
|
chain: &[[u8; CERT_LEN]],
|
|
) -> bool {
|
|
if ct_eq(pinned, current) {
|
|
return true;
|
|
}
|
|
if chain.is_empty() || chain.len() > MAX_CHAIN {
|
|
return false;
|
|
}
|
|
let mut key: Vec<u8> = pinned.to_vec();
|
|
let mut started = false;
|
|
for cert_bytes in chain {
|
|
let cert = match RotationCert::from_bytes(cert_bytes) {
|
|
Ok(cert) => cert,
|
|
Err(_) => return false,
|
|
};
|
|
if !started {
|
|
if key != cert.old_pub {
|
|
continue; // a link predating the key we hold
|
|
}
|
|
started = true;
|
|
} else if key != cert.old_pub {
|
|
return false; // the chain is not continuous
|
|
}
|
|
if !cert.verify(username) {
|
|
return false;
|
|
}
|
|
key = cert.new_pub.to_vec();
|
|
}
|
|
started && ct_eq(&key, current)
|
|
}
|
|
|
|
/// How `trust_key` left the contact.
|
|
#[derive(PartialEq, Eq, Debug)]
|
|
pub enum TrustChange {
|
|
/// Trust on first use; `unverified` marks a session with no pin (sec 8).
|
|
New { unverified: bool },
|
|
/// A signed chain confirmed a key change (sec 7).
|
|
Rotated,
|
|
/// The key matches what we already hold.
|
|
None,
|
|
}
|
|
|
|
/// Resolves a contact and applies sec 8's trust rules. A key change without
|
|
/// a valid chain is an error here: the user confirms out of band and
|
|
/// imports, rather than clicking through.
|
|
#[cfg(feature = "store")]
|
|
pub fn trust_key(
|
|
store: &Store,
|
|
transport: &mut dyn Transport,
|
|
addr: &Address,
|
|
) -> Result<([u8; KEY_LEN], TrustChange), Error> {
|
|
let (identity, chain) = resolve(transport, &addr.user)?;
|
|
let known = store.contact(&addr.short())?;
|
|
match known {
|
|
None => {
|
|
store.save_contact(&addr.short(), &identity, false)?;
|
|
Ok((
|
|
identity,
|
|
TrustChange::New {
|
|
unverified: !transport.pinned(),
|
|
},
|
|
))
|
|
}
|
|
Some((old, _)) if ct_eq(&old, &identity) => Ok((identity, TrustChange::None)),
|
|
Some((old, verified)) => {
|
|
if walk_chain(&addr.user, &old, &identity, &chain) {
|
|
// The displaced key is the only local record that this
|
|
// contact rotated, so it is kept (sec 8) before the current
|
|
// one moves on. `until` is milliseconds, the unit the export
|
|
// format carries.
|
|
store.save_history(&addr.short(), &old, crate::store::now() * 1000)?;
|
|
store.save_contact(&addr.short(), &identity, verified)?;
|
|
Ok((identity, TrustChange::Rotated))
|
|
} else {
|
|
Err(Error::KeyChanged {
|
|
address: addr.short(),
|
|
known: old,
|
|
offered: identity,
|
|
})
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// AUTH (sec 4): sign the handshake hash and push the accept-token set. The
|
|
/// trailing block carries the mailbox's tokens (sec 5.8): `sync = 0` leaves
|
|
/// the stored set untouched, `sync = 1` replaces it with exactly what
|
|
/// follows.
|
|
#[cfg(feature = "store")]
|
|
pub fn authenticate(
|
|
transport: &mut dyn Transport,
|
|
username: &str,
|
|
account: &Account,
|
|
store: &Store,
|
|
) -> Result<u16, Error> {
|
|
let (sync, tokens) = store.token_set(account)?;
|
|
let mut body = Vec::new();
|
|
string_field(&mut body, username.as_bytes())?;
|
|
body.extend_from_slice(&account.me().pk());
|
|
body.extend_from_slice(
|
|
&account
|
|
.me()
|
|
.sign(&[LABEL_AUTH, &transport.bind().h[..]].concat()),
|
|
);
|
|
body.push(sync);
|
|
body.extend_from_slice(&(tokens.len() as u16).to_be_bytes());
|
|
for token in &tokens {
|
|
body.extend_from_slice(token);
|
|
}
|
|
let response = transport.request(OP_AUTH, &body)?;
|
|
expect_ok(response.status, "authentication")?;
|
|
let held = if response.body.len() >= 2 {
|
|
u16::from_be_bytes(response.body[..2].try_into().unwrap())
|
|
} else {
|
|
0
|
|
};
|
|
Ok(held)
|
|
}
|
|
|
|
/// What pushing the accept-token set achieved (sec 5.8).
|
|
#[cfg(feature = "store")]
|
|
pub enum Pushed {
|
|
/// Not registered yet; the set travels with the first fetch instead.
|
|
NotRegistered,
|
|
/// The server now holds this many tokens.
|
|
Held(u16),
|
|
}
|
|
|
|
/// An accept or a block only takes effect once the server holds the changed
|
|
/// set, so it is pushed now rather than at the next fetch.
|
|
#[cfg(feature = "store")]
|
|
pub fn push_tokens(
|
|
store: &Store,
|
|
account: &Account,
|
|
dial: Option<&str>,
|
|
timeout: u64,
|
|
) -> Result<Pushed, Error> {
|
|
let addr = match account_with_dial(store, dial) {
|
|
Ok(addr) => addr,
|
|
Err(Error::NotRegistered) => return Ok(Pushed::NotRegistered),
|
|
Err(e) => return Err(e),
|
|
};
|
|
let mut session = connect(store, &addr, true, timeout)?;
|
|
let held = authenticate(session.transport(), &addr.user, account, store)?;
|
|
session.close();
|
|
Ok(Pushed::Held(held))
|
|
}
|
|
|
|
/// REGISTER (sec 6.1): binds this identity to a username, with proof of
|
|
/// possession over the server's static key so the attestation cannot be
|
|
/// replayed to another server.
|
|
#[cfg(feature = "store")]
|
|
pub fn register(
|
|
store: &Store,
|
|
addr: &Address,
|
|
account: &Account,
|
|
invite: Option<&str>,
|
|
timeout: u64,
|
|
) -> Result<(), Error> {
|
|
let mut session = connect(store, addr, true, timeout)?;
|
|
let transport = session.transport();
|
|
let me = account.me();
|
|
let mut body = Vec::new();
|
|
string_field(&mut body, addr.user.as_bytes())?;
|
|
body.extend_from_slice(&me.pk());
|
|
body.extend_from_slice(
|
|
&me.sign(
|
|
&[
|
|
LABEL_REGISTER,
|
|
&transport.bind().server_static[..],
|
|
addr.user.as_bytes(),
|
|
&me.pk()[..],
|
|
]
|
|
.concat(),
|
|
),
|
|
);
|
|
let token = invite.unwrap_or("").as_bytes();
|
|
string_field(&mut body, token)?;
|
|
body.push(0); // no rotation certificate on a plain registration
|
|
let response = transport.request(OP_REGISTER, &body)?;
|
|
expect_ok(response.status, &format!("registering {}", addr.short()))?;
|
|
session.close();
|
|
store.set_account(addr)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// A completed rotation: the new key and the certificate the server now
|
|
/// holds, ready to be pushed to contacts as an ordinary message (sec 7).
|
|
#[cfg(feature = "store")]
|
|
pub struct Rotated {
|
|
pub new_pk: [u8; KEY_LEN],
|
|
pub cert: [u8; CERT_LEN],
|
|
}
|
|
|
|
/// Advances the rotation index and pushes the certificate (sec 7). The master
|
|
/// is untouched; only the index moves, and the superseded key stays
|
|
/// derivable from it (sec 2).
|
|
#[cfg(feature = "store")]
|
|
pub fn rotate(
|
|
store: &Store,
|
|
account: &Account,
|
|
dial: Option<&str>,
|
|
timeout: u64,
|
|
) -> Result<Rotated, Error> {
|
|
let addr = account_with_dial(store, dial)?;
|
|
if account.index() as usize >= MAX_CHAIN {
|
|
return Err(Error::ChainLimit {
|
|
index: account.index() + 1,
|
|
});
|
|
}
|
|
let old = account.me();
|
|
let new = Identity::from_seed(crate::account::identity_seed(
|
|
account.master(),
|
|
account.index() + 1,
|
|
));
|
|
let cert = RotationCert::build(old, &new, &addr.user, now());
|
|
|
|
let mut session = connect(store, &addr, true, timeout)?;
|
|
let transport = session.transport();
|
|
let mut body = Vec::new();
|
|
string_field(&mut body, addr.user.as_bytes())?;
|
|
body.extend_from_slice(&new.pk());
|
|
body.extend_from_slice(
|
|
&new.sign(
|
|
&[
|
|
LABEL_REGISTER,
|
|
&transport.bind().server_static[..],
|
|
addr.user.as_bytes(),
|
|
&new.pk()[..],
|
|
]
|
|
.concat(),
|
|
),
|
|
);
|
|
body.push(0); // no invite token on a rotation
|
|
body.push(CERT_LEN as u8);
|
|
body.extend_from_slice(&cert);
|
|
let response = transport.request(OP_REGISTER, &body)?;
|
|
expect_ok(response.status, "rotating")?;
|
|
session.close();
|
|
|
|
store.set_rotations(account.index() + 1)?;
|
|
Ok(Rotated {
|
|
new_pk: new.pk(),
|
|
cert,
|
|
})
|
|
}
|
|
|
|
/// Recovers the rotation index, and with it every superseded key, from the
|
|
/// master alone (sec 2). The accepted set is gone, so sync is disabled until
|
|
/// the user rebuilds it: an empty set must not replace the server's (sec 4).
|
|
/// The pieces are public for hosts that want intermediate UI states:
|
|
/// `connect` + `resolve`, `find_rotation_index`, then the four `set_*`
|
|
/// calls below.
|
|
#[cfg(feature = "store")]
|
|
pub fn restore(store: &Store, master: &[u8; KEY_LEN], addr: &Address, timeout: u64) -> Result<u32, Error> {
|
|
let mut session = connect(store, addr, false, timeout)?;
|
|
let (identity, _) = resolve(session.transport(), &addr.user)?;
|
|
session.close();
|
|
|
|
let Some(index) = crate::account::find_rotation_index(master, &identity) else {
|
|
return Err(Error::Other(format!(
|
|
"the key bound to {} is not derived from this master within {MAX_CHAIN} rotations",
|
|
addr.short()
|
|
)));
|
|
};
|
|
store.set_account(addr)?;
|
|
store.set_cursor(0, &[0u8; ID_LEN])?;
|
|
store.set_rotations(index)?;
|
|
store.set_sync_ok(false)?;
|
|
Ok(index)
|
|
}
|
|
|
|
/// One message to send, as the CLI assembled it.
|
|
#[cfg(feature = "store")]
|
|
pub struct SendDraft<'a> {
|
|
pub address: &'a Address,
|
|
pub text: String,
|
|
pub subject: Option<&'a str>,
|
|
/// An id this message replies to, as 64 hex characters.
|
|
pub reply_to: Option<&'a str>,
|
|
/// Extra `Key: value` frontmatter fields.
|
|
pub headers: &'a [(&'a str, &'a str)],
|
|
/// Omit the Reply-To field carrying this address (sec 5.7).
|
|
pub anonymous: bool,
|
|
pub no_pad: bool,
|
|
}
|
|
|
|
/// What `send` did, for the caller to report.
|
|
#[cfg(feature = "store")]
|
|
pub struct Sent {
|
|
pub id: [u8; ID_LEN],
|
|
pub bytes: usize,
|
|
/// An accept token for their mailbox was attached (sec 5.8).
|
|
pub token_used: bool,
|
|
/// How the recipient key was obtained, if it was learned now.
|
|
pub change: TrustChange,
|
|
/// A non-fatal oddity worth showing the user.
|
|
pub warning: Option<String>,
|
|
}
|
|
|
|
/// Seals and delivers one message (sec 5, sec 6.1). Recipient selection
|
|
/// prefers a key we already trust: a self-certifying address, then a stored
|
|
/// contact, and only then RESOLVE with trust on first use.
|
|
#[cfg(feature = "store")]
|
|
pub fn send(
|
|
store: &Store,
|
|
account: &Account,
|
|
draft: &SendDraft<'_>,
|
|
timeout: u64,
|
|
) -> Result<Sent, Error> {
|
|
let addr = draft.address;
|
|
let me = account.me();
|
|
|
|
let (recipient, change) = if let Some(key) = addr.identity {
|
|
store.save_contact(&addr.short(), &key, true)?;
|
|
(key, TrustChange::None)
|
|
} else if let Some((known, _)) = store.contact(&addr.short())? {
|
|
(known, TrustChange::None)
|
|
} else {
|
|
let mut session = connect(store, addr, false, timeout)?;
|
|
let (key, change) = trust_key(store, session.transport(), addr)?;
|
|
session.close();
|
|
(key, change)
|
|
};
|
|
|
|
// Frontmatter, in the order the reference client emits it.
|
|
let mut fields: Vec<(String, String)> = Vec::new();
|
|
if let Some(subject) = draft.subject {
|
|
fields.push(("Subject".into(), subject.to_string()));
|
|
}
|
|
if let Some(reply_to) = draft.reply_to {
|
|
fields.push(("In-Reply-To".into(), reply_to.to_string()));
|
|
}
|
|
for (key, value) in draft.headers {
|
|
fields.push(((*key).to_string(), (*value).to_string()));
|
|
}
|
|
// sec 5.7: a signed reply address lets a first-time recipient answer us.
|
|
if let Some(home) = store.account()? {
|
|
if !draft.anonymous {
|
|
fields.push(("Reply-To".into(), home.uri(&me.pk())));
|
|
}
|
|
}
|
|
// sec 5.8: hand an accepted correspondent the token for our own mailbox.
|
|
if let Some(identity) = store.accepted_identity(&addr.short())? {
|
|
fields.push(("Accept".into(), b32(&account.token_for(&identity))));
|
|
}
|
|
let borrowed: Vec<(&str, &str)> = fields
|
|
.iter()
|
|
.map(|(k, v)| (k.as_str(), v.as_str()))
|
|
.collect();
|
|
let body = build_frontmatter(&borrowed, &draft.text).into_bytes();
|
|
|
|
let envelope = seal(me, &recipient, &body, now(), !draft.no_pad)?;
|
|
let mid = message_id(&envelope);
|
|
// sec 5.6: the sent copy is sealed to ourselves, not to the recipient —
|
|
// the recipient's envelope uses an ephemeral key that is gone, so the
|
|
// only envelope its sender can ever open again is one sealed to their
|
|
// own long-term key. The stored id stays the delivered message's.
|
|
let self_copy = seal(me, &me.pk(), &body, now(), !draft.no_pad)?;
|
|
|
|
// sec 5.8: our token for their mailbox, if they have given us one.
|
|
let mac = match store.token_of(&addr.short())? {
|
|
Some(token) => hmac_sha256(&token, &[LABEL_MAC, &mid[..]].concat()).to_vec(),
|
|
None => Vec::new(),
|
|
};
|
|
|
|
let mut session = connect(store, addr, false, timeout)?;
|
|
let mut wire = vec![mac.len() as u8];
|
|
wire.extend_from_slice(&mac);
|
|
wire.extend_from_slice(&envelope);
|
|
let response = session.transport().request(OP_SEND, &wire)?;
|
|
expect_ok(response.status, &format!("sending to {}", addr.short()))?;
|
|
session.close();
|
|
// The server echoes the id; it is derived, so ours is authoritative.
|
|
let warning = if !response.body.is_empty() && response.body != mid.to_vec() {
|
|
Some("server returned an id we did not derive; it is not authoritative".to_string())
|
|
} else {
|
|
None
|
|
};
|
|
|
|
// sec 5.6: the ephemeral is gone, so keep a copy sealed to ourselves.
|
|
store.store_sent(&mid, &addr.short(), &self_copy, now())?;
|
|
Ok(Sent {
|
|
id: mid,
|
|
bytes: envelope.len(),
|
|
token_used: !mac.is_empty(),
|
|
change,
|
|
warning,
|
|
})
|
|
}
|
|
|
|
/// One rejected fetch record: the id and why it did not enter the inbox.
|
|
#[cfg(feature = "store")]
|
|
pub struct Rejected {
|
|
pub id: [u8; ID_LEN],
|
|
pub reason: String,
|
|
}
|
|
|
|
/// What `fetch` did, for the caller to report.
|
|
#[cfg(feature = "store")]
|
|
pub struct Fetched {
|
|
pub total: u32,
|
|
pub stored: u32,
|
|
pub rejected: Vec<Rejected>,
|
|
/// Messages were deleted from the server rather than kept behind a
|
|
/// cursor (the default).
|
|
pub acknowledged: bool,
|
|
/// `cancel` was raised mid-fetch; the summary covers what completed.
|
|
/// Every envelope the summary counted is accounted for locally, so
|
|
/// fetching again continues rather than repeats.
|
|
pub cancelled: bool,
|
|
}
|
|
|
|
/// The knobs a host with a UI needs on `fetch_with`. Progress needs no
|
|
/// callback: each envelope is committed to the `Store` as it is verified, so
|
|
/// a concurrent reader (the store is `Send + Sync`) can list incrementally.
|
|
#[cfg(feature = "store")]
|
|
pub struct FetchOptions<'a> {
|
|
/// Remember the cursor instead of acknowledging with DELETE.
|
|
pub keep: bool,
|
|
/// Forget the cursor and page again.
|
|
pub reset: bool,
|
|
/// Network timeout in seconds.
|
|
pub timeout: u64,
|
|
/// Dial this address instead of the account host, for hosts that
|
|
/// resolve DNS themselves; the host keeps its identity roles.
|
|
pub dial: Option<&'a str>,
|
|
/// Checked between pages and before each envelope; when raised,
|
|
/// `fetch_with` stops and returns the partial summary with `cancelled`
|
|
/// set. What was processed is accounted for locally — cursor moved or
|
|
/// messages deleted — so fetching again continues rather than repeats.
|
|
pub cancel: Option<&'a std::sync::atomic::AtomicBool>,
|
|
}
|
|
|
|
#[cfg(feature = "store")]
|
|
impl FetchOptions<'_> {
|
|
/// The defaults a plain fetch uses.
|
|
fn new(keep: bool, reset: bool, timeout: u64) -> Self {
|
|
FetchOptions {
|
|
keep,
|
|
reset,
|
|
timeout,
|
|
dial: None,
|
|
cancel: None,
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Retrieves, verifies and stores mail (sec 6.1), paging forward by
|
|
/// `(received_at, id)` until a response is empty. Acknowledging deletes what
|
|
/// it takes, so it always pages from the start; `keep` instead remembers the
|
|
/// cursor. A message that fails verification is left on the server, so a
|
|
/// client-side bug cannot lose mail.
|
|
#[cfg(feature = "store")]
|
|
pub fn fetch(
|
|
store: &Store,
|
|
account: &Account,
|
|
keep: bool,
|
|
reset: bool,
|
|
timeout: u64,
|
|
) -> Result<Fetched, Error> {
|
|
fetch_with(store, account, FetchOptions::new(keep, reset, timeout))
|
|
}
|
|
|
|
/// `fetch` with the full option set, for hosts that need cancellation.
|
|
#[cfg(feature = "store")]
|
|
pub fn fetch_with(
|
|
store: &Store,
|
|
account: &Account,
|
|
opts: FetchOptions<'_>,
|
|
) -> Result<Fetched, Error> {
|
|
let addr = account_with_dial(store, opts.dial)?;
|
|
if opts.reset {
|
|
store.set_cursor(0, &[0u8; ID_LEN])?;
|
|
}
|
|
let (mut after_time, mut after_id) = if opts.keep {
|
|
store.cursor()?
|
|
} else {
|
|
(0, [0u8; ID_LEN])
|
|
};
|
|
|
|
let mut summary = Fetched {
|
|
total: 0,
|
|
stored: 0,
|
|
rejected: Vec::new(),
|
|
acknowledged: !opts.keep,
|
|
cancelled: false,
|
|
};
|
|
let cancelled = || opts.cancel.is_some_and(|c| c.load(std::sync::atomic::Ordering::Relaxed));
|
|
let mut session = connect(store, &addr, true, opts.timeout)?;
|
|
let transport = session.transport();
|
|
authenticate(transport, &addr.user, account, store)?;
|
|
|
|
loop {
|
|
if cancelled() {
|
|
summary.cancelled = true;
|
|
break;
|
|
}
|
|
let mut body = Vec::with_capacity(8 + ID_LEN);
|
|
body.extend_from_slice(&after_time.to_be_bytes());
|
|
body.extend_from_slice(&after_id);
|
|
let response = transport.request(OP_FETCH, &body)?;
|
|
expect_ok(response.status, "fetching")?;
|
|
let mut r = Reader::new(&response.body);
|
|
let count = r.u16()? as usize;
|
|
if count == 0 {
|
|
break;
|
|
}
|
|
let mut acked: Vec<[u8; ID_LEN]> = Vec::new();
|
|
let mut page_cancelled = false;
|
|
for _ in 0..count {
|
|
// Checked before each envelope as well as between pages, so a
|
|
// cancel lands within one envelope; the cursor and acks below
|
|
// then cover exactly the envelopes that were processed.
|
|
if cancelled() {
|
|
page_cancelled = true;
|
|
summary.cancelled = true;
|
|
break;
|
|
}
|
|
let mid: [u8; ID_LEN] = r.take(ID_LEN)?.try_into().unwrap();
|
|
let received_at = r.i64()?;
|
|
let flags = r.u8()?;
|
|
let envelope_len = r.u32()? as usize;
|
|
let envelope = r.take(envelope_len)?.to_vec();
|
|
after_time = received_at;
|
|
after_id = mid;
|
|
summary.total += 1;
|
|
|
|
let opened = match verify_and_open(account, &mid, &envelope) {
|
|
Ok(opened) => opened,
|
|
Err(reason) => {
|
|
summary.rejected.push(Rejected {
|
|
id: mid,
|
|
reason: reason.to_string(),
|
|
});
|
|
continue;
|
|
}
|
|
};
|
|
// A message we have already had once is not stored again, even
|
|
// if we deleted it locally in the meantime (sec 10).
|
|
if !store.seen(&mid)? {
|
|
let requests_tier = flags & FLAG_REQUESTS != 0;
|
|
store.store_inbox(&mid, &envelope, received_at, requests_tier, opts.keep)?;
|
|
learn_token(store, &opened)?;
|
|
summary.stored += 1;
|
|
}
|
|
acked.push(mid);
|
|
}
|
|
// A cancelled page is not fully parsed, so trailing bytes are
|
|
// expected and not an error; everything processed is accounted for.
|
|
if !page_cancelled {
|
|
r.done()?;
|
|
}
|
|
if opts.keep {
|
|
store.set_cursor(after_time, &after_id)?;
|
|
} else if !acked.is_empty() {
|
|
delete_ids(transport, &acked)?;
|
|
}
|
|
if page_cancelled {
|
|
break;
|
|
}
|
|
}
|
|
session.close();
|
|
if !opts.keep {
|
|
store.set_cursor(0, &[0u8; ID_LEN])?;
|
|
}
|
|
Ok(summary)
|
|
}
|
|
|
|
#[cfg(feature = "store")]
|
|
fn verify_and_open(
|
|
account: &Account,
|
|
mid: &[u8; ID_LEN],
|
|
envelope: &[u8],
|
|
) -> Result<Opened, Error> {
|
|
if message_id(envelope) != *mid {
|
|
return Err(Error::Other("id does not match the envelope".into()));
|
|
}
|
|
unseal(account.keys(), envelope, now())
|
|
}
|
|
|
|
/// Files the accept token a verified payload carried, under the address that
|
|
/// signed it (sec 5.8). The signature has already been checked by `unseal`,
|
|
/// so the attribution is the signer's own claim.
|
|
#[cfg(feature = "store")]
|
|
fn learn_token(store: &Store, opened: &Opened) -> Result<(), Error> {
|
|
let text = String::from_utf8_lossy(&opened.body).into_owned();
|
|
let (fields, _) = parse_frontmatter(&text);
|
|
let Some(token) = accept_field(&fields) else {
|
|
return Ok(());
|
|
};
|
|
let reply_to = fields.get("reply-to").map(String::as_str);
|
|
match store.address_of(&opened.sender, reply_to)? {
|
|
// No address to send to, so no use for a token.
|
|
None => Ok(()),
|
|
Some(address) => store.save_token(&address, &token),
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "store")]
|
|
fn delete_ids(transport: &mut dyn Transport, ids: &[[u8; ID_LEN]]) -> Result<u16, Error> {
|
|
let mut body = Vec::with_capacity(2 + ids.len() * ID_LEN);
|
|
body.extend_from_slice(&(ids.len() as u16).to_be_bytes());
|
|
for id in ids {
|
|
body.extend_from_slice(id);
|
|
}
|
|
let response = transport.request(OP_DELETE, &body)?;
|
|
expect_ok(response.status, "acknowledging")?;
|
|
let removed = if response.body.len() >= 2 {
|
|
u16::from_be_bytes(response.body[..2].try_into().unwrap())
|
|
} else {
|
|
0
|
|
};
|
|
Ok(removed)
|
|
}
|
|
|
|
/// DELETE (sec 6.1) over an authenticated session: remove ids from the
|
|
/// server explicitly. Unknown ids are not an error.
|
|
#[cfg(feature = "store")]
|
|
pub fn delete(
|
|
store: &Store,
|
|
account: &Account,
|
|
ids: &[[u8; ID_LEN]],
|
|
dial: Option<&str>,
|
|
timeout: u64,
|
|
) -> Result<u16, Error> {
|
|
let addr = account_with_dial(store, dial)?;
|
|
let mut session = connect(store, &addr, true, timeout)?;
|
|
let transport = session.transport();
|
|
authenticate(transport, &addr.user, account, store)?;
|
|
let removed = delete_ids(transport, ids)?;
|
|
session.close();
|
|
Ok(removed)
|
|
}
|
|
|
|
/// An opened, described message for listing and reading.
|
|
#[cfg(feature = "store")]
|
|
pub struct Described {
|
|
pub sender: [u8; KEY_LEN],
|
|
pub time: i64,
|
|
pub fields: BTreeMap<String, String>,
|
|
pub text: String,
|
|
/// The contact address we know the signer by, or a fingerprint.
|
|
pub from: String,
|
|
pub subject: String,
|
|
}
|
|
|
|
#[cfg(feature = "store")]
|
|
impl Described {
|
|
/// Whether this payload carried its signer's accept token (sec 5.8):
|
|
/// machinery, not content.
|
|
pub fn carries_token(&self) -> bool {
|
|
self.fields.contains_key("accept")
|
|
}
|
|
}
|
|
|
|
/// Opens one sealed message and splits its body into frontmatter and text.
|
|
#[cfg(feature = "store")]
|
|
pub fn describe(store: &Store, account: &Account, stored: &Stored) -> Result<Described, Error> {
|
|
let opened = unseal(account.keys(), &stored.envelope, now())?;
|
|
let text = String::from_utf8_lossy(&opened.body).into_owned();
|
|
let (fields, body_text) = parse_frontmatter(&text);
|
|
let from = store
|
|
.address_of(&opened.sender, fields.get("reply-to").map(String::as_str))?
|
|
.unwrap_or_else(|| format!("<{}…>", &b32(&opened.sender)[..20]));
|
|
Ok(Described {
|
|
sender: opened.sender,
|
|
time: opened.time,
|
|
subject: fields.get("subject").cloned().unwrap_or_default(),
|
|
fields,
|
|
text: body_text,
|
|
from,
|
|
})
|
|
}
|
|
|
|
/// The accept MAC a sender attaches to SEND (sec 5.8), exposed for tests.
|
|
pub fn accept_mac(token: &[u8; TOKEN_LEN], mid: &[u8; ID_LEN]) -> [u8; 32] {
|
|
hmac_sha256(token, &[LABEL_MAC, mid].concat())
|
|
}
|
|
|
|
/// Proof of possession over a server's static key (sec 6.1), for tests.
|
|
#[cfg(all(test, feature = "store"))]
|
|
fn register_signed(server_static: &[u8], username: &str, identity: &[u8]) -> Vec<u8> {
|
|
[LABEL_REGISTER, server_static, username.as_bytes(), identity].concat()
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::account::{identity_seed, verify_sig, Account, Identity, RotationCert};
|
|
use crate::transport::{TransportBindValues, SIG_LEN};
|
|
|
|
fn master() -> [u8; 32] {
|
|
core::array::from_fn(|i| i as u8)
|
|
}
|
|
|
|
fn identity(n: u32) -> Identity {
|
|
Identity::from_seed(identity_seed(&master(), n))
|
|
}
|
|
|
|
#[test]
|
|
fn walk_chain_accepts_signed_progressions() {
|
|
let old = identity(0);
|
|
let new = identity(1);
|
|
let cert = RotationCert::build(&old, &new, "alice", 1700000001);
|
|
assert!(walk_chain("alice", &old.pk(), &old.pk(), &[]));
|
|
assert!(walk_chain("alice", &old.pk(), &new.pk(), &[cert]));
|
|
// The username is covered, so the same cert proves nothing for bob.
|
|
assert!(!walk_chain("bob", &old.pk(), &new.pk(), &[cert]));
|
|
// No chain, no acceptance.
|
|
assert!(!walk_chain("alice", &old.pk(), &new.pk(), &[]));
|
|
}
|
|
|
|
#[test]
|
|
fn walk_chain_requires_continuity_and_terminates_at_current() {
|
|
let keys: Vec<Identity> = (0..4).map(identity).collect();
|
|
let cert01 = RotationCert::build(&keys[0], &keys[1], "alice", 1);
|
|
let cert12 = RotationCert::build(&keys[1], &keys[2], "alice", 2);
|
|
let cert23 = RotationCert::build(&keys[2], &keys[3], "alice", 3);
|
|
// A full chain from 0 to 3.
|
|
assert!(walk_chain(
|
|
"alice",
|
|
&keys[0].pk(),
|
|
&keys[3].pk(),
|
|
&[cert01, cert12, cert23]
|
|
));
|
|
// Holding key 1: the earlier link is skipped, the rest must continue.
|
|
assert!(walk_chain(
|
|
"alice",
|
|
&keys[1].pk(),
|
|
&keys[3].pk(),
|
|
&[cert01, cert12, cert23]
|
|
));
|
|
// A gap between the held key and the chain never validates.
|
|
assert!(!walk_chain(
|
|
"alice",
|
|
&keys[0].pk(),
|
|
&keys[3].pk(),
|
|
&[cert12, cert23]
|
|
));
|
|
// The chain must end at the key RESOLVE returned, not merely agree
|
|
// partway.
|
|
assert!(!walk_chain(
|
|
"alice",
|
|
&keys[0].pk(),
|
|
&keys[2].pk(),
|
|
&[cert01, cert12, cert23]
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn walk_chain_enforces_the_length_limit() {
|
|
let keys: Vec<Identity> = (0..=16).map(identity).collect();
|
|
let certs: Vec<[u8; CERT_LEN]> = (0..16)
|
|
.map(|n| RotationCert::build(&keys[n], &keys[n + 1], "alice", n as i64))
|
|
.collect();
|
|
// Sixteen links walk the whole chain from the first key.
|
|
assert!(walk_chain("alice", &keys[0].pk(), &keys[16].pk(), &certs));
|
|
// Seventeen links is over-long and needs user confirmation, whatever
|
|
// key we hold (sec 7).
|
|
let mut over = certs.clone();
|
|
let extra = identity(17);
|
|
over.push(RotationCert::build(&keys[16], &extra, "alice", 17));
|
|
assert!(!walk_chain("alice", &keys[0].pk(), &extra.pk(), &over));
|
|
assert!(!walk_chain("alice", &keys[1].pk(), &extra.pk(), &over));
|
|
}
|
|
|
|
#[test]
|
|
fn accept_mac_matches_the_reference_vector() {
|
|
// The reference client's own vector: token_for(pk_1) over the id of
|
|
// the reference envelope (see the message module's vectors).
|
|
let v = crate::vectors::load();
|
|
let account = Account::new(master(), 0).unwrap();
|
|
let token = account.token_for(&identity(1).pk());
|
|
let mid: [u8; 32] = crate::crypto::unb32(&v.accept_mac.message_id_b32)
|
|
.unwrap()
|
|
.try_into()
|
|
.unwrap();
|
|
assert_eq!(
|
|
data_encoding::HEXLOWER.encode(&accept_mac(&token, &mid)),
|
|
v.accept_mac.mac_hex
|
|
);
|
|
}
|
|
|
|
#[cfg(feature = "store")]
|
|
#[test]
|
|
fn bind_values_feed_the_auth_and_register_signatures() {
|
|
let bind = TransportBindValues::tcp(&[2u8; 32], &[7u8; 32]).unwrap();
|
|
let account = Account::new(master(), 0).unwrap();
|
|
let me = account.me();
|
|
let auth = me.sign(&[LABEL_AUTH, &bind.h[..]].concat());
|
|
let register = me.sign(®ister_signed(&bind.server_static, "alice", &me.pk()));
|
|
assert_eq!(auth.len(), SIG_LEN);
|
|
assert!(verify_sig(
|
|
&me.pk(),
|
|
&auth,
|
|
&[LABEL_AUTH, &bind.h[..]].concat()
|
|
));
|
|
assert!(verify_sig(
|
|
&me.pk(),
|
|
®ister,
|
|
®ister_signed(&bind.server_static, "alice", &me.pk())
|
|
));
|
|
}
|
|
}
|