2026-09-26 11:52:22 +03:00
|
|
|
//! Per-connection dispatch and the six wire operations.
|
|
|
|
|
|
2026-09-27 09:39:01 +03:00
|
|
|
use crate::crypto::{b32, ct_eq, hmac_sha256, message_id, verify};
|
2026-09-26 11:52:22 +03:00
|
|
|
use crate::proto::{
|
|
|
|
|
valid_username, ProtocolError, Reader, AUTH_FAILED, AUTH_REQUIRED, BAD_VERSION, CERT_LEN,
|
|
|
|
|
ENVELOPE_MAGIC, ENVELOPE_MIN, ENVELOPE_VERSION, FETCH_BUDGET, ID_LEN, KEY_LEN, LABEL_AUTH,
|
2026-09-27 09:39:01 +03:00
|
|
|
LABEL_MAC, LABEL_REGISTER, LABEL_ROTATE, MAC_LEN, MALFORMED, MAX_CHAIN, NOT_PERMITTED, OK,
|
|
|
|
|
OP_AUTH, OP_DELETE, OP_FETCH, OP_REGISTER, OP_RESOLVE, OP_SEND, QUOTA_EXCEEDED, RATE_LIMITED,
|
|
|
|
|
TOKEN_LEN, TOO_LARGE, UNKNOWN_USER,
|
2026-09-26 11:52:22 +03:00
|
|
|
};
|
|
|
|
|
use crate::ratelimit::RateLimiter;
|
|
|
|
|
use crate::store::Store;
|
|
|
|
|
|
|
|
|
|
pub struct ServerConfig {
|
|
|
|
|
pub max_envelope: usize,
|
2026-09-27 09:39:01 +03:00
|
|
|
pub main_quota: i64,
|
|
|
|
|
pub requests_quota: i64,
|
|
|
|
|
pub max_tokens: u16,
|
2026-09-26 11:52:22 +03:00
|
|
|
pub invite_token: Option<Vec<u8>>,
|
2026-09-27 09:39:01 +03:00
|
|
|
pub server_static: [u8; KEY_LEN],
|
2026-09-26 11:52:22 +03:00
|
|
|
pub conn_limiter: RateLimiter,
|
|
|
|
|
pub send_limiter: RateLimiter,
|
2026-09-27 09:39:01 +03:00
|
|
|
pub token_limiter: RateLimiter,
|
2026-09-26 11:52:22 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// A parse failure (-> MALFORMED) or a storage failure (-> INTERNAL_ERROR).
|
|
|
|
|
pub enum HandlerError {
|
|
|
|
|
Protocol(ProtocolError),
|
|
|
|
|
Store(rusqlite::Error),
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl From<ProtocolError> for HandlerError {
|
|
|
|
|
fn from(e: ProtocolError) -> Self {
|
|
|
|
|
HandlerError::Protocol(e)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl From<rusqlite::Error> for HandlerError {
|
|
|
|
|
fn from(e: rusqlite::Error) -> Self {
|
|
|
|
|
HandlerError::Store(e)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type OpResult = Result<(u8, Vec<u8>), HandlerError>;
|
|
|
|
|
|
|
|
|
|
pub struct Session<'a> {
|
|
|
|
|
config: &'a ServerConfig,
|
|
|
|
|
store: Store,
|
|
|
|
|
peer_ip: String,
|
|
|
|
|
handshake_hash: Vec<u8>,
|
|
|
|
|
username: Option<String>,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<'a> Session<'a> {
|
|
|
|
|
pub fn new(config: &'a ServerConfig, store: Store, peer_ip: String, handshake_hash: Vec<u8>) -> Self {
|
|
|
|
|
Session {
|
|
|
|
|
config,
|
|
|
|
|
store,
|
|
|
|
|
peer_ip,
|
|
|
|
|
handshake_hash,
|
|
|
|
|
username: None,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Dispatches one frame, always producing a status to send back, never
|
|
|
|
|
/// panicking or propagating errors to the caller: a bad frame or a
|
|
|
|
|
/// storage error both become a response, and the caller decides
|
|
|
|
|
/// separately whether to keep the connection open.
|
|
|
|
|
pub fn dispatch(&mut self, op: u8, body: &[u8]) -> (u8, Vec<u8>) {
|
|
|
|
|
if matches!(op, OP_FETCH | OP_DELETE) && self.username.is_none() {
|
|
|
|
|
return (AUTH_REQUIRED, Vec::new());
|
|
|
|
|
}
|
|
|
|
|
let mut r = Reader::new(body);
|
|
|
|
|
let result = match op {
|
|
|
|
|
OP_AUTH => self.op_auth(&mut r),
|
|
|
|
|
OP_RESOLVE => self.op_resolve(&mut r),
|
|
|
|
|
OP_SEND => self.op_send(&mut r),
|
|
|
|
|
OP_FETCH => self.op_fetch(&mut r),
|
|
|
|
|
OP_DELETE => self.op_delete(&mut r),
|
|
|
|
|
OP_REGISTER => self.op_register(&mut r),
|
|
|
|
|
_ => return (MALFORMED, Vec::new()),
|
|
|
|
|
};
|
|
|
|
|
match result {
|
|
|
|
|
Ok(response) => response,
|
|
|
|
|
Err(HandlerError::Protocol(e)) => {
|
|
|
|
|
log::info!("bad body from {}: {}", self.peer_ip, e);
|
|
|
|
|
(MALFORMED, Vec::new())
|
|
|
|
|
}
|
|
|
|
|
Err(HandlerError::Store(e)) => {
|
|
|
|
|
log::error!("storage error from {}: {}", self.peer_ip, e);
|
|
|
|
|
(crate::proto::INTERNAL_ERROR, Vec::new())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn read_str(r: &mut Reader) -> Result<String, ProtocolError> {
|
|
|
|
|
let len = r.u8()? as usize;
|
|
|
|
|
let bytes = r.take(len)?;
|
|
|
|
|
std::str::from_utf8(bytes)
|
|
|
|
|
.map(str::to_string)
|
|
|
|
|
.map_err(|_| ProtocolError::new("invalid utf-8"))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn op_auth(&mut self, r: &mut Reader) -> OpResult {
|
|
|
|
|
let username = Self::read_str(r)?;
|
|
|
|
|
let identity = r.take(KEY_LEN)?.to_vec();
|
|
|
|
|
let signature = r.take(64)?.to_vec();
|
2026-09-27 09:39:01 +03:00
|
|
|
let sync = r.u8()?;
|
|
|
|
|
let count = r.u16()? as usize;
|
|
|
|
|
let mut tokens = Vec::with_capacity(count);
|
|
|
|
|
for _ in 0..count {
|
|
|
|
|
tokens.push(r.take(TOKEN_LEN)?.to_vec());
|
|
|
|
|
}
|
2026-09-26 11:52:22 +03:00
|
|
|
r.done()?;
|
|
|
|
|
|
2026-09-27 09:39:01 +03:00
|
|
|
// sync = 0 leaves the stored set untouched, so it carries no tokens.
|
|
|
|
|
if sync > 1 || (sync == 0 && count != 0) {
|
|
|
|
|
return Ok((MALFORMED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-26 11:52:22 +03:00
|
|
|
let bound = self.store.identity_of(&username)?;
|
|
|
|
|
// A wrong username and a wrong signature are both AUTH_FAILED: telling
|
|
|
|
|
// them apart would turn this into an account-existence oracle.
|
|
|
|
|
if bound.as_deref() != Some(identity.as_slice()) {
|
|
|
|
|
return Ok((AUTH_FAILED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let mut msg = LABEL_AUTH.to_vec();
|
|
|
|
|
msg.extend_from_slice(&self.handshake_hash);
|
|
|
|
|
if !verify(&identity, &signature, &msg) {
|
|
|
|
|
return Ok((AUTH_FAILED, Vec::new()));
|
|
|
|
|
}
|
2026-09-27 09:39:01 +03:00
|
|
|
|
|
|
|
|
if sync == 1 {
|
|
|
|
|
if count > self.config.max_tokens as usize {
|
|
|
|
|
// Leaves the session unauthenticated, per SPEC.md sec 4.
|
|
|
|
|
return Ok((TOO_LARGE, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
self.store.set_tokens(&username, &tokens)?;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let accepted = self.store.token_count(&username)?;
|
2026-09-26 11:52:22 +03:00
|
|
|
self.username = Some(username);
|
2026-09-27 09:39:01 +03:00
|
|
|
Ok((OK, accepted.to_be_bytes().to_vec()))
|
2026-09-26 11:52:22 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn op_resolve(&mut self, r: &mut Reader) -> OpResult {
|
|
|
|
|
let username = Self::read_str(r)?;
|
|
|
|
|
r.done()?;
|
|
|
|
|
let identity = match self.store.identity_of(&username)? {
|
|
|
|
|
Some(id) => id,
|
|
|
|
|
None => return Ok((UNKNOWN_USER, Vec::new())),
|
|
|
|
|
};
|
|
|
|
|
let chain = self.store.chain(&username)?;
|
|
|
|
|
let mut out = identity;
|
|
|
|
|
out.push(chain.len() as u8);
|
|
|
|
|
for cert in chain {
|
|
|
|
|
out.extend_from_slice(&cert);
|
|
|
|
|
}
|
|
|
|
|
Ok((OK, out))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn op_send(&mut self, r: &mut Reader) -> OpResult {
|
2026-09-27 09:39:01 +03:00
|
|
|
let mac_len = r.u8()? as usize;
|
|
|
|
|
if mac_len != 0 && mac_len != MAC_LEN {
|
|
|
|
|
return Ok((MALFORMED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let mac = r.take(mac_len)?.to_vec();
|
2026-09-26 11:52:22 +03:00
|
|
|
let envelope = r.rest().to_vec();
|
2026-09-27 09:39:01 +03:00
|
|
|
|
2026-09-26 11:52:22 +03:00
|
|
|
if !self.config.send_limiter.allow(&self.peer_ip) {
|
|
|
|
|
return Ok((RATE_LIMITED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
if envelope.len() > self.config.max_envelope {
|
|
|
|
|
return Ok((TOO_LARGE, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
if envelope.len() < ENVELOPE_MIN || &envelope[..4] != ENVELOPE_MAGIC {
|
|
|
|
|
return Ok((MALFORMED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
if envelope[4] != ENVELOPE_VERSION {
|
|
|
|
|
return Ok((BAD_VERSION, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let recipient = &envelope[5..37];
|
|
|
|
|
let username = match self.store.username_for_key(recipient)? {
|
|
|
|
|
Some(u) => u,
|
|
|
|
|
None => return Ok((UNKNOWN_USER, Vec::new())),
|
|
|
|
|
};
|
2026-09-27 09:39:01 +03:00
|
|
|
|
|
|
|
|
// The ciphertext is never inspected; the server cannot read it.
|
|
|
|
|
let mid = message_id(&envelope);
|
|
|
|
|
|
|
|
|
|
// A matching accept token (SPEC.md sec 5.8) puts the envelope in the
|
|
|
|
|
// mailbox's main tier; anything else lands in the smaller, short-lived
|
|
|
|
|
// requests tier. A failed match is never reported to the sender.
|
|
|
|
|
let matched_token = if mac.len() == MAC_LEN {
|
|
|
|
|
let mut msg = LABEL_MAC.to_vec();
|
|
|
|
|
msg.extend_from_slice(&mid);
|
|
|
|
|
self.store
|
|
|
|
|
.tokens_of(&username)?
|
|
|
|
|
.into_iter()
|
|
|
|
|
.find(|t| ct_eq(&hmac_sha256(t, &msg), &mac))
|
|
|
|
|
} else {
|
|
|
|
|
None
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
if let Some(token) = &matched_token {
|
|
|
|
|
if !self.config.token_limiter.allow(&b32(token)) {
|
|
|
|
|
return Ok((RATE_LIMITED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let unsolicited = matched_token.is_none();
|
|
|
|
|
let tier_quota = if unsolicited {
|
|
|
|
|
self.config.requests_quota
|
|
|
|
|
} else {
|
|
|
|
|
self.config.main_quota
|
|
|
|
|
};
|
2026-09-26 11:52:22 +03:00
|
|
|
let keys = self.store.keys_of(&username)?;
|
2026-09-27 09:39:01 +03:00
|
|
|
let used = self.store.mailbox_bytes(&keys, unsolicited)?;
|
|
|
|
|
if used + envelope.len() as i64 > tier_quota {
|
2026-09-26 11:52:22 +03:00
|
|
|
return Ok((QUOTA_EXCEEDED, Vec::new()));
|
|
|
|
|
}
|
2026-09-27 09:39:01 +03:00
|
|
|
|
|
|
|
|
self.store
|
|
|
|
|
.store_message(&mid, recipient, &envelope, unsolicited)?;
|
2026-09-26 11:52:22 +03:00
|
|
|
Ok((OK, mid.to_vec()))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn op_fetch(&mut self, r: &mut Reader) -> OpResult {
|
2026-09-27 09:39:01 +03:00
|
|
|
let after_received_at = r.i64()?;
|
|
|
|
|
let after_id = r.take(ID_LEN)?.to_vec();
|
2026-09-26 11:52:22 +03:00
|
|
|
r.done()?;
|
|
|
|
|
let username = self.username.as_ref().expect("AUTH_REQUIRED gate above");
|
|
|
|
|
let keys = self.store.keys_of(username)?;
|
2026-09-27 09:39:01 +03:00
|
|
|
let records = self
|
|
|
|
|
.store
|
|
|
|
|
.pending(&keys, after_received_at, &after_id, FETCH_BUDGET)?;
|
2026-09-26 11:52:22 +03:00
|
|
|
|
|
|
|
|
let mut out = Vec::new();
|
|
|
|
|
out.extend_from_slice(&(records.len() as u16).to_be_bytes());
|
2026-09-27 09:39:01 +03:00
|
|
|
for (mid, received_at, unsolicited, envelope) in records {
|
2026-09-26 11:52:22 +03:00
|
|
|
out.extend_from_slice(&mid);
|
|
|
|
|
out.extend_from_slice(&received_at.to_be_bytes());
|
2026-09-27 09:39:01 +03:00
|
|
|
out.push(if unsolicited { 1 } else { 0 });
|
2026-09-26 11:52:22 +03:00
|
|
|
out.extend_from_slice(&(envelope.len() as u32).to_be_bytes());
|
|
|
|
|
out.extend_from_slice(&envelope);
|
|
|
|
|
}
|
|
|
|
|
Ok((OK, out))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn op_delete(&mut self, r: &mut Reader) -> OpResult {
|
|
|
|
|
let count = r.u16()? as usize;
|
|
|
|
|
let mut ids = Vec::with_capacity(count);
|
|
|
|
|
for _ in 0..count {
|
|
|
|
|
ids.push(r.take(ID_LEN)?.to_vec());
|
|
|
|
|
}
|
|
|
|
|
r.done()?;
|
|
|
|
|
let username = self.username.as_ref().expect("AUTH_REQUIRED gate above");
|
|
|
|
|
|
|
|
|
|
if ids.is_empty() {
|
|
|
|
|
return Ok((OK, 0u16.to_be_bytes().to_vec()));
|
|
|
|
|
}
|
|
|
|
|
// Scoped to the caller's own keys, so ids cannot be used to probe or
|
|
|
|
|
// delete another mailbox.
|
|
|
|
|
let keys = self.store.keys_of(username)?;
|
|
|
|
|
let removed = self.store.delete(&keys, &ids)?;
|
|
|
|
|
Ok((OK, (removed as u16).to_be_bytes().to_vec()))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn op_register(&mut self, r: &mut Reader) -> OpResult {
|
|
|
|
|
let username = Self::read_str(r)?;
|
|
|
|
|
let identity = r.take(KEY_LEN)?.to_vec();
|
2026-09-27 09:39:01 +03:00
|
|
|
let signature = r.take(64)?.to_vec();
|
2026-09-26 11:52:22 +03:00
|
|
|
let token_len = r.u8()? as usize;
|
|
|
|
|
let token = r.take(token_len)?.to_vec();
|
|
|
|
|
let cert_len = r.u8()? as usize;
|
|
|
|
|
let cert = r.take(cert_len)?.to_vec();
|
|
|
|
|
r.done()?;
|
|
|
|
|
|
|
|
|
|
if !valid_username(&username) {
|
|
|
|
|
return Ok((MALFORMED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
// identity is exactly KEY_LEN bytes by construction (Reader::take
|
|
|
|
|
// enforces it); no separate curve-point validity check is needed.
|
|
|
|
|
|
2026-09-27 09:39:01 +03:00
|
|
|
// Proof of possession, required on every registration and rotation
|
|
|
|
|
// (SPEC.md sec 6.1). Binding server_static stops the attestation from
|
|
|
|
|
// being replayed against another server.
|
|
|
|
|
let mut pop_msg = LABEL_REGISTER.to_vec();
|
|
|
|
|
pop_msg.extend_from_slice(&self.config.server_static);
|
|
|
|
|
pop_msg.extend_from_slice(username.as_bytes());
|
|
|
|
|
pop_msg.extend_from_slice(&identity);
|
|
|
|
|
if !verify(&identity, &signature, &pop_msg) {
|
|
|
|
|
return Ok((AUTH_FAILED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-26 11:52:22 +03:00
|
|
|
if let Some(expected) = &self.config.invite_token {
|
|
|
|
|
if &token != expected {
|
|
|
|
|
return Ok((NOT_PERMITTED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if cert.is_empty() {
|
|
|
|
|
if !self.store.register(&username, &identity)? {
|
|
|
|
|
return Ok((NOT_PERMITTED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
return Ok((OK, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if cert.len() != CERT_LEN {
|
|
|
|
|
return Ok((MALFORMED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let old_pub = &cert[..32];
|
|
|
|
|
let new_pub = &cert[32..64];
|
|
|
|
|
let when = &cert[64..72];
|
2026-09-27 09:39:01 +03:00
|
|
|
let sig_old = &cert[72..136];
|
|
|
|
|
let sig_new = &cert[136..200];
|
2026-09-26 11:52:22 +03:00
|
|
|
if new_pub != identity.as_slice() {
|
|
|
|
|
return Ok((MALFORMED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let bound = match self.store.identity_of(&username)? {
|
|
|
|
|
Some(b) => b,
|
|
|
|
|
None => return Ok((UNKNOWN_USER, Vec::new())),
|
|
|
|
|
};
|
|
|
|
|
// Only the currently bound key may hand the username on.
|
|
|
|
|
if bound != old_pub {
|
|
|
|
|
return Ok((NOT_PERMITTED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let mut msg = LABEL_ROTATE.to_vec();
|
2026-09-27 09:39:01 +03:00
|
|
|
msg.extend_from_slice(username.as_bytes());
|
2026-09-26 11:52:22 +03:00
|
|
|
msg.extend_from_slice(old_pub);
|
|
|
|
|
msg.extend_from_slice(new_pub);
|
|
|
|
|
msg.extend_from_slice(when);
|
2026-09-27 09:39:01 +03:00
|
|
|
if !verify(old_pub, sig_old, &msg) || !verify(new_pub, sig_new, &msg) {
|
2026-09-26 11:52:22 +03:00
|
|
|
return Ok((AUTH_FAILED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
let chain = self.store.chain(&username)?;
|
|
|
|
|
if chain.len() >= MAX_CHAIN {
|
|
|
|
|
return Ok((NOT_PERMITTED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
if !self.store.rotate(&username, new_pub, &cert, chain.len())? {
|
|
|
|
|
return Ok((NOT_PERMITTED, Vec::new()));
|
|
|
|
|
}
|
|
|
|
|
Ok((OK, Vec::new()))
|
|
|
|
|
}
|
|
|
|
|
}
|