feat: serve multiple domains with per-domain keys, ports and databases
This commit is contained in:
parent
931933c423
commit
dfb57b265c
8 changed files with 891 additions and 120 deletions
332
src/config.rs
Normal file
332
src/config.rs
Normal file
|
|
@ -0,0 +1,332 @@
|
|||
//! Multi-domain TOML configuration: one `[defaults]` table inherited by
|
||||
//! every `[domains.<name>]` table.
|
||||
//!
|
||||
//! bunshin routes by port, not by name: a domain is a listener with its own
|
||||
//! static key, mailbox database and policy knobs. Names exist for logs,
|
||||
//! defaults inheritance and the derived database path. Per-domain settings
|
||||
//! that matter are `key` and `port` (required), the database, the invite
|
||||
//! token and the quota/retention knobs; the remaining knobs are operational
|
||||
//! and belong in `[defaults]`.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
use serde::Deserialize;
|
||||
|
||||
use crate::proto::valid_username;
|
||||
|
||||
// Fallbacks mirroring the CLI flags' defaults.
|
||||
const DEFAULT_HOST: &str = "127.0.0.1";
|
||||
const DEFAULT_MAX_ENVELOPE: usize = 768 << 10;
|
||||
const DEFAULT_QUOTA: i64 = 64 << 20;
|
||||
const DEFAULT_REQUESTS_QUOTA: i64 = 2 << 20;
|
||||
const DEFAULT_RETENTION_DAYS: i64 = 30;
|
||||
const DEFAULT_REQUESTS_RETENTION_DAYS: i64 = 7;
|
||||
const DEFAULT_MAX_TOKENS: u16 = 1024;
|
||||
const DEFAULT_RATE_CONNECTIONS: u32 = 120;
|
||||
const DEFAULT_RATE_SENDS: u32 = 60;
|
||||
const DEFAULT_RATE_TOKENS: u32 = 30;
|
||||
|
||||
/// One resolved domain, ready to serve.
|
||||
pub struct DomainConfig {
|
||||
pub name: String,
|
||||
pub key_path: String,
|
||||
pub db_path: String,
|
||||
pub host: String,
|
||||
pub port: u16,
|
||||
pub max_envelope: usize,
|
||||
pub quota: i64,
|
||||
pub requests_quota: i64,
|
||||
pub retention_days: i64,
|
||||
pub requests_retention_days: i64,
|
||||
pub max_tokens: u16,
|
||||
pub invite_token: Option<Vec<u8>>,
|
||||
pub rate_connections: u32,
|
||||
pub rate_sends: u32,
|
||||
pub rate_tokens: u32,
|
||||
}
|
||||
|
||||
/// The union of every setting either table accepts; `deny_unknown_fields`
|
||||
/// catches typos, which on a quota would otherwise silently fall back to a
|
||||
/// built-in default. The same shape serves `[defaults]` and each domain: the
|
||||
/// presence checks in `load` tell them apart.
|
||||
#[derive(Deserialize, Default)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct Table {
|
||||
key: Option<String>,
|
||||
port: Option<u16>,
|
||||
host: Option<String>,
|
||||
db: Option<String>,
|
||||
max_envelope: Option<usize>,
|
||||
quota: Option<i64>,
|
||||
requests_quota: Option<i64>,
|
||||
retention_days: Option<i64>,
|
||||
requests_retention_days: Option<i64>,
|
||||
max_tokens: Option<u16>,
|
||||
invite_token: Option<String>,
|
||||
invite_token_file: Option<String>,
|
||||
rate_connections: Option<u32>,
|
||||
rate_sends: Option<u32>,
|
||||
rate_tokens: Option<u32>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
struct ConfigFile {
|
||||
#[serde(default)]
|
||||
defaults: Table,
|
||||
domains: BTreeMap<String, Table>,
|
||||
}
|
||||
|
||||
/// Resolves `invite_token`/`invite_token_file` into one token; exactly one
|
||||
/// may be set, or the mailbox is open to registration.
|
||||
fn invite_of(table: &Table, ctx: &str) -> anyhow::Result<Option<Vec<u8>>> {
|
||||
anyhow::ensure!(
|
||||
table.invite_token.is_none() || table.invite_token_file.is_none(),
|
||||
"{ctx}: set only one of invite_token or invite_token_file"
|
||||
);
|
||||
if let Some(token) = &table.invite_token {
|
||||
return Ok(Some(token.clone().into_bytes()));
|
||||
}
|
||||
if let Some(path) = &table.invite_token_file {
|
||||
return std::fs::read(path)
|
||||
.map(Some)
|
||||
.map_err(|e| anyhow::anyhow!("{ctx}: cannot read invite token file {path}: {e}"));
|
||||
}
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
/// Loads and validates a config file. Every key file must exist and every
|
||||
/// port must be bindable before any domain serves, so validation here is
|
||||
/// about the file's own consistency; key readability is checked by the
|
||||
/// server before it binds anything.
|
||||
pub fn load(path: &str) -> anyhow::Result<Vec<DomainConfig>> {
|
||||
let raw = std::fs::read_to_string(path)
|
||||
.map_err(|e| anyhow::anyhow!("cannot read config {path}: {e}"))?;
|
||||
let file: ConfigFile =
|
||||
toml::from_str(&raw).map_err(|e| anyhow::anyhow!("invalid config {path}: {e}"))?;
|
||||
|
||||
anyhow::ensure!(
|
||||
!file.domains.is_empty(),
|
||||
"{path}: no [domains.<name>] tables"
|
||||
);
|
||||
anyhow::ensure!(
|
||||
file.defaults.key.is_none() && file.defaults.port.is_none(),
|
||||
"{path}: [defaults] cannot set key or port"
|
||||
);
|
||||
let default_invite = invite_of(&file.defaults, "[defaults]")?;
|
||||
|
||||
let mut seen_ports = std::collections::HashSet::new();
|
||||
let mut seen_dbs = std::collections::HashSet::new();
|
||||
let mut domains = Vec::with_capacity(file.domains.len());
|
||||
for (name, table) in &file.domains {
|
||||
anyhow::ensure!(
|
||||
valid_username(name),
|
||||
"{path}: domain name {name} must be 1-63 bytes of [a-z0-9._-]"
|
||||
);
|
||||
let (key, port) = match (&table.key, table.port) {
|
||||
(Some(key), Some(port)) => (key.clone(), port),
|
||||
_ => anyhow::bail!("{path}: domains.{name} must set both key and port"),
|
||||
};
|
||||
let host = table
|
||||
.host
|
||||
.clone()
|
||||
.or_else(|| file.defaults.host.clone())
|
||||
.unwrap_or_else(|| DEFAULT_HOST.to_string());
|
||||
anyhow::ensure!(
|
||||
seen_ports.insert((host.clone(), port)),
|
||||
"{path}: domains.{name} reuses {host}:{port}"
|
||||
);
|
||||
let db_path = table
|
||||
.db
|
||||
.clone()
|
||||
.or_else(|| file.defaults.db.clone())
|
||||
.unwrap_or_else(|| format!("mail-{name}.db"));
|
||||
anyhow::ensure!(
|
||||
seen_dbs.insert(db_path.clone()),
|
||||
"{path}: domains.{name} shares database {db_path} with another domain"
|
||||
);
|
||||
|
||||
let invite = invite_of(table, &format!("domains.{name}"))?.or(default_invite.clone());
|
||||
domains.push(DomainConfig {
|
||||
name: name.clone(),
|
||||
key_path: key,
|
||||
db_path,
|
||||
host,
|
||||
port,
|
||||
max_envelope: table
|
||||
.max_envelope
|
||||
.or(file.defaults.max_envelope)
|
||||
.unwrap_or(DEFAULT_MAX_ENVELOPE),
|
||||
quota: table.quota.or(file.defaults.quota).unwrap_or(DEFAULT_QUOTA),
|
||||
requests_quota: table
|
||||
.requests_quota
|
||||
.or(file.defaults.requests_quota)
|
||||
.unwrap_or(DEFAULT_REQUESTS_QUOTA),
|
||||
retention_days: table
|
||||
.retention_days
|
||||
.or(file.defaults.retention_days)
|
||||
.unwrap_or(DEFAULT_RETENTION_DAYS),
|
||||
requests_retention_days: table
|
||||
.requests_retention_days
|
||||
.or(file.defaults.requests_retention_days)
|
||||
.unwrap_or(DEFAULT_REQUESTS_RETENTION_DAYS),
|
||||
max_tokens: table
|
||||
.max_tokens
|
||||
.or(file.defaults.max_tokens)
|
||||
.unwrap_or(DEFAULT_MAX_TOKENS),
|
||||
invite_token: invite,
|
||||
rate_connections: table
|
||||
.rate_connections
|
||||
.or(file.defaults.rate_connections)
|
||||
.unwrap_or(DEFAULT_RATE_CONNECTIONS),
|
||||
rate_sends: table
|
||||
.rate_sends
|
||||
.or(file.defaults.rate_sends)
|
||||
.unwrap_or(DEFAULT_RATE_SENDS),
|
||||
rate_tokens: table
|
||||
.rate_tokens
|
||||
.or(file.defaults.rate_tokens)
|
||||
.unwrap_or(DEFAULT_RATE_TOKENS),
|
||||
});
|
||||
}
|
||||
Ok(domains)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn write_config(name: &str, body: &str) -> String {
|
||||
let path = std::env::temp_dir()
|
||||
.join(format!("bunshin-config-{name}-{}.toml", std::process::id()))
|
||||
.to_string_lossy()
|
||||
.into_owned();
|
||||
std::fs::write(&path, body).unwrap();
|
||||
path
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn minimal_domain_uses_built_in_defaults() {
|
||||
let path = write_config("minimal", "[domains.example]\nkey = \"k\"\nport = 1961\n");
|
||||
let domains = load(&path).unwrap();
|
||||
assert_eq!(domains.len(), 1);
|
||||
let d = &domains[0];
|
||||
assert_eq!(d.name, "example");
|
||||
assert_eq!(d.key_path, "k");
|
||||
assert_eq!(d.db_path, "mail-example.db");
|
||||
assert_eq!(d.host, "127.0.0.1");
|
||||
assert_eq!(d.max_envelope, 768 << 10);
|
||||
assert_eq!(d.quota, 64 << 20);
|
||||
assert_eq!(d.rate_tokens, 30);
|
||||
assert!(d.invite_token.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn domain_overrides_inherit_from_defaults() {
|
||||
let path = write_config(
|
||||
"inherit",
|
||||
"[defaults]\nquota = 1000\nhost = \"0.0.0.0\"\nrate_sends = 5\n\
|
||||
[domains.a]\nkey = \"ka\"\nport = 1\nquota = 2000\n\
|
||||
[domains.b]\nkey = \"kb\"\nport = 2\n",
|
||||
);
|
||||
let domains = load(&path).unwrap();
|
||||
assert_eq!(domains[0].quota, 2000);
|
||||
assert_eq!(domains[0].host, "0.0.0.0");
|
||||
assert_eq!(domains[0].rate_sends, 5);
|
||||
assert_eq!(domains[1].quota, 1000);
|
||||
assert_eq!(domains[1].host, "0.0.0.0");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_setting_is_rejected() {
|
||||
let path = write_config("typo", "[domains.a]\nkey = \"k\"\nport = 1\nqouta = 5\n");
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn duplicate_port_is_rejected() {
|
||||
let path = write_config(
|
||||
"dup-port",
|
||||
"[domains.a]\nkey = \"ka\"\nport = 1\n\
|
||||
[domains.b]\nkey = \"kb\"\nport = 1\n",
|
||||
);
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn duplicate_port_on_different_hosts_is_allowed() {
|
||||
let path = write_config(
|
||||
"dup-port-hosts",
|
||||
"[domains.a]\nkey = \"ka\"\nport = 1\n\
|
||||
[domains.b]\nkey = \"kb\"\nport = 1\nhost = \"127.0.0.2\"\n",
|
||||
);
|
||||
assert!(load(&path).is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn shared_database_is_rejected() {
|
||||
let path = write_config(
|
||||
"dup-db",
|
||||
"[domains.a]\nkey = \"ka\"\nport = 1\ndb = \"shared.db\"\n\
|
||||
[domains.b]\nkey = \"kb\"\nport = 2\ndb = \"shared.db\"\n",
|
||||
);
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invalid_domain_name_is_rejected() {
|
||||
let path = write_config("bad-name", "[domains.Example]\nkey = \"k\"\nport = 1\n");
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn defaults_cannot_set_key_or_port() {
|
||||
let path = write_config(
|
||||
"defaults-port",
|
||||
"[defaults]\nport = 1\n[domains.a]\nkey = \"k\"\nport = 1\n",
|
||||
);
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_key_or_port_is_rejected() {
|
||||
let no_port = write_config("no-port", "[domains.a]\nkey = \"k\"\n");
|
||||
assert!(load(&no_port).is_err());
|
||||
let no_key = write_config("no-key", "[domains.a]\nport = 1\n");
|
||||
assert!(load(&no_key).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_config_is_rejected() {
|
||||
let path = write_config("empty", "[defaults]\nquota = 1\n");
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invite_token_and_file_together_are_rejected() {
|
||||
let path = write_config(
|
||||
"invite-conflict",
|
||||
"[domains.a]\nkey = \"k\"\nport = 1\n\
|
||||
invite_token = \"t\"\ninvite_token_file = \"f\"\n",
|
||||
);
|
||||
assert!(load(&path).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invite_token_file_is_read() {
|
||||
let token_path = std::env::temp_dir()
|
||||
.join(format!("bunshin-config-token-{}.txt", std::process::id()))
|
||||
.to_string_lossy()
|
||||
.into_owned();
|
||||
std::fs::write(&token_path, "sekrit").unwrap();
|
||||
let path = write_config(
|
||||
"invite-file",
|
||||
&format!("[domains.a]\nkey = \"k\"\nport = 1\ninvite_token_file = \"{token_path}\"\n"),
|
||||
);
|
||||
let domains = load(&path).unwrap();
|
||||
assert_eq!(
|
||||
domains[0].invite_token.as_deref(),
|
||||
Some(b"sekrit".as_slice())
|
||||
);
|
||||
}
|
||||
}
|
||||
143
src/harness.rs
143
src/harness.rs
|
|
@ -20,9 +20,9 @@ use snow::{Builder, TransportState};
|
|||
|
||||
use crate::crypto::generate_static_key;
|
||||
use crate::proto::{
|
||||
AUTH_REQUIRED, ENVELOPE_MAGIC, ENVELOPE_VERSION, ID_LEN, LABEL_AUTH, LABEL_REGISTER,
|
||||
MALFORMED, MAX_FRAME, NOISE_PARAMS, NOISE_PAYLOAD, OP_AUTH, OP_FETCH, OP_REGISTER, OP_RESOLVE,
|
||||
OP_SEND, PROLOGUE, TOO_LARGE,
|
||||
AUTH_REQUIRED, ENVELOPE_MAGIC, ENVELOPE_VERSION, ID_LEN, LABEL_AUTH, LABEL_REGISTER, MALFORMED,
|
||||
MAX_FRAME, NOISE_PARAMS, NOISE_PAYLOAD, OP_AUTH, OP_FETCH, OP_REGISTER, OP_RESOLVE, OP_SEND,
|
||||
PROLOGUE, TOO_LARGE, UNKNOWN_USER,
|
||||
};
|
||||
use crate::ratelimit::RateLimiter;
|
||||
use crate::server::handle_connection;
|
||||
|
|
@ -70,13 +70,7 @@ impl Target {
|
|||
// The harness deliberately produces hostile connections;
|
||||
// a panic in one must not take the accept loop with it.
|
||||
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
|
||||
handle_connection(
|
||||
stream,
|
||||
config,
|
||||
&static_private,
|
||||
&db_path,
|
||||
"127.0.0.1",
|
||||
)
|
||||
handle_connection(stream, config, &static_private, &db_path, "127.0.0.1")
|
||||
}));
|
||||
});
|
||||
}
|
||||
|
|
@ -92,7 +86,8 @@ impl Target {
|
|||
fn assert_alive(&self) {
|
||||
let mut client = Client::connect(self.port);
|
||||
client.handshake().expect("server survived");
|
||||
let (status, _) = client.request(OP_RESOLVE, &[7, b'h', b'a', b'r', b'n', b'e', b's', b's']);
|
||||
let (status, _) =
|
||||
client.request(OP_RESOLVE, &[7, b'h', b'a', b'r', b'n', b'e', b's', b's']);
|
||||
assert_ne!(status, 0xFF, "server answered with a status byte");
|
||||
}
|
||||
}
|
||||
|
|
@ -143,7 +138,9 @@ impl Client {
|
|||
let mut packet = vec![0u8; payload.len() + 16];
|
||||
let n = transport.write_message(payload, &mut packet).unwrap();
|
||||
packet.truncate(n);
|
||||
self.stream.write_all(&(packet.len() as u16).to_be_bytes()).unwrap();
|
||||
self.stream
|
||||
.write_all(&(packet.len() as u16).to_be_bytes())
|
||||
.unwrap();
|
||||
self.stream.write_all(&packet).unwrap();
|
||||
}
|
||||
|
||||
|
|
@ -226,7 +223,10 @@ fn short_garbage_handshake_is_refused_and_server_survives() {
|
|||
client.write_raw(&(10u16).to_be_bytes());
|
||||
client.write_raw(&[0xA5; 10]);
|
||||
let mut sink = [0u8; 16];
|
||||
assert!(client.stream.read(&mut sink).unwrap_or(0) == 0, "server closed");
|
||||
assert!(
|
||||
client.stream.read(&mut sink).unwrap_or(0) == 0,
|
||||
"server closed"
|
||||
);
|
||||
target.assert_alive();
|
||||
}
|
||||
|
||||
|
|
@ -299,7 +299,10 @@ fn zero_length_frame_gets_malformed_then_close() {
|
|||
assert_eq!(op, 0);
|
||||
assert_eq!(body[0], MALFORMED);
|
||||
let mut sink = [0u8; 16];
|
||||
assert!(client.stream.read(&mut sink).unwrap_or(0) == 0, "closed after bad frame");
|
||||
assert!(
|
||||
client.stream.read(&mut sink).unwrap_or(0) == 0,
|
||||
"closed after bad frame"
|
||||
);
|
||||
target.assert_alive();
|
||||
}
|
||||
|
||||
|
|
@ -414,7 +417,11 @@ fn registered_flow_survives_absurd_future_cursor() {
|
|||
let (op, resp) = client.request(OP_FETCH, &body);
|
||||
assert_eq!(op, OP_FETCH);
|
||||
assert_eq!(resp[0], 0, "status OK");
|
||||
assert_eq!(u16::from_be_bytes(resp[1..3].try_into().unwrap()), 0, "empty page");
|
||||
assert_eq!(
|
||||
u16::from_be_bytes(resp[1..3].try_into().unwrap()),
|
||||
0,
|
||||
"empty page"
|
||||
);
|
||||
|
||||
// And the registered name resolves — the store came through intact.
|
||||
let (op, resp) = client.request(OP_RESOLVE, &[7, b'h', b'a', b'r', b'n', b'e', b's', b's']);
|
||||
|
|
@ -423,7 +430,6 @@ fn registered_flow_survives_absurd_future_cursor() {
|
|||
target.assert_alive();
|
||||
}
|
||||
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Long-lived multi-session driver: the last agreed stress-campaign phase.
|
||||
// A dozen authenticated sessions live on one server simultaneously, run
|
||||
|
|
@ -615,3 +621,108 @@ fn partial_frame_is_dropped_after_the_header_deadline() {
|
|||
);
|
||||
target.assert_alive();
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Multi-domain: one process, one listener/key/database per domain. The
|
||||
// well-behaved client drives the real serve_domains path (config file, real
|
||||
// key files, real listeners) and must see two isolated namespaces.
|
||||
|
||||
/// Reserves a free localhost port by binding to :0 and dropping the socket.
|
||||
fn free_port() -> u16 {
|
||||
std::net::TcpListener::bind("127.0.0.1:0")
|
||||
.unwrap()
|
||||
.local_addr()
|
||||
.unwrap()
|
||||
.port()
|
||||
}
|
||||
|
||||
/// Connects to a domain's listener, tolerating the window between
|
||||
/// serve_domains starting and binding its socket.
|
||||
fn connect_domain(port: u16) -> TcpStream {
|
||||
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
|
||||
loop {
|
||||
match TcpStream::connect(("127.0.0.1", port)) {
|
||||
Ok(stream) => return stream,
|
||||
Err(_) if std::time::Instant::now() < deadline => {
|
||||
std::thread::sleep(std::time::Duration::from_millis(20));
|
||||
}
|
||||
Err(e) => panic!("cannot connect to domain port {port}: {e}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn domains_isolate_namespaces() {
|
||||
let dir = std::env::temp_dir()
|
||||
.join(format!("bunshin-domains-{}", std::process::id()))
|
||||
.join("isolate");
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
let (key_a, pub_a) = generate_static_key();
|
||||
let (key_b, _pub_b) = generate_static_key();
|
||||
std::fs::write(dir.join("a.key"), key_a).unwrap();
|
||||
std::fs::write(dir.join("b.key"), key_b).unwrap();
|
||||
let port_a = free_port();
|
||||
let port_b = free_port();
|
||||
let config_path = dir.join("bunshin.toml");
|
||||
std::fs::write(
|
||||
&config_path,
|
||||
format!(
|
||||
"[domains.a]\nkey = \"{}\"\nport = {port_a}\ndb = \"{}\"\n\
|
||||
[domains.b]\nkey = \"{}\"\nport = {port_b}\ndb = \"{}\"\n",
|
||||
dir.join("a.key").display(),
|
||||
dir.join("a.db").display(),
|
||||
dir.join("b.key").display(),
|
||||
dir.join("b.db").display()
|
||||
),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let domains = crate::config::load(config_path.to_str().unwrap()).unwrap();
|
||||
std::thread::spawn(move || crate::server::serve_domains(domains).unwrap());
|
||||
|
||||
let register = |mut client: Client, public: &[u8; 32]| {
|
||||
client.handshake().unwrap();
|
||||
let key = SigningKey::from_bytes(&[11u8; 32]);
|
||||
let verifying = key.verifying_key();
|
||||
let identity: &[u8] = verifying.as_bytes();
|
||||
let mut pop = LABEL_REGISTER.to_vec();
|
||||
pop.extend_from_slice(public);
|
||||
pop.extend_from_slice(b"alice");
|
||||
pop.extend_from_slice(identity);
|
||||
let mut body = vec![5];
|
||||
body.extend_from_slice(b"alice");
|
||||
body.extend_from_slice(identity);
|
||||
body.extend_from_slice(&key.sign(&pop).to_bytes());
|
||||
body.push(0);
|
||||
body.push(0);
|
||||
let (op, resp) = client.request(OP_REGISTER, &body);
|
||||
assert_eq!((op, resp[0]), (OP_REGISTER, 0), "REGISTER accepted");
|
||||
client
|
||||
};
|
||||
|
||||
// alice exists only on domain a; the same wire bytes against domain b
|
||||
// must find no such user.
|
||||
let mut client = register(
|
||||
Client {
|
||||
stream: connect_domain(port_a),
|
||||
transport: None,
|
||||
handshake_hash: Vec::new(),
|
||||
},
|
||||
&pub_a,
|
||||
);
|
||||
let (op, resp) = client.request(OP_RESOLVE, &[5, b'a', b'l', b'i', b'c', b'e']);
|
||||
assert_eq!((op, resp[0]), (OP_RESOLVE, 0), "RESOLVE on own domain");
|
||||
|
||||
let mut client = Client {
|
||||
stream: connect_domain(port_b),
|
||||
transport: None,
|
||||
handshake_hash: Vec::new(),
|
||||
};
|
||||
client.handshake().unwrap();
|
||||
let (op, resp) = client.request(OP_RESOLVE, &[5, b'a', b'l', b'i', b'c', b'e']);
|
||||
assert_eq!(
|
||||
(op, resp[0]),
|
||||
(OP_RESOLVE, UNKNOWN_USER),
|
||||
"b has no alice: databases are isolated"
|
||||
);
|
||||
}
|
||||
|
|
|
|||
48
src/main.rs
48
src/main.rs
|
|
@ -2,6 +2,7 @@
|
|||
|
||||
mod bind;
|
||||
mod channel;
|
||||
mod config;
|
||||
mod crypto;
|
||||
#[cfg(test)]
|
||||
mod harness;
|
||||
|
|
@ -42,6 +43,11 @@ enum Command {
|
|||
},
|
||||
/// Run the mailbox server
|
||||
Serve {
|
||||
/// Serve one domain per [domains.<name>] table; every domain has its
|
||||
/// own key, port and mailbox database
|
||||
#[arg(long, conflicts_with_all = ["key", "db", "host", "port", "max_envelope", "quota", "requests_quota", "retention_days", "requests_retention_days", "max_tokens", "invite_token", "rate_connections", "rate_sends", "rate_tokens"])]
|
||||
#[cfg_attr(feature = "rns", arg(conflicts_with_all = ["enabled", "rns_key", "rns_max_envelope", "rns_fetch_budget", "rns_max_links", "rns_rate_link_requests", "rns_rate_link_bytes", "rns_link_idle", "rns_udp", "rns_udp_forward"]))]
|
||||
config: Option<String>,
|
||||
#[arg(long, default_value = "server.key")]
|
||||
key: String,
|
||||
#[arg(long, default_value = "mail.db")]
|
||||
|
|
@ -132,6 +138,7 @@ fn main() -> anyhow::Result<()> {
|
|||
Command::RnsKeygen { key, force } => cmd_rns_keygen(&key, force),
|
||||
#[cfg(not(feature = "rns"))]
|
||||
Command::Serve {
|
||||
config,
|
||||
key,
|
||||
db,
|
||||
host,
|
||||
|
|
@ -146,24 +153,30 @@ fn main() -> anyhow::Result<()> {
|
|||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
} => server::run(server::ServeArgs {
|
||||
key_path: key,
|
||||
db_path: db,
|
||||
host,
|
||||
port,
|
||||
max_envelope,
|
||||
quota,
|
||||
requests_quota,
|
||||
retention_days,
|
||||
requests_retention_days,
|
||||
max_tokens,
|
||||
invite_token,
|
||||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
}),
|
||||
} => {
|
||||
if let Some(path) = config {
|
||||
return server::serve_domains(config::load(&path)?);
|
||||
}
|
||||
server::run(server::ServeArgs {
|
||||
key_path: key,
|
||||
db_path: db,
|
||||
host,
|
||||
port,
|
||||
max_envelope,
|
||||
quota,
|
||||
requests_quota,
|
||||
retention_days,
|
||||
requests_retention_days,
|
||||
max_tokens,
|
||||
invite_token,
|
||||
rate_connections,
|
||||
rate_sends,
|
||||
rate_tokens,
|
||||
})
|
||||
}
|
||||
#[cfg(feature = "rns")]
|
||||
Command::Serve {
|
||||
config,
|
||||
key,
|
||||
db,
|
||||
host,
|
||||
|
|
@ -180,6 +193,9 @@ fn main() -> anyhow::Result<()> {
|
|||
rate_tokens,
|
||||
rns,
|
||||
} => {
|
||||
if let Some(path) = config {
|
||||
return server::serve_domains(config::load(&path)?);
|
||||
}
|
||||
// One process serves both carriers (RNS.md sec 7.4): shared
|
||||
// Store, shared config, single purge loop in server::run.
|
||||
if rns.enabled {
|
||||
|
|
|
|||
117
src/server.rs
117
src/server.rs
|
|
@ -1,11 +1,12 @@
|
|||
//! TCP accept loop, per-connection handling, and the background purge loop.
|
||||
|
||||
use std::net::TcpStream;
|
||||
use std::net::{TcpListener, TcpStream};
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::bind::TransportBindValues;
|
||||
use crate::channel::{handshake, Channel};
|
||||
use crate::config::DomainConfig;
|
||||
use crate::crypto::{b32, derive_public};
|
||||
use crate::proto::{FETCH_BUDGET, HEADER_TIMEOUT_SECS, KEY_LEN, MALFORMED, PURGE_INTERVAL_SECS};
|
||||
use crate::ratelimit::RateLimiter;
|
||||
|
|
@ -29,41 +30,97 @@ pub struct ServeArgs {
|
|||
pub rate_tokens: u32,
|
||||
}
|
||||
|
||||
/// The single-domain CLI path: one anonymous domain.
|
||||
pub fn run(args: ServeArgs) -> anyhow::Result<()> {
|
||||
let static_key = std::fs::read(&args.key_path)
|
||||
.map_err(|e| anyhow::anyhow!("cannot read server key: {e}"))?;
|
||||
anyhow::ensure!(
|
||||
static_key.len() == KEY_LEN,
|
||||
"server key must be {KEY_LEN} raw bytes, got {}",
|
||||
static_key.len()
|
||||
);
|
||||
|
||||
let key_array: [u8; KEY_LEN] = static_key.clone().try_into().unwrap();
|
||||
let server_static = derive_public(&key_array);
|
||||
|
||||
let config = Arc::new(ServerConfig {
|
||||
serve_domains(vec![DomainConfig {
|
||||
name: "default".to_string(),
|
||||
key_path: args.key_path,
|
||||
db_path: args.db_path,
|
||||
host: args.host,
|
||||
port: args.port,
|
||||
max_envelope: args.max_envelope,
|
||||
fetch_budget: FETCH_BUDGET,
|
||||
main_quota: args.quota,
|
||||
quota: args.quota,
|
||||
requests_quota: args.requests_quota,
|
||||
retention_days: args.retention_days,
|
||||
requests_retention_days: args.requests_retention_days,
|
||||
max_tokens: args.max_tokens,
|
||||
invite_token: args.invite_token.map(String::into_bytes),
|
||||
conn_limiter: RateLimiter::new(args.rate_connections),
|
||||
send_limiter: RateLimiter::new(args.rate_sends),
|
||||
token_limiter: RateLimiter::new(args.rate_tokens),
|
||||
});
|
||||
rate_connections: args.rate_connections,
|
||||
rate_sends: args.rate_sends,
|
||||
rate_tokens: args.rate_tokens,
|
||||
}])
|
||||
}
|
||||
|
||||
let main_retention_secs = args.retention_days * 86400;
|
||||
let requests_retention_secs = args.requests_retention_days * 86400;
|
||||
let purge_db_path = args.db_path.clone();
|
||||
std::thread::spawn(move || {
|
||||
purge_loop(purge_db_path, main_retention_secs, requests_retention_secs)
|
||||
});
|
||||
/// Serves every domain: one listener, key, mailbox database, config and
|
||||
/// purge loop per domain, nothing shared between them.
|
||||
pub fn serve_domains(domains: Vec<DomainConfig>) -> anyhow::Result<()> {
|
||||
// Read every key and bind every listener first, so a missing key or a
|
||||
// taken port fails the whole process before any domain serves.
|
||||
let mut bound = Vec::with_capacity(domains.len());
|
||||
for domain in &domains {
|
||||
let static_key = read_static_key(&domain.key_path)?;
|
||||
let listener = TcpListener::bind((domain.host.as_str(), domain.port))?;
|
||||
let config = Arc::new(ServerConfig {
|
||||
max_envelope: domain.max_envelope,
|
||||
fetch_budget: FETCH_BUDGET,
|
||||
main_quota: domain.quota,
|
||||
requests_quota: domain.requests_quota,
|
||||
max_tokens: domain.max_tokens,
|
||||
invite_token: domain.invite_token.clone(),
|
||||
conn_limiter: RateLimiter::new(domain.rate_connections),
|
||||
send_limiter: RateLimiter::new(domain.rate_sends),
|
||||
token_limiter: RateLimiter::new(domain.rate_tokens),
|
||||
});
|
||||
bound.push((domain, static_key, listener, config));
|
||||
}
|
||||
|
||||
let listener = std::net::TcpListener::bind((args.host.as_str(), args.port))?;
|
||||
log::info!("listening on {}:{}", args.host, args.port);
|
||||
log::info!("server public key: {}", b32(&server_static));
|
||||
for (domain, static_key, listener, config) in bound {
|
||||
log::info!(
|
||||
"domain {}: listening on {}:{}",
|
||||
domain.name,
|
||||
domain.host,
|
||||
domain.port
|
||||
);
|
||||
log::info!(
|
||||
"domain {}: server public key: {}",
|
||||
domain.name,
|
||||
b32(&derive_public(&static_key))
|
||||
);
|
||||
|
||||
let purge_db_path = domain.db_path.clone();
|
||||
let main_retention_secs = domain.retention_days * 86400;
|
||||
let requests_retention_secs = domain.requests_retention_days * 86400;
|
||||
std::thread::spawn(move || {
|
||||
purge_loop(purge_db_path, main_retention_secs, requests_retention_secs)
|
||||
});
|
||||
|
||||
let db_path = domain.db_path.clone();
|
||||
std::thread::spawn(move || accept_loop(listener, config, static_key, db_path));
|
||||
}
|
||||
|
||||
// Each domain's accept loop runs in its own thread; nothing fails here.
|
||||
loop {
|
||||
std::thread::sleep(Duration::from_secs(3600));
|
||||
}
|
||||
}
|
||||
|
||||
fn read_static_key(path: &str) -> anyhow::Result<[u8; KEY_LEN]> {
|
||||
let key =
|
||||
std::fs::read(path).map_err(|e| anyhow::anyhow!("cannot read server key {path}: {e}"))?;
|
||||
anyhow::ensure!(
|
||||
key.len() == KEY_LEN,
|
||||
"server key {path} must be {KEY_LEN} raw bytes, got {}",
|
||||
key.len()
|
||||
);
|
||||
Ok(key.try_into().unwrap())
|
||||
}
|
||||
|
||||
fn accept_loop(
|
||||
listener: TcpListener,
|
||||
config: Arc<ServerConfig>,
|
||||
static_key: [u8; KEY_LEN],
|
||||
db_path: String,
|
||||
) {
|
||||
for incoming in listener.incoming() {
|
||||
let stream = match incoming {
|
||||
Ok(s) => s,
|
||||
|
|
@ -83,15 +140,13 @@ pub fn run(args: ServeArgs) -> anyhow::Result<()> {
|
|||
}
|
||||
|
||||
let config = Arc::clone(&config);
|
||||
let static_key = key_array;
|
||||
let db_path = args.db_path.clone();
|
||||
let db_path = db_path.clone();
|
||||
std::thread::spawn(move || {
|
||||
if let Err(e) = handle_connection(stream, &config, &static_key, &db_path, &peer_ip) {
|
||||
log::info!("connection error from {peer_ip}: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Also the hostile harness's entry point: it drives real connections
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue