oxedyne/fe2o3/fe2o3_steel/src/srv/server.rs
35.1 KiB, 188 runs
created by r1870400018:993, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | use crate::{ |
| 2 | app::mail::{ |
| 3 | AppMailHandler, |
| 4 | run_outbound_worker, |
| 5 | }, |
| 6 | srv::{ |
| 7 | cert::Certificate, |
| 8 | cfg::MailConfig, |
| 9 | context::{ |
| 10 | Protocol, |
| 11 | ServerContext, |
| 12 | }, |
| 13 | http::{ |
| 14 | handle_redirect, |
| 15 | run_redirect_listener, |
| 16 | }, |
| 17 | mail::{ |
| 18 | build_smtp_servers, |
| 19 | run_imap_listener, |
| 20 | run_smtp_listener, |
| 21 | AppImapServer, |
| 22 | ListenerTally, |
| 23 | }, |
| 24 | stop, |
| 25 | }, |
| 26 | }; |
| 27 | |
| 28 | use oxedyne_fe2o3_mail::{ |
| 29 | maildir::MaildirStore, |
| 30 | outbound::OutboundSpool, |
| 31 | passwd::PasswdFileUserStore, |
| 32 | }; |
| 33 | use oxedyne_fe2o3_net::{ |
| 34 | dkim::DkimSigner, |
| 35 | imap::server::ImapServer, |
| 36 | smtp::client::OutboundClient, |
| 37 | tls::{ |
| 38 | BoundedTlsAcceptor, |
| 39 | Handshake, |
| 40 | }, |
| 41 | }; |
| 42 | |
| 43 | use oxedyne_fe2o3_core::{ |
| 44 | prelude::*, |
| 45 | file as core_file, |
| 46 | }; |
| 47 | use oxedyne_fe2o3_iop_crypto::enc::Encrypter; |
| 48 | use oxedyne_fe2o3_iop_db::api::Database; |
| 49 | use oxedyne_fe2o3_iop_hash::api::Hasher; |
| 50 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 51 | use oxedyne_fe2o3_net::{ |
| 52 | http::handler::WebHandler, |
| 53 | ws::handler::WebSocketHandler, |
| 54 | }; |
| 55 | |
| 56 | use std::{ |
| 57 | net::SocketAddr, |
| 58 | sync::{ |
| 59 | atomic::{ |
| 60 | AtomicUsize, |
| 61 | Ordering, |
| 62 | }, |
| 63 | Arc, |
| 64 | }, |
| 65 | time::Duration, |
| 66 | }; |
| 67 | |
| 68 | use tokio::net::TcpListener; |
| 69 | use tokio_rustls::TlsAcceptor; |
| 70 | |
| 71 | pub const PERSIST_INTERVAL_SECS: u64 = 60; |
| 72 | |
| 73 | pub const DRAIN_SECS: u64 = 5; |
| 74 | |
| 75 | // How often the address guard is swept for idle per-IP records, and how long a |
| 76 | // record may go unseen before an idle sweep may evict it. Ten minutes of silence |
| 77 | // is comfortably longer than any keep-alive, so a genuinely active address is |
| 78 | // never forgotten out from under its per-IP cap. |
| 79 | pub const GUARD_SWEEP_INTERVAL_SECS: u64 = 60; |
| 80 | pub const GUARD_SWEEP_IDLE_SECS: u64 = 600; |
| 81 | |
| 82 | |
| 83 | pub struct Server< |
| 84 | const UIDL: usize, |
| 85 | UID: NumIdDat<UIDL> + 'static, |
| 86 | ENC: Encrypter, |
| 87 | KH: Hasher, |
| 88 | DB: Database<UIDL, UID, ENC, KH>, |
| 89 | WH: WebHandler, |
| 90 | WSH: WebSocketHandler, |
| 91 | > { |
| 92 | pub context: ServerContext<UIDL, UID, ENC, KH, DB, WH, WSH>, |
| 93 | } |
| 94 | |
| 95 | impl< |
| 96 | const UIDL: usize, |
| 97 | UID: NumIdDat<UIDL> + 'static, |
| 98 | ENC: Encrypter + 'static, |
| 99 | KH: Hasher + 'static, |
| 100 | DB: Database<UIDL, UID, ENC, KH> + 'static, |
| 101 | WH: WebHandler + 'static, |
| 102 | WSH: WebSocketHandler + 'static, |
| 103 | > |
| 104 | Server<UIDL, UID, ENC, KH, DB, WH, WSH> |
| 105 | { |
| 106 | pub fn new( |
| 107 | context: ServerContext<UIDL, UID, ENC, KH, DB, WH, WSH>, |
| 108 | ) |
| 109 | -> Self |
| 110 | { |
| 111 | Self { context } |
| 112 | } |
| 113 | |
| 114 | fn primary_db_handle(&self) -> Option<(Arc<std::sync::RwLock<DB>>, UID)> { |
| 115 | let default_vhost = match &self.context.protocol { |
| 116 | Protocol::Web { default_vhost, .. } => default_vhost.clone(), |
| 117 | }; |
| 118 | self.context.db_for_vhost(&default_vhost) |
| 119 | } |
| 120 | |
| 121 | pub async fn start(&self) -> Outcome<()> { |
| 122 | |
| 123 | let dev_mode = match &self.context.protocol { |
| 124 | Protocol::Web { dev_mode, .. } => *dev_mode, |
| 125 | }; |
| 126 | |
| 127 | let loaded = res!(Certificate::load( |
| 128 | &self.context.cfg, |
| 129 | &self.context.root, |
| 130 | dev_mode, |
| 131 | )); |
| 132 | |
| 133 | // If ACME is enabled, spawn the renewer task. It drives the |
| 134 | // initial issuance (if the cache is empty) and then loops with |
| 135 | // a 24-hour tick, re-issuing whenever the cached cert is older |
| 136 | // than the renewal threshold. |
| 137 | if let Some(renewer) = loaded.acme_renewer { |
| 138 | tokio::spawn(async move { |
| 139 | if let Err(e) = renewer.run_forever().await { |
| 140 | error!(err!(e, |
| 141 | "ACME renewer task exited."; |
| 142 | Init, Network)); |
| 143 | } |
| 144 | }); |
| 145 | } |
| 146 | |
| 147 | // Spawn the periodic traffic sampler if the traffic |
| 148 | // recorder is configured. Samples counters into the |
| 149 | // recorder's bounded history ring every |
| 150 | // DEFAULT_SAMPLE_INTERVAL_SECS seconds so the dashboard |
| 151 | // traffic view has time-series data to chart. Lives here |
| 152 | // (inside the async Server::start) rather than in the sync |
| 153 | // AppShellContext::start_server so the current tokio |
| 154 | // runtime is active when the spawn happens. |
| 155 | if let Some(recorder) = self.context.traffic.clone() { |
| 156 | tokio::spawn(async move { |
| 157 | let mut ticker = tokio::time::interval( |
| 158 | std::time::Duration::from_secs( |
| 159 | crate::srv::admin::traffic::DEFAULT_SAMPLE_INTERVAL_SECS, |
| 160 | ), |
| 161 | ); |
| 162 | // Skip the first immediate tick so the first |
| 163 | // real sample is taken after one full interval. |
| 164 | ticker.tick().await; |
| 165 | loop { |
| 166 | ticker.tick().await; |
| 167 | if let Err(e) = recorder.sample_now() { |
| 168 | warn!("traffic sampler: {}", e); |
| 169 | } |
| 170 | } |
| 171 | }); |
| 172 | } |
| 173 | |
| 174 | // Spawn the periodic host sampler if the admin state is |
| 175 | // configured. Reads `/proc/*` via `fe2o3_sys::Snapshot` |
| 176 | // and pushes a new entry into the sampler's bounded |
| 177 | // history ring on every tick, so the dashboard's host |
| 178 | // resource strip has real CPU / memory / disk / network |
| 179 | // data to render. Lives next to the traffic sampler so |
| 180 | // both get the same late-spawn treatment. |
| 181 | // |
| 182 | // Before the sampler starts, load any persisted sparkline |
| 183 | // points the previous run saved to ozone and seed the |
| 184 | // sampler so the Overview chart does not reset to blank |
| 185 | // across a restart. |
| 186 | if let Some(admin) = self.context.admin_state.clone() { |
| 187 | let sampler = admin.host_sampler.clone(); |
| 188 | |
| 189 | // Start-up restore. Pick the default vhost's database |
| 190 | // (the primary / the vhost the dashboard is attached |
| 191 | // to) and pull the derived history back. |
| 192 | let primary_db = self.primary_db_handle(); |
| 193 | if let Some((db, _uid)) = primary_db.as_ref() { |
| 194 | let db_guard = db.read(); |
| 195 | if let Ok(db_g) = db_guard { |
| 196 | match crate::srv::admin::persist::load_host_points(&*db_g) { |
| 197 | Ok(points) if !points.is_empty() => { |
| 198 | info!("admin persist: restored {} derived host points.", |
| 199 | points.len()); |
| 200 | if let Err(e) = sampler.seed_persisted(points) { |
| 201 | warn!("host sampler seed: {}", e); |
| 202 | } |
| 203 | }, |
| 204 | Ok(_) => { |
| 205 | info!("admin persist: no previous host history to restore."); |
| 206 | }, |
| 207 | Err(e) => { |
| 208 | warn!("admin persist: load host points failed: {}", e); |
| 209 | }, |
| 210 | } |
| 211 | } |
| 212 | } |
| 213 | |
| 214 | let sampler_for_task = sampler.clone(); |
| 215 | tokio::spawn(async move { |
| 216 | let mut ticker = tokio::time::interval( |
| 217 | std::time::Duration::from_secs( |
| 218 | crate::srv::admin::host_sampler::DEFAULT_SAMPLE_INTERVAL_SECS, |
| 219 | ), |
| 220 | ); |
| 221 | // Prime with an immediate read so the dashboard |
| 222 | // has a non-empty ring on the first request. |
| 223 | if let Err(e) = sampler_for_task.sample_now() { |
| 224 | warn!("host sampler prime: {}", e); |
| 225 | } |
| 226 | ticker.tick().await; |
| 227 | loop { |
| 228 | ticker.tick().await; |
| 229 | if let Err(e) = sampler_for_task.sample_now() { |
| 230 | warn!("host sampler: {}", e); |
| 231 | } |
| 232 | } |
| 233 | }); |
| 234 | |
| 235 | // Periodic saver: write the merged derived history to |
| 236 | // ozone every PERSIST_INTERVAL_SECS so a crash or |
| 237 | // restart loses at most one tick's worth of history. |
| 238 | // Spawned in addition to -- not instead of -- the |
| 239 | // live sampler task so a slow disk cannot throttle |
| 240 | // the sampling cadence. |
| 241 | if let Some((db, uid)) = primary_db { |
| 242 | let sampler_for_save = sampler.clone(); |
| 243 | tokio::spawn(async move { |
| 244 | let mut ticker = tokio::time::interval( |
| 245 | std::time::Duration::from_secs(PERSIST_INTERVAL_SECS), |
| 246 | ); |
| 247 | ticker.tick().await; |
| 248 | loop { |
| 249 | ticker.tick().await; |
| 250 | let points = match sampler_for_save.merged_derived_history() { |
| 251 | Ok(p) => p, |
| 252 | Err(e) => { |
| 253 | warn!("admin persist: merged history failed: {}", e); |
| 254 | continue; |
| 255 | }, |
| 256 | }; |
| 257 | if points.is_empty() { |
| 258 | continue; |
| 259 | } |
| 260 | let db_guard = db.read(); |
| 261 | let db_g = match db_guard { |
| 262 | Ok(g) => g, |
| 263 | Err(_) => { |
| 264 | warn!("admin persist: db read lock poisoned."); |
| 265 | continue; |
| 266 | }, |
| 267 | }; |
| 268 | if let Err(e) = crate::srv::admin::persist::save_host_points( |
| 269 | &*db_g, |
| 270 | uid, |
| 271 | &points, |
| 272 | ) { |
| 273 | warn!("admin persist: save host points failed: {}", e); |
| 274 | } |
| 275 | } |
| 276 | }); |
| 277 | } |
| 278 | } |
| 279 | |
| 280 | // Spawn the plaintext HTTP redirect listener if configured. |
| 281 | // This binds a separate port (typically 80) and responds to every |
| 282 | // incoming HTTP request with a 301 to the HTTPS equivalent, so |
| 283 | // browsers that do not default to HTTPS-first mode can still reach |
| 284 | // the site by typing a bare hostname. |
| 285 | let http_port = self.context.cfg.server_port_tcp_plaintext; |
| 286 | if http_port != 0 { |
| 287 | let address = self.context.cfg.server_address.clone(); |
| 288 | let https_port = self.context.cfg.server_port_tcp; |
| 289 | tokio::spawn(async move { |
| 290 | if let Err(e) = run_redirect_listener( |
| 291 | address, |
| 292 | http_port, |
| 293 | https_port, |
| 294 | ).await { |
| 295 | error!(err!(e, |
| 296 | "Plaintext HTTP redirect listener exited."; |
| 297 | Init, Network)); |
| 298 | } |
| 299 | }); |
| 300 | } |
| 301 | |
| 302 | // Spawn the localhost plain-HTTP admin listener if |
| 303 | // configured. This binds 127.0.0.1:<port> for /admin only, |
| 304 | // intended to be reached via SSH tunnel when the public |
| 305 | // TLS chain is broken or the operator wants emergency |
| 306 | // access. Bound only to loopback by design; there is no |
| 307 | // network-exposed knob. |
| 308 | let admin_local_port = self.context.cfg.admin_local_port; |
| 309 | if admin_local_port != 0 { |
| 310 | let ctx_for_local = self.context.clone(); |
| 311 | tokio::spawn(async move { |
| 312 | if let Err(e) = ctx_for_local.run_admin_local_listener( |
| 313 | admin_local_port, |
| 314 | ).await { |
| 315 | error!(err!(e, |
| 316 | "Admin localhost listener exited."; |
| 317 | Init, Network)); |
| 318 | } |
| 319 | }); |
| 320 | } |
| 321 | |
| 322 | let tls_acceptor = TlsAcceptor::from(Arc::new(loaded.server_config)); |
| 323 | |
| 324 | // Bound every TLS handshake -- HTTPS and the mail listeners share one |
| 325 | // rustls server config, so they share one bound. `max_tls_handshakes` |
| 326 | // caps how many handshakes run at once (the CPU-exhaustion vector on a |
| 327 | // single-vCPU box) and `tls_handshake_timeout_ms` deadlines each one so |
| 328 | // a drip-fed handshake cannot pin a permit and turn the cap into a |
| 329 | // slowloris amplifier. Both zero is the inert default: a bare acceptor. |
| 330 | let hs_deadline = if self.context.cfg.tls_handshake_timeout_ms == 0 { |
| 331 | None |
| 332 | } else { |
| 333 | Some(Duration::from_millis(self.context.cfg.tls_handshake_timeout_ms)) |
| 334 | }; |
| 335 | let bounded_acceptor = BoundedTlsAcceptor::new( |
| 336 | tls_acceptor, |
| 337 | self.context.cfg.max_tls_handshakes as usize, |
| 338 | hs_deadline, |
| 339 | ); |
| 340 | |
| 341 | // Spawn the mail listeners only when the mail server is enabled. A site that |
| 342 | // merely sends newsletters wants a DKIM identity and an outbound client, not |
| 343 | // to bind SMTP-receive, submission and IMAP and become an MX -- so sending is |
| 344 | // built separately (see the newsletter sender in `app/server.rs`), and this |
| 345 | // gate keeps a send-only host from standing up a mail server it never asked |
| 346 | // for. |
| 347 | if let Some(mail_cfg) = res!(self.context.cfg.get_mail()) { |
| 348 | if mail_cfg.enabled { |
| 349 | // Counted into the health body's `mail_down`, so a listener that |
| 350 | // failed to bind is visible to every watcher, not only in this log. |
| 351 | let tally = self.context.admin_state.as_ref().map(|a| a.mail.clone()); |
| 352 | if let Err(e) = spawn_mail_listeners( |
| 353 | &mail_cfg, |
| 354 | &self.context.root, |
| 355 | bounded_acceptor.clone(), |
| 356 | &self.context.cfg.server_address, |
| 357 | tally, |
| 358 | ).await { |
| 359 | error!(err!(e, |
| 360 | "Failed to spawn mail listeners."; |
| 361 | Init, Network)); |
| 362 | } |
| 363 | } |
| 364 | } |
| 365 | |
| 366 | // Build the bind address from the (now honoured) server_cfg. |
| 367 | let addr: SocketAddr = { |
| 368 | let ip: std::net::IpAddr = match self.context.cfg.server_address.parse() { |
| 369 | Ok(ip) => ip, |
| 370 | Err(e) => return Err(err!(e, |
| 371 | "Invalid server_address '{}' in config.", |
| 372 | self.context.cfg.server_address; |
| 373 | Invalid, Input, Network)), |
| 374 | }; |
| 375 | SocketAddr::new(ip, self.context.cfg.server_port_tcp) |
| 376 | }; |
| 377 | let listener = res!(TcpListener::bind(&addr).await, IO, Network); |
| 378 | info!("Listening on: {}", addr); |
| 379 | |
| 380 | // Shared address guard, cloned once per accept. Cheap -- it's an Arc. |
| 381 | let addr_guard = self.context.admin_state.as_ref() |
| 382 | .map(|a| a.addr_guard.clone()); |
| 383 | |
| 384 | // The rolling count of connections dropped at admission, surfaced as |
| 385 | // `dropped_1m` in the health body. `None` when the admin dashboard is |
| 386 | // not configured, in which case nothing reads the counter anyway. |
| 387 | let dropped = self.context.admin_state.as_ref() |
| 388 | .map(|a| a.dropped.clone()); |
| 389 | |
| 390 | // The total-connection ceiling. 0 -- the inert default -- never fires. |
| 391 | let max_conn = self.context.cfg.max_conn as usize; |
| 392 | |
| 393 | // A background sweep of the address guard, so a distributed flood cannot |
| 394 | // accumulate idle per-address records without bound: each source below the |
| 395 | // rate floor never throttles, so nothing else reclaims its `Monitor` log, |
| 396 | // and a botnet of such addresses is a memory-exhaustion vector on its own. |
| 397 | // The sweep evicts only an idle record with no live connection, so it can |
| 398 | // never reset a per-IP concurrency cap out from under a held connection, |
| 399 | // and a re-created record starts clean. It runs whenever the guard exists. |
| 400 | if let Some(guard) = addr_guard.clone() { |
| 401 | tokio::spawn(async move { |
| 402 | let idle = Duration::from_secs(GUARD_SWEEP_IDLE_SECS); |
| 403 | loop { |
| 404 | tokio::time::sleep( |
| 405 | Duration::from_secs(GUARD_SWEEP_INTERVAL_SECS)).await; |
| 406 | match guard.sweep_idle(idle) { |
| 407 | Ok(n) if n > 0 => debug!( |
| 408 | "addr guard swept {} idle record(s).", n), |
| 409 | Ok(_) => (), |
| 410 | Err(e) => warn!("addr guard sweep failed: {}", e), |
| 411 | } |
| 412 | } |
| 413 | }); |
| 414 | } |
| 415 | |
| 416 | // Connections being served right now. Read by the wind-up below and as |
| 417 | // the total-connection ceiling above, so an idle server pays one atomic |
| 418 | // per connection for it. |
| 419 | let inflight = Arc::new(AtomicUsize::new(0)); |
| 420 | |
| 421 | loop { |
| 422 | // Two ways out of a wait that used to have none. Both arms are |
| 423 | // cancel-safe -- tokio documents `accept` as such, and a stop |
| 424 | // wait dropped part way through has consumed nothing -- which |
| 425 | // is what makes selecting on them legitimate rather than a way |
| 426 | // of losing every other connection. |
| 427 | let accepted = tokio::select! { |
| 428 | () = stop::wait() => break, |
| 429 | accepted = listener.accept() => accepted, |
| 430 | }; |
| 431 | let (stream, src_addr) = match accepted { |
| 432 | Ok(pair) => pair, |
| 433 | Err(e) => { |
| 434 | error!(err!(e, "TCP connection aborted."; IO, Network)); |
| 435 | continue; |
| 436 | } |
| 437 | }; |
| 438 | |
| 439 | // Address-guard *rate* check runs before the TLS handshake so that a |
| 440 | // blacklisted attacker costs the server only a TCP SYN/ACK. |
| 441 | // Absent admin state (unusual -- only happens when the admin |
| 442 | // dashboard is not configured) the guard is skipped entirely. |
| 443 | if let Some(guard) = addr_guard.as_ref() { |
| 444 | match guard.check(&src_addr.ip()) { |
| 445 | Ok(decision) if decision.should_drop() => { |
| 446 | debug!("addr guard dropped TCP from {}: {:?}", |
| 447 | src_addr, decision); |
| 448 | drop(stream); |
| 449 | if let Some(d) = &dropped { d.incr(); } |
| 450 | continue; |
| 451 | } |
| 452 | Ok(_) => (), |
| 453 | Err(e) => { |
| 454 | warn!("addr guard error for {}: {}", src_addr, e); |
| 455 | } |
| 456 | } |
| 457 | } |
| 458 | |
| 459 | // Total-connection ceiling. The accept loop is the sole incrementer |
| 460 | // and `InFlight::begin` below runs before the next accept, so the |
| 461 | // read here cannot overshoot the cap. A distributed flood now hits a |
| 462 | // hard limit instead of spawning unbounded tasks. |
| 463 | if max_conn != 0 && inflight.load(Ordering::Relaxed) >= max_conn { |
| 464 | debug!("total-connection cap {} reached; dropping TCP from {}.", |
| 465 | max_conn, src_addr); |
| 466 | drop(stream); |
| 467 | if let Some(d) = &dropped { d.incr(); } |
| 468 | continue; |
| 469 | } |
| 470 | |
| 471 | // Per-IP *concurrency* permit, orthogonal to the rate check above. |
| 472 | // `None` means the address is already at its concurrency cap: drop |
| 473 | // the connection. The permit is moved into the serving task so it |
| 474 | // lives exactly as long as the connection and decrements on drop. |
| 475 | let conn_permit = match addr_guard.as_ref() { |
| 476 | Some(guard) => match guard.acquire(&src_addr.ip()) { |
| 477 | Ok(Some(permit)) => Some(permit), |
| 478 | Ok(None) => { |
| 479 | debug!("per-IP concurrency cap reached; dropping TCP from {}.", |
| 480 | src_addr); |
| 481 | drop(stream); |
| 482 | if let Some(d) = &dropped { d.incr(); } |
| 483 | continue; |
| 484 | } |
| 485 | Err(e) => { |
| 486 | warn!("addr guard acquire error for {}: {}", src_addr, e); |
| 487 | None |
| 488 | } |
| 489 | }, |
| 490 | None => None, |
| 491 | }; |
| 492 | |
| 493 | // ── Per-connection processing ──────────────────────── |
| 494 | // |
| 495 | // Spawn immediately so the accept loop is never blocked |
| 496 | // by a slow client, an incomplete TLS handshake, or a |
| 497 | // scanner that connects without sending data. Without |
| 498 | // this, a single hung peek() or TLS accept() would |
| 499 | // prevent all new connections from being accepted. |
| 500 | let context_clone = self.context.clone(); |
| 501 | let bounded_conn = bounded_acceptor.clone(); |
| 502 | let dropped_conn = dropped.clone(); |
| 503 | // Counted from here until the task ends, so that a stop can |
| 504 | // wait for whatever is mid-response rather than cutting it |
| 505 | // off. Taken before the spawn: taking it inside would leave a |
| 506 | // connection uncounted for as long as the runtime took to get |
| 507 | // round to the task, which is exactly when a wind-up is busiest. |
| 508 | let counted = stop::InFlight::begin(&inflight); |
| 509 | tokio::spawn(async move { |
| 510 | let _counted = counted; |
| 511 | // Held for the life of the connection; decrements the per-IP |
| 512 | // concurrency count when this task ends, however it ends. |
| 513 | let _conn_permit = conn_permit; |
| 514 | // Peek at first bytes to detect TLS handshake, deadlined to the |
| 515 | // same bound as the handshake so a client that opens a TLS port |
| 516 | // and never speaks cannot pin the task and its per-IP slot. |
| 517 | let mut peek_buf = [0u8; 5]; |
| 518 | let peeked = match bounded_conn.deadline() { |
| 519 | Some(d) => match tokio::time::timeout( |
| 520 | d, stream.peek(&mut peek_buf)).await |
| 521 | { |
| 522 | Ok(r) => r, |
| 523 | Err(_) => { |
| 524 | debug!("peek from {} did not arrive within the \ |
| 525 | handshake deadline; dropping.", src_addr); |
| 526 | if let Some(dc) = &dropped_conn { dc.incr(); } |
| 527 | return; |
| 528 | } |
| 529 | }, |
| 530 | None => stream.peek(&mut peek_buf).await, |
| 531 | }; |
| 532 | match peeked { |
| 533 | Ok(n) if n >= 5 && peek_buf[0] == 0x16 && peek_buf[1] == 0x03 => { |
| 534 | match bounded_conn.accept(stream).await { |
| 535 | Handshake::Ok(tls_stream) => { |
| 536 | // Extract SNI now, before we hand |
| 537 | // ownership of the stream to the handler. |
| 538 | let sni = tls_stream.get_ref().1.server_name() |
| 539 | .map(|s| s.to_string()); |
| 540 | match &context_clone.protocol { |
| 541 | Protocol::Web { .. } => { |
| 542 | if let Err(e) = context_clone.handle_https( |
| 543 | tls_stream, |
| 544 | sni, |
| 545 | src_addr, |
| 546 | ).await { |
| 547 | error!(err!(e, |
| 548 | "Error handling HTTPS connection."; |
| 549 | IO, Network)); |
| 550 | } |
| 551 | } |
| 552 | } |
| 553 | } |
| 554 | Handshake::TimedOut => { |
| 555 | debug!("TLS handshake from {} exceeded the \ |
| 556 | deadline; dropping.", src_addr); |
| 557 | if let Some(dc) = &dropped_conn { dc.incr(); } |
| 558 | } |
| 559 | Handshake::Failed(e) => { |
| 560 | error!(err!(e, |
| 561 | "TLS handshake aborted."; |
| 562 | IO, Network, Init)); |
| 563 | } |
| 564 | } |
| 565 | } |
| 566 | _ => { |
| 567 | // Non-TLS connection on the TLS port: redirect |
| 568 | // to HTTPS using the incoming Host header so the |
| 569 | // redirect target matches whatever the client |
| 570 | // typed. |
| 571 | let https_port = context_clone.cfg.server_port_tcp; |
| 572 | if let Err(e) = handle_redirect( |
| 573 | stream, |
| 574 | src_addr, |
| 575 | https_port, |
| 576 | ).await { |
| 577 | error!(err!(e, |
| 578 | "Failed to redirect plaintext HTTP on TLS port."; |
| 579 | IO, Network, Write)); |
| 580 | } |
| 581 | } |
| 582 | } |
| 583 | }); |
| 584 | } |
| 585 | |
| 586 | // ── The wind-up ────────────────────────────────────────── |
| 587 | // |
| 588 | // Reached only by the stop arm above: the loop has no other |
| 589 | // break, and an accept that errors goes round again as it always |
| 590 | // did. The listener is dropped here, so the port is given back |
| 591 | // before anything slow happens and a restart is not refused by |
| 592 | // the machine still holding it. |
| 593 | drop(listener); |
| 594 | info!("Stop requested; no longer accepting connections on {}.", addr); |
| 595 | |
| 596 | // What is already in flight is allowed to finish, within a stated |
| 597 | // bound rather than for as long as it likes. See `DRAIN_SECS` for |
| 598 | // what the bound costs and why it is not longer. |
| 599 | let left = stop::drain(&inflight, Duration::from_secs(DRAIN_SECS)).await; |
| 600 | if left > 0 { |
| 601 | warn!("{} connection(s) were still open after {} seconds and are \ |
| 602 | being dropped. Most of these will be idle ones waiting to be \ |
| 603 | reused, which costs nothing to drop; anything still being \ |
| 604 | answered after that long is a large download or a WebSocket.", |
| 605 | left, DRAIN_SECS); |
| 606 | } else { |
| 607 | info!("Every connection in flight finished."); |
| 608 | } |
| 609 | |
| 610 | // The other listeners this server spawned -- the plaintext |
| 611 | // redirect, the localhost admin port, the mail listeners -- are |
| 612 | // tasks on the runtime the caller shuts down next, and go with |
| 613 | // it. None of them holds a store, and the worst a cut costs is a |
| 614 | // redirect that has to be asked for again or an SMTP session the |
| 615 | // sender will retry; a message is only acknowledged once it is on |
| 616 | // disk, so nothing accepted is lost. |
| 617 | Ok(()) |
| 618 | } |
| 619 | } |
| 620 | |
| 621 | |
| 622 | async fn spawn_mail_listeners( |
| 623 | cfg: &MailConfig, |
| 624 | root: &oxedyne_fe2o3_core::path::NormPathBuf, |
| 625 | tls_acceptor: BoundedTlsAcceptor, |
| 626 | bind_address: &str, |
| 627 | tally: Option<Arc<ListenerTally>>, |
| 628 | ) |
| 629 | -> Outcome<()> |
| 630 | { |
| 631 | use oxedyne_fe2o3_core::path::NormalPath; |
| 632 | use std::path::{Path, PathBuf}; |
| 633 | |
| 634 | // Every listener asked for is counted before anything below can fail, so a |
| 635 | // set-up that dies before binding shows each of them as down rather than |
| 636 | // showing no mail server at all. |
| 637 | if let Some(t) = &tally { |
| 638 | for port in [cfg.smtp_port, cfg.submission_port, cfg.imap_port] { |
| 639 | if port != 0 { |
| 640 | t.want(); |
| 641 | } |
| 642 | } |
| 643 | } |
| 644 | |
| 645 | // Resolve paths under the app root. |
| 646 | let resolve = |rel: &str| -> PathBuf { |
| 647 | let p = Path::new(rel); |
| 648 | if p.is_absolute() { |
| 649 | p.to_path_buf() |
| 650 | } else { |
| 651 | root.clone().join(p.normalise()).absolute().as_pathbuf() |
| 652 | } |
| 653 | }; |
| 654 | |
| 655 | // Maildir storage root. |
| 656 | let maildir_root = if cfg.maildir_root.is_empty() { |
| 657 | return Err(err!( |
| 658 | "MailConfig: maildir_root must be set."; Invalid, Input, Missing)); |
| 659 | } else { |
| 660 | resolve(&cfg.maildir_root) |
| 661 | }; |
| 662 | if !maildir_root.exists() { |
| 663 | if let Err(e) = std::fs::create_dir_all(&maildir_root) { |
| 664 | return Err(err!(e, |
| 665 | "Creating maildir root {:?}.", maildir_root; |
| 666 | IO, File, Init)); |
| 667 | } |
| 668 | } |
| 669 | let store = res!(MaildirStore::new(maildir_root, cfg.hostname.clone())); |
| 670 | |
| 671 | // User store. |
| 672 | if cfg.users_file_rel.is_empty() { |
| 673 | return Err(err!( |
| 674 | "MailConfig: users_file_rel must be set."; Invalid, Input, Missing)); |
| 675 | } |
| 676 | let users = PasswdFileUserStore::new(resolve(&cfg.users_file_rel)); |
| 677 | |
| 678 | // Outbound spool. |
| 679 | if cfg.spool_dir_rel.is_empty() { |
| 680 | return Err(err!( |
| 681 | "MailConfig: spool_dir_rel must be set."; Invalid, Input, Missing)); |
| 682 | } |
| 683 | let spool = res!(OutboundSpool::new(resolve(&cfg.spool_dir_rel))); |
| 684 | |
| 685 | // DKIM signers. Shared with the alerter, which signs its own mail: an |
| 686 | // unsigned alert from a domain that signs everything else is exactly the |
| 687 | // message a spam filter is entitled to distrust, and it is the one message |
| 688 | // that has to arrive. |
| 689 | let dkim = res!(load_dkim_signers(cfg, root)); |
| 690 | |
| 691 | let handler = AppMailHandler { |
| 692 | store: store.clone(), |
| 693 | users: users.clone(), |
| 694 | spool: spool.clone(), |
| 695 | dkim, |
| 696 | local_domains: Arc::new(cfg.local_domains.clone()), |
| 697 | }; |
| 698 | |
| 699 | let hostname = Arc::new(cfg.hostname.clone()); |
| 700 | |
| 701 | // Build the SMTP servers. |
| 702 | let (recv_server, sub_server) = build_smtp_servers( |
| 703 | handler, |
| 704 | users.clone(), |
| 705 | Some(tls_acceptor.clone()), |
| 706 | hostname.clone(), |
| 707 | ); |
| 708 | |
| 709 | // Bind addresses. |
| 710 | let bind_ip: std::net::IpAddr = match bind_address.parse() { |
| 711 | Ok(ip) => ip, |
| 712 | Err(e) => return Err(err!(e, |
| 713 | "Invalid server_address '{}'.", bind_address; |
| 714 | Invalid, Input, Network)), |
| 715 | }; |
| 716 | |
| 717 | // SMTP receive on port 25. |
| 718 | if cfg.smtp_port != 0 { |
| 719 | let addr = SocketAddr::new(bind_ip, cfg.smtp_port); |
| 720 | let server = recv_server.clone(); |
| 721 | let tally = tally.clone(); |
| 722 | tokio::spawn(async move { |
| 723 | if let Err(e) = run_smtp_listener(addr, server, tally).await { |
| 724 | error!(err!(e, "SMTP receive listener exited."; IO, Network)); |
| 725 | } |
| 726 | }); |
| 727 | } |
| 728 | // SMTP submission on port 587. |
| 729 | if cfg.submission_port != 0 { |
| 730 | let addr = SocketAddr::new(bind_ip, cfg.submission_port); |
| 731 | let server = sub_server.clone(); |
| 732 | let tally = tally.clone(); |
| 733 | tokio::spawn(async move { |
| 734 | if let Err(e) = run_smtp_listener(addr, server, tally).await { |
| 735 | error!(err!(e, "SMTP submission listener exited."; IO, Network)); |
| 736 | } |
| 737 | }); |
| 738 | } |
| 739 | // IMAP on port 993. |
| 740 | if cfg.imap_port != 0 { |
| 741 | let addr = SocketAddr::new(bind_ip, cfg.imap_port); |
| 742 | let server: AppImapServer = ImapServer { |
| 743 | store, |
| 744 | users, |
| 745 | hostname: hostname.clone(), |
| 746 | }; |
| 747 | let acceptor = tls_acceptor.clone(); |
| 748 | let tally = tally.clone(); |
| 749 | tokio::spawn(async move { |
| 750 | if let Err(e) = run_imap_listener(addr, acceptor, server, tally).await { |
| 751 | error!(err!(e, "IMAP listener exited."; IO, Network)); |
| 752 | } |
| 753 | }); |
| 754 | } |
| 755 | // Outbound delivery worker. |
| 756 | let client = res!(OutboundClient::with_system_roots(cfg.hostname.clone())); |
| 757 | tokio::spawn(async move { |
| 758 | if let Err(e) = run_outbound_worker(spool, client).await { |
| 759 | error!(err!(e, "Outbound delivery worker exited."; IO, Network)); |
| 760 | } |
| 761 | }); |
| 762 | |
| 763 | Ok(()) |
| 764 | } |
| 765 | |
| 766 | pub fn load_dkim_signers( |
| 767 | cfg: &MailConfig, |
| 768 | root: &oxedyne_fe2o3_core::path::NormPathBuf, |
| 769 | ) |
| 770 | -> Outcome<Vec<Arc<DkimSigner>>> |
| 771 | { |
| 772 | use oxedyne_fe2o3_core::path::NormalPath; |
| 773 | use std::path::{Path, PathBuf}; |
| 774 | |
| 775 | let resolve = |rel: &str| -> PathBuf { |
| 776 | let p = Path::new(rel); |
| 777 | if p.is_absolute() { |
| 778 | p.to_path_buf() |
| 779 | } else { |
| 780 | root.clone().join(p.normalise()).absolute().as_pathbuf() |
| 781 | } |
| 782 | }; |
| 783 | |
| 784 | let mut dkim: Vec<Arc<DkimSigner>> = Vec::new(); |
| 785 | |
| 786 | let domain = if cfg.dkim_domain.is_empty() { |
| 787 | cfg.hostname.clone() |
| 788 | } else { |
| 789 | cfg.dkim_domain.clone() |
| 790 | }; |
| 791 | |
| 792 | // The ed25519 key, generated in tree if absent. |
| 793 | if !cfg.dkim_key_file.is_empty() { |
| 794 | let path = resolve(&cfg.dkim_key_file); |
| 795 | let selector = if cfg.dkim_selector.is_empty() { |
| 796 | "default".to_string() |
| 797 | } else { |
| 798 | cfg.dkim_selector.clone() |
| 799 | }; |
| 800 | let bytes = if path.exists() { |
| 801 | // Tighten a key that predates `save_secret` before reading it. |
| 802 | res!(core_file::restrict_secret(&path)); |
| 803 | match std::fs::read(&path) { |
| 804 | Ok(b) => b, |
| 805 | Err(e) => return Err(err!(e, |
| 806 | "Reading DKIM key {:?}.", path; |
| 807 | IO, File, Read)), |
| 808 | } |
| 809 | } else { |
| 810 | info!("DKIM: no key at {:?}, generating a fresh ed25519 pair.", path); |
| 811 | let s = res!(DkimSigner::generate(domain.clone(), selector.clone())); |
| 812 | if let Some(parent) = path.parent() { |
| 813 | // Key directory: 0700, not the default create mode, and the |
| 814 | // error propagates rather than being swallowed -- a failure |
| 815 | // here means the save below fails too, but at a less specific |
| 816 | // error site. |
| 817 | res!(core_file::create_secret_dir(parent)); |
| 818 | } |
| 819 | // Key material: 0600 whatever the umask, not the default create mode. |
| 820 | if let Err(e) = core_file::save_secret(&path, s.pkcs8_bytes()) { |
| 821 | return Err(err!(e, |
| 822 | "Writing DKIM key {:?}.", path; |
| 823 | IO, File, Write)); |
| 824 | } |
| 825 | s.pkcs8_bytes().to_vec() |
| 826 | }; |
| 827 | let signer = res!(DkimSigner::from_pkcs8(&bytes, domain.clone(), selector.clone())); |
| 828 | info!("DKIM ({}) TXT record for {}._domainkey.{}: {}", |
| 829 | signer.algorithm(), selector, domain, signer.dns_txt_record()); |
| 830 | dkim.push(Arc::new(signer)); |
| 831 | } |
| 832 | |
| 833 | // The RSA key. Never generated -- `ring` will not generate RSA keys, and |
| 834 | // hand-rolling the arithmetic to do so is not a road worth taking to save |
| 835 | // one command: |
| 836 | // |
| 837 | // openssl genpkey -algorithm RSA -pkeyopt rsa_keygen_bits:2048 \ |
| 838 | // -outform DER -out mail/dkim_rsa.key |
| 839 | if !cfg.dkim_rsa_key_file.is_empty() { |
| 840 | let path = resolve(&cfg.dkim_rsa_key_file); |
| 841 | let selector = if cfg.dkim_rsa_selector.is_empty() { |
| 842 | "rsa".to_string() |
| 843 | } else { |
| 844 | cfg.dkim_rsa_selector.clone() |
| 845 | }; |
| 846 | if !path.exists() { |
| 847 | return Err(err!( |
| 848 | "MailConfig: dkim_rsa_key_file {:?} does not exist. An RSA DKIM \ |
| 849 | key is generated once, offline: `openssl genpkey -algorithm RSA \ |
| 850 | -pkeyopt rsa_keygen_bits:2048 -outform DER -out {:?}`. Steel will \ |
| 851 | not generate it, because ring will not.", |
| 852 | path, path; |
| 853 | Init, Missing, File)); |
| 854 | } |
| 855 | let bytes = match std::fs::read(&path) { |
| 856 | Ok(b) => b, |
| 857 | Err(e) => return Err(err!(e, |
| 858 | "Reading RSA DKIM key {:?}.", path; |
| 859 | IO, File, Read)), |
| 860 | }; |
| 861 | let signer = res!(DkimSigner::from_pkcs8(&bytes, domain.clone(), selector.clone())); |
| 862 | info!("DKIM ({}) TXT record for {}._domainkey.{}: {}", |
| 863 | signer.algorithm(), selector, domain, signer.dns_txt_record()); |
| 864 | dkim.push(Arc::new(signer)); |
| 865 | } |
| 866 | |
| 867 | if dkim.is_empty() { |
| 868 | warn!("Mail is enabled with no DKIM key. Outbound mail will be unsigned, \ |
| 869 | and receivers will judge it on SPF alone."); |
| 870 | } |
| 871 | Ok(dkim) |
| 872 | } |