oxedyne/fe2o3/fe2o3_steel/src/app/server.rs
48.6 KiB, 278 runs
created by r1870400018:955, 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::{ |
| 3 | self, |
| 4 | constant as app_const, |
| 5 | https::AppWebHandler, |
| 6 | repl::AppShellContext, |
| 7 | }, |
| 8 | srv::{ |
| 9 | admin::{ |
| 10 | state::AdminState, |
| 11 | traffic::TrafficRecorder, |
| 12 | }, |
| 13 | alert::{ |
| 14 | AlertEvent, |
| 15 | Alerter, |
| 16 | }, |
| 17 | cert::Certificate, |
| 18 | cfg::ServerConfig, |
| 19 | constant as srv_const, |
| 20 | context::{ |
| 21 | new_db, |
| 22 | Protocol, |
| 23 | ServerContext, |
| 24 | VhostDbSpec, |
| 25 | VhostDbs, |
| 26 | VhostRuntime, |
| 27 | }, |
| 28 | dev::{ |
| 29 | cfg::DevConfig, |
| 30 | refresh::DevRefreshManager, |
| 31 | }, |
| 32 | id, |
| 33 | server::Server, |
| 34 | stop, |
| 35 | tiles::{ |
| 36 | TileService, |
| 37 | TileSource, |
| 38 | }, |
| 39 | ws::{ |
| 40 | handler::AppWebSocketHandler, |
| 41 | syntax::WebSocketSyntax, |
| 42 | }, |
| 43 | }, |
| 44 | }; |
| 45 | |
| 46 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 47 | use oxedyne_fe2o3_hash::{ |
| 48 | csum::ChecksumScheme, |
| 49 | hash::HashScheme, |
| 50 | }; |
| 51 | use oxedyne_fe2o3_o3db_sync::O3db; |
| 52 | |
| 53 | use oxedyne_fe2o3_core::{ |
| 54 | prelude::*, |
| 55 | log::{ |
| 56 | bot::FileConfig, |
| 57 | console::{ |
| 58 | LoggerConsole, |
| 59 | StdoutLoggerConsole, |
| 60 | }, |
| 61 | }, |
| 62 | path::NormalPath, |
| 63 | }; |
| 64 | use oxedyne_fe2o3_jdat::{ |
| 65 | prelude::*, |
| 66 | }; |
| 67 | use oxedyne_fe2o3_syntax::{ |
| 68 | msg::{ |
| 69 | MsgCmd, |
| 70 | }, |
| 71 | }; |
| 72 | use oxedyne_fe2o3_tui::lib_tui::{ |
| 73 | repl::{ |
| 74 | Evaluation, |
| 75 | ShellConfig, |
| 76 | }, |
| 77 | }; |
| 78 | |
| 79 | use std::{ |
| 80 | collections::HashMap, |
| 81 | path::{ |
| 82 | Path, |
| 83 | PathBuf, |
| 84 | }, |
| 85 | sync::{ |
| 86 | Arc, |
| 87 | RwLock, |
| 88 | }, |
| 89 | time::Duration, |
| 90 | }; |
| 91 | |
| 92 | use tokio; |
| 93 | |
| 94 | |
| 95 | pub const RUNTIME_STOP_SECS: u64 = 2; |
| 96 | |
| 97 | |
| 98 | impl AppShellContext { |
| 99 | |
| 100 | pub fn start_server( |
| 101 | &mut self, |
| 102 | _shell_cfg: &ShellConfig, |
| 103 | cmd: Option<&MsgCmd>, |
| 104 | ) |
| 105 | -> Outcome<Evaluation> |
| 106 | { |
| 107 | let root_path = Path::new(&self.app_cfg.app_root) |
| 108 | .normalise() |
| 109 | .absolute(); |
| 110 | |
| 111 | info!("Reading server config..."); |
| 112 | let mut server_cfg = res!(ServerConfig::from_datmap(self.app_cfg.server_cfg.clone())); |
| 113 | |
| 114 | // Resolve the health token's `{file:...}` reference at load, from a path |
| 115 | // outside the tree, so the secret compared on both sides is the file's |
| 116 | // bytes and not the literal reference text -- and so it is readable while |
| 117 | // the box is sealed. A token stored verbatim would make the path the |
| 118 | // secret, which is exactly the leak the file-homing rule exists to stop. |
| 119 | if !server_cfg.health_token.is_empty() { |
| 120 | server_cfg.health_token = res!(crate::srv::cfg::ApiRoute::resolve_file_refs( |
| 121 | &server_cfg.health_token, root_path.as_ref())); |
| 122 | } |
| 123 | // A stamp with an unusable name or a relative path is a refusal to start, not a |
| 124 | // warning: served past, it would be a field that is quietly never in the body, and an |
| 125 | // absent field trips no watcher's threshold, so the job it watches could stop unheard. |
| 126 | let health_stamps = res!(server_cfg.get_health_stamps()); |
| 127 | if !health_stamps.is_empty() { |
| 128 | info!("The health body reports the age of {} stamp(s): {}.", |
| 129 | health_stamps.len(), |
| 130 | health_stamps.iter() |
| 131 | .map(|s| fmt!("{} <- {:?}", s.field, s.path)) |
| 132 | .collect::<Vec<_>>().join(", ")); |
| 133 | } |
| 134 | |
| 135 | info!("Reading dev config..."); |
| 136 | let dev_cfg = res!(DevConfig::from_datmap(self.app_cfg.dev_cfg.clone())); |
| 137 | |
| 138 | // ┌───────────────────────┐ |
| 139 | // │ Determine mode. │ |
| 140 | // └───────────────────────┘ |
| 141 | let mut dev_mode = false; |
| 142 | if let Some(msg_cmd) = cmd { |
| 143 | if msg_cmd.has_arg("dev") { |
| 144 | dev_mode = true; |
| 145 | info!("Running in development mode."); |
| 146 | } |
| 147 | } |
| 148 | |
| 149 | // Ensure compatibility with existing websites. |
| 150 | res!(app::dev::ensure_compatibility(&root_path)); |
| 151 | |
| 152 | info!("Validating server config..."); |
| 153 | match server_cfg.validate(&root_path) { |
| 154 | Ok(()) => info!("Server configuration validated successfully."), |
| 155 | Err(e) => { |
| 156 | warn!("Server configuration validation issues: {}", e); |
| 157 | info!("Continuing with available routes..."); |
| 158 | } |
| 159 | } |
| 160 | |
| 161 | // A handshake semaphore without a deadline is a slowloris amplifier -- a |
| 162 | // permit held across a drip-fed handshake is never returned -- so this |
| 163 | // combination is a hard refusal to start, not a warning to serve past. |
| 164 | // (The general `validate` above is advisory in this harness; admission |
| 165 | // safety is not.) |
| 166 | if server_cfg.max_tls_handshakes > 0 && server_cfg.tls_handshake_timeout_ms == 0 { |
| 167 | return Err(err!( |
| 168 | "Refusing to start: max_tls_handshakes is {} but \ |
| 169 | tls_handshake_timeout_ms is 0. Bounding handshake concurrency \ |
| 170 | without a per-handshake deadline turns the bound into a \ |
| 171 | denial-of-service vector. Set tls_handshake_timeout_ms (e.g. 10000).", |
| 172 | server_cfg.max_tls_handshakes; |
| 173 | Configuration, Invalid, Input)); |
| 174 | } |
| 175 | |
| 176 | if dev_mode { |
| 177 | info!("Validating development config..."); |
| 178 | match dev_cfg.validate(&root_path) { |
| 179 | Ok(()) => info!("Development configuration validated successfully."), |
| 180 | Err(e) => { |
| 181 | warn!("Development configuration issues: {}", e); |
| 182 | info!("Some development features may be disabled."); |
| 183 | } |
| 184 | } |
| 185 | } |
| 186 | |
| 187 | if self.stat.first && !dev_mode { |
| 188 | return Ok(Evaluation::Error(fmt!( |
| 189 | "You should update values in {} before running the server in production mode.", |
| 190 | app_const::CONFIG_NAME, |
| 191 | ))); |
| 192 | } |
| 193 | |
| 194 | // ┌───────────────────────┐ |
| 195 | // │ Reconfigure logging. │ |
| 196 | // └───────────────────────┘ |
| 197 | |
| 198 | let mut log_cfg = log_get_config!(); |
| 199 | let mut logger_console = StdoutLoggerConsole::new(); |
| 200 | let logger_console_thread = logger_console.go(); |
| 201 | log_cfg.console = Some(logger_console_thread.chan.clone()); |
| 202 | (log_cfg.level, _) = res!(self.app_cfg.server_log_level()); |
| 203 | log_cfg.file = Some(FileConfig::new( |
| 204 | PathBuf::from(&root_path).join("www").join("logs"), |
| 205 | self.app_cfg.app_name.clone(), |
| 206 | "log".to_string(), |
| 207 | 0, |
| 208 | Some(1_048_576), // Activate multiple log file archiving using this max size. |
| 209 | )); |
| 210 | debug!("log_cfg = {:?}", log_cfg); |
| 211 | log_set_config!(log_cfg); |
| 212 | println!("Server now logging at {:?}", log_get_file_path!()); |
| 213 | info!("┌───────────────────────┐"); |
| 214 | info!("│ New server session. │"); |
| 215 | info!("└───────────────────────┘"); |
| 216 | |
| 217 | // ┌───────────────────────┐ |
| 218 | // │ Parse vhosts + ACME. │ |
| 219 | // └───────────────────────┘ |
| 220 | |
| 221 | let vhosts_cfg = res!(server_cfg.get_vhosts()); |
| 222 | let acme_cfg = res!(server_cfg.get_acme()); |
| 223 | |
| 224 | // ┌───────────────────────┐ |
| 225 | // │ Note per-vhost dbs. │ |
| 226 | // └───────────────────────┘ |
| 227 | |
| 228 | // One Ozone instance per vhost that has a `db_dir_rel` configured. |
| 229 | // Pure-redirect vhosts (no webroot, no database) simply have no |
| 230 | // entry. The map is keyed by each vhost's canonical (primary) |
| 231 | // hostname in lowercase, matching what `db_for_vhost` looks up on |
| 232 | // request dispatch. |
| 233 | // |
| 234 | // The databases are *not* opened here. Ozone is encrypted with the |
| 235 | // wallet master key, and Steel starts sealed -- no passphrase has |
| 236 | // been supplied, so no key exists yet. We record what to open and |
| 237 | // leave the map empty; `open_dbs_on_unseal` fills it in once an |
| 238 | // admin unseals, which may be seconds later (the operator typed a |
| 239 | // passphrase at the shell) or hours later (a cold restart, unsealed |
| 240 | // from the dashboard on someone's phone). Either way the listeners |
| 241 | // bind and the static vhosts serve in the meantime. |
| 242 | let uid = id::Uid::new(0); |
| 243 | let vhost_dbs: VhostDbs<{ id::UID_LEN }, id::Uid, O3db< |
| 244 | { id::UID_LEN }, |
| 245 | id::Uid, |
| 246 | EncryptionScheme, |
| 247 | HashScheme, |
| 248 | HashScheme, |
| 249 | ChecksumScheme, |
| 250 | >> = Arc::new(RwLock::new(HashMap::new())); |
| 251 | |
| 252 | let mut db_specs: Vec<VhostDbSpec> = Vec::new(); |
| 253 | for vh in &vhosts_cfg { |
| 254 | let primary_lc = vh.primary_hostname().to_lowercase(); |
| 255 | let db_dir = match res!(vh.get_db_dir(&root_path)) { |
| 256 | Some(p) => p, |
| 257 | None => { |
| 258 | info!("Vhost '{}' has no database configured, skipping.", |
| 259 | primary_lc); |
| 260 | continue; |
| 261 | } |
| 262 | }; |
| 263 | info!("Vhost '{}' has a database at {:?}; it will be opened on unseal.", |
| 264 | primary_lc, db_dir); |
| 265 | db_specs.push(VhostDbSpec { |
| 266 | vhost_key: primary_lc, |
| 267 | db_dir, |
| 268 | }); |
| 269 | } |
| 270 | |
| 271 | // ┌───────────────────────┐ |
| 272 | // │ Certificates. │ |
| 273 | // └───────────────────────┘ |
| 274 | |
| 275 | if dev_mode { |
| 276 | // Dev mode uses a single shared self-signed cert. |
| 277 | let tls_dir = res!(server_cfg.get_tls_dir(&root_path, true)); |
| 278 | let cert_path = tls_dir.join("fullchain.pem"); |
| 279 | let key_path = tls_dir.join("privkey.pem"); |
| 280 | debug!("dev tls_dir = {:?}", tls_dir); |
| 281 | if !cert_path.exists() || !key_path.exists() { |
| 282 | info!("Development certificates not found -- generating self-signed cert."); |
| 283 | res!(Certificate::new_dev( |
| 284 | &server_cfg, |
| 285 | &root_path, |
| 286 | )); |
| 287 | } |
| 288 | } else if acme_cfg.enabled { |
| 289 | info!("ACME is enabled; certificates will be issued on start-up via {}.", |
| 290 | acme_cfg.directory_url); |
| 291 | } else { |
| 292 | // Production without ACME: require per-vhost cert files on disk. |
| 293 | let tls_dir = res!(server_cfg.get_tls_dir(&root_path, false)); |
| 294 | for vh in &vhosts_cfg { |
| 295 | let primary = vh.primary_hostname(); |
| 296 | let vh_dir = tls_dir.join(primary); |
| 297 | let cert_path = vh_dir.join("fullchain.pem"); |
| 298 | let key_path = vh_dir.join("privkey.pem"); |
| 299 | if !cert_path.exists() || !key_path.exists() { |
| 300 | return Ok(Evaluation::Error(fmt!( |
| 301 | "Missing certificate for vhost '{}'. ACME is disabled, so \ |
| 302 | Steel expects {:?} and {:?} to already exist. Either enable \ |
| 303 | ACME in the acme section of config.jdat, or install the \ |
| 304 | required PEM files and restart.", |
| 305 | primary, cert_path, key_path, |
| 306 | ))); |
| 307 | } |
| 308 | } |
| 309 | } |
| 310 | |
| 311 | if dev_mode { |
| 312 | info!("Connect via: https://localhost:{}", server_cfg.server_port_tcp); |
| 313 | } else if let Some(first) = vhosts_cfg.first() { |
| 314 | info!("Connect via: https://{}", first.primary_hostname()); |
| 315 | } |
| 316 | |
| 317 | // ┌───────────────────────┐ |
| 318 | // │ Refresh the css and │ |
| 319 | // │ javascript bundles. │ |
| 320 | // └───────────────────────┘ |
| 321 | |
| 322 | let js_bundles_map = if dev_mode && dev_cfg.has_js_bundling(&root_path) { |
| 323 | info!("JavaScript bundling enabled."); |
| 324 | res!(dev_cfg.get_js_bundles_map(&root_path)) |
| 325 | } else { |
| 326 | info!("JavaScript bundling disabled or not configured."); |
| 327 | Vec::new() |
| 328 | }; |
| 329 | |
| 330 | let js_import_aliases = res!(dev_cfg.get_js_import_aliases(&root_path)); |
| 331 | |
| 332 | let css_paths = if dev_mode && dev_cfg.has_css_bundling(&root_path) { |
| 333 | info!("CSS compilation enabled."); |
| 334 | res!(dev_cfg.get_css_paths(&root_path)) |
| 335 | } else { |
| 336 | info!("CSS compilation disabled or not configured."); |
| 337 | (PathBuf::new(), PathBuf::new()) |
| 338 | }; |
| 339 | |
| 340 | let refresh_manager = Arc::new(DevRefreshManager::new( |
| 341 | &root_path, |
| 342 | js_bundles_map, |
| 343 | js_import_aliases, |
| 344 | css_paths, |
| 345 | )); |
| 346 | res!(refresh_manager.refresh()); |
| 347 | |
| 348 | // ┌───────────────────────┐ |
| 349 | // │ Initialise the dev │ |
| 350 | // │ refresh functionality.│ |
| 351 | // └───────────────────────┘ |
| 352 | |
| 353 | let rt = res!(tokio::runtime::Runtime::new()); |
| 354 | |
| 355 | let ws_handler = AppWebSocketHandler::new( |
| 356 | if dev_mode { |
| 357 | let manager_clone = refresh_manager.clone(); |
| 358 | // The file watcher is synchronous blocking code (notify's |
| 359 | // event loop). On single-CPU hosts Tokio's default runtime |
| 360 | // has a single worker, so a blocking task spawned with |
| 361 | // `rt.spawn` would hog that worker and starve the accept |
| 362 | // loop. Use `spawn_blocking` to run it on the dedicated |
| 363 | // blocking thread pool instead. |
| 364 | rt.spawn_blocking(move || { |
| 365 | debug!("Starting dev refresh file watcher."); |
| 366 | if let Err(e) = manager_clone.watch() { |
| 367 | error!(err!(e, |
| 368 | "Failed to start development file watcher."; |
| 369 | Init)); |
| 370 | } |
| 371 | }); |
| 372 | Some(refresh_manager) |
| 373 | } else { |
| 374 | None |
| 375 | } |
| 376 | ); |
| 377 | |
| 378 | // ┌───────────────────────┐ |
| 379 | // │ Start server. │ |
| 380 | // └───────────────────────┘ |
| 381 | |
| 382 | let ws_syntax = res!(WebSocketSyntax::new( |
| 383 | &self.app_cfg.app_human_name, |
| 384 | &app_const::VERSION, |
| 385 | &self.app_cfg.app_description, |
| 386 | )); |
| 387 | |
| 388 | // Build per-vhost runtimes: one AppWebHandler per vhost, each with its |
| 389 | // own public_dir / static_routes / default index files. The resulting |
| 390 | // map is keyed by every alias hostname so SNI dispatch finds the right |
| 391 | // runtime regardless of which alias the client asked for. |
| 392 | let mut vhost_map: HashMap<String, Arc<VhostRuntime< |
| 393 | AppWebHandler<HashMap<String, oxedyne_fe2o3_core::file::OsPath>>, |
| 394 | _, |
| 395 | >>> = HashMap::new(); |
| 396 | let mut default_vhost_key: Option<String> = None; |
| 397 | |
| 398 | // Build a TLS client config for outbound API proxy requests. |
| 399 | // Reads the system CA bundle so Steel can talk to any upstream. |
| 400 | let tls_client = match build_outbound_tls_client() { |
| 401 | Ok(cfg) => Some(cfg), |
| 402 | Err(e) => { |
| 403 | warn!("{}", e); |
| 404 | info!("Outbound API proxy will be unavailable."); |
| 405 | None |
| 406 | } |
| 407 | }; |
| 408 | |
| 409 | // Build the admin dashboard runtime. AdminState holds a |
| 410 | // shared handle to the wallet (so dashboard login uses the |
| 411 | // same admin list the CLI sees), the wallet's on-disk path |
| 412 | // (so the admin-management UI can call Wallet::save), the |
| 413 | // recovered master key (so it can call Wallet::enrol |
| 414 | // without re-prompting), an AES-256-GCM cipher pre-keyed |
| 415 | // with a SHA3-256 derivation of the master key (for |
| 416 | // session cookies), and the shared TrafficRecorder. Both |
| 417 | // AdminState and ServerContext hold the same Arc to the |
| 418 | // recorder so dashboard reads and request-pipeline writes |
| 419 | // see one consistent view. |
| 420 | let traffic = TrafficRecorder::new_shared(0); |
| 421 | // The sampler also reads what the health body needs beyond its snapshot -- |
| 422 | // the `health_residents` processes on a slower cadence of its own, and how |
| 423 | // full the app root's filesystem is -- so a request formats figures rather |
| 424 | // than walking `/proc`. Resident names were checked at validation. |
| 425 | let host_sampler = crate::srv::admin::host_sampler::HostSampler::new_shared_for_health( |
| 426 | res!(server_cfg.get_health_residents()), |
| 427 | Some(PathBuf::from(&root_path)), |
| 428 | ); |
| 429 | let addr_guard = res!(crate::srv::admin::guard::new_shared_with( |
| 430 | server_cfg.get_addr_guard_settings(), |
| 431 | )); |
| 432 | // Build the tighter auth-path guard from the general |
| 433 | // settings, then override the rps cap with the auth-specific |
| 434 | // value. Everything else (throttle spacing, cooldown, strike |
| 435 | // count) stays aligned with the general guard so operators |
| 436 | // only need to tune one number for the common case. |
| 437 | let mut auth_settings = server_cfg.get_addr_guard_settings(); |
| 438 | auth_settings.rps_max = server_cfg.auth_rps_max; |
| 439 | let auth_guard = res!(crate::srv::admin::guard::new_shared_with( |
| 440 | auth_settings, |
| 441 | )); |
| 442 | // Durable whitelist: addresses that bypass the rate guard entirely, |
| 443 | // written in config so they survive a restart rather than living only in |
| 444 | // the in-memory map. Applied to both the general and the auth guard, so a |
| 445 | // whitelisted address is never throttled on either path. Parsed already |
| 446 | // at validation, so an unparseable entry has failed before here. |
| 447 | let whitelist_ips = res!(server_cfg.get_whitelist_ips()); |
| 448 | for ip in &whitelist_ips { |
| 449 | res!(addr_guard.whitelist(ip)); |
| 450 | res!(auth_guard.whitelist(ip)); |
| 451 | } |
| 452 | if !whitelist_ips.is_empty() { |
| 453 | info!("Whitelisted {} address(es) from config; they bypass the rate guard \ |
| 454 | across restarts.", whitelist_ips.len()); |
| 455 | } |
| 456 | // The periodic traffic and host samplers are spawned inside |
| 457 | // Server::start (not here) because this function runs in a |
| 458 | // sync context -- the tokio runtime `rt` has been built but |
| 459 | // not entered yet, so calling `tokio::spawn` here would |
| 460 | // panic with "there is no reactor running". Server::start |
| 461 | // is invoked via `rt.block_on` which makes the runtime |
| 462 | // current for its body, so any tokio::spawn call from |
| 463 | // there works. |
| 464 | let wallet_path_for_admin = Path::new(&self.app_cfg.app_root) |
| 465 | .join(app_const::WALLET_NAME); |
| 466 | // Pull the signed-admin-login configuration off the primary |
| 467 | // vhost, if any. A Steel process can host any number of |
| 468 | // vhosts, but the admin dashboard is a single cross-vhost |
| 469 | // surface, so the admin_keys list and head_injection_url |
| 470 | // are sourced from the canonical vhost (the first entry in |
| 471 | // the config). This mirrors how the wallet and master key |
| 472 | // are already a process-wide concern. |
| 473 | let (admin_keys_cfg, head_injection_url_cfg) = { |
| 474 | let vhosts = res!(server_cfg.get_vhosts()); |
| 475 | match vhosts.into_iter().next() { |
| 476 | Some(v) => (v.admin_keys, v.head_injection_url), |
| 477 | None => (Vec::new(), None), |
| 478 | } |
| 479 | }; |
| 480 | // Operator alerting. The public hostname of the primary vhost is what |
| 481 | // an alert names, and what its `/admin` link points at. |
| 482 | // What an alert calls this machine. |
| 483 | // |
| 484 | // THE MACHINE, not the first website it happens to serve. This used to |
| 485 | // be `vhosts.first()`, so a host serving several sites announced itself |
| 486 | // as whichever one sorted first in the config -- karri reported an |
| 487 | // outage on jarrah as "seen from need2know.ai", naming a website that |
| 488 | // has nothing to do with either end of the sentence. An operator reads |
| 489 | // this half-awake and has to know which box to open. |
| 490 | // |
| 491 | // The kernel's own answer, read the way `fe2o3_sys` reads everything |
| 492 | // else about a host, with the first vhost kept only as a fallback for a |
| 493 | // system that does not carry one. |
| 494 | let alert_host = { |
| 495 | let machine = std::fs::read_to_string("/proc/sys/kernel/hostname") |
| 496 | .map(|s| s.trim().to_string()) |
| 497 | .unwrap_or_default(); |
| 498 | if !machine.is_empty() { |
| 499 | machine |
| 500 | } else { |
| 501 | let vhosts = res!(server_cfg.get_vhosts()); |
| 502 | match vhosts.first() { |
| 503 | Some(v) => v.primary_hostname().to_string(), |
| 504 | None => String::new(), |
| 505 | } |
| 506 | } |
| 507 | }; |
| 508 | let alerter = match res!(server_cfg.get_alerts()) { |
| 509 | Some(mut cfg) => { |
| 510 | // Expand `{file:...}` in the submission credential, so the |
| 511 | // password is not sitting in config.jdat in the clear. |
| 512 | res!(cfg.resolve_secrets(root_path.as_ref())); |
| 513 | // Sign alerts with the host's own DKIM keys, when it has any. |
| 514 | // An unsigned message from a domain that signs everything else |
| 515 | // is what a spam filter is entitled to distrust, and the alert |
| 516 | // is the one message that has to arrive. |
| 517 | let dkim = match res!(server_cfg.get_mail()) { |
| 518 | Some(mail_cfg) => res!( |
| 519 | crate::srv::server::load_dkim_signers(&mail_cfg, &root_path)), |
| 520 | None => Vec::new(), |
| 521 | }; |
| 522 | let to = cfg.to.join(", "); |
| 523 | let via = match &cfg.submission { |
| 524 | Some(s) => fmt!("via {}:{}", s.host, s.port), |
| 525 | None => fmt!("direct to the recipient's MX"), |
| 526 | }; |
| 527 | let signed = if dkim.is_empty() { |
| 528 | fmt!("unsigned") |
| 529 | } else { |
| 530 | fmt!("DKIM-signed with {} key(s)", dkim.len()) |
| 531 | }; |
| 532 | let a = res!(Alerter::new( |
| 533 | cfg, alert_host.clone(), dkim, tls_client.clone())); |
| 534 | info!("Alerting enabled; operator alerts go to {} ({}, {}).", |
| 535 | to, via, signed); |
| 536 | a |
| 537 | } |
| 538 | None => { |
| 539 | info!("Alerting is not configured. Steel will not tell anybody \ |
| 540 | when it comes up sealed or when its passphrase is guessed at."); |
| 541 | None |
| 542 | } |
| 543 | }; |
| 544 | |
| 545 | // Watching the other machines. This is the half of alerting that a host |
| 546 | // cannot do for itself: a machine that has died sends nothing, and |
| 547 | // silence reads exactly like health. See `srv::watch`. |
| 548 | // |
| 549 | // Requires an alerter, because a watcher with nothing to report through |
| 550 | // is a thread that discovers an outage and keeps it to itself. Requires |
| 551 | // an outbound TLS client for the same practical reason. Both absences |
| 552 | // are said out loud rather than logged at debug, since the operator who |
| 553 | // configured a watch believes they are covered. |
| 554 | // |
| 555 | // What the watcher sees goes into the shared fleet rings too, which the |
| 556 | // dashboard's Fleet page draws. A host that watches nobody still gets an |
| 557 | // empty fleet, since the page always has this host's own row to show. |
| 558 | let fleet = match res!(server_cfg.get_watch()) { |
| 559 | Some(mut wcfg) => { |
| 560 | // Resolve each peer's health-token `{file:...}` at load, the same |
| 561 | // reason the server's own token is resolved: the secret must be |
| 562 | // the file's bytes, not the reference text, or the watcher would |
| 563 | // present the literal `{file:/x}` string as its token. |
| 564 | for p in &mut wcfg.peers { |
| 565 | if let Some(tok) = &p.token { |
| 566 | if !tok.is_empty() { |
| 567 | p.token = Some(res!(crate::srv::cfg::ApiRoute::resolve_file_refs( |
| 568 | tok, root_path.as_ref()))); |
| 569 | } |
| 570 | } |
| 571 | } |
| 572 | let fleet = crate::srv::fleet::Fleet::new_shared( |
| 573 | alert_host.clone(), Some(&wcfg)); |
| 574 | match (&alerter, &tls_client) { |
| 575 | (Some(a), Some(tls)) => { |
| 576 | let w = res!(crate::srv::watch::Watcher::new( |
| 577 | Arc::new(wcfg), |
| 578 | Arc::new(a.clone()), |
| 579 | tls.clone(), |
| 580 | alert_host.to_string(), |
| 581 | fleet.clone(), |
| 582 | )); |
| 583 | // `rt.spawn`, NOT `tokio::spawn`. This function is sync: the |
| 584 | // runtime is built at the top and is not current until |
| 585 | // `block_on` further down, so `tokio::spawn` here panics with |
| 586 | // "there is no reactor running" -- exactly as the comment |
| 587 | // beside the samplers above says it will. That panic took a |
| 588 | // live server into a restart loop on 2026-08-10, because the |
| 589 | // code compiles perfectly and fails only when a server is |
| 590 | // actually started. A change to this function is not tested |
| 591 | // until something has been started with it. |
| 592 | rt.spawn(w.run()); |
| 593 | }, |
| 594 | (None, _) => warn!("A watch list is configured but alerting is not, \ |
| 595 | so nothing would be told about a peer going down. The watcher \ |
| 596 | was not started."), |
| 597 | (_, None) => warn!("A watch list is configured but this Steel has no \ |
| 598 | outbound TLS client, so it cannot probe anything. The watcher \ |
| 599 | was not started."), |
| 600 | } |
| 601 | fleet |
| 602 | }, |
| 603 | None => crate::srv::fleet::Fleet::new_shared(alert_host.clone(), None), |
| 604 | }; |
| 605 | |
| 606 | // Any failure here disables the newsletter with a warning; see `newsletter_sender`. |
| 607 | let mail_sender = newsletter_sender(res!(server_cfg.get_mail_any()), &root_path); |
| 608 | |
| 609 | let admin_state = res!(AdminState::new( |
| 610 | self.wallet.clone(), |
| 611 | wallet_path_for_admin, |
| 612 | self.db_enc_key.clone(), |
| 613 | db_specs.len(), |
| 614 | alerter, |
| 615 | traffic.clone(), |
| 616 | host_sampler.clone(), |
| 617 | addr_guard.clone(), |
| 618 | auth_guard.clone(), |
| 619 | admin_keys_cfg, |
| 620 | head_injection_url_cfg, |
| 621 | )).with_fleet(fleet).with_health_stamps(health_stamps); |
| 622 | let admin_state = Arc::new(admin_state); |
| 623 | info!("Admin dashboard runtime initialised \ |
| 624 | (traffic ring capacity {}; host sampler capacity {}).", |
| 625 | traffic.capacity(), |
| 626 | host_sampler.history_capacity()); |
| 627 | if admin_state.seal_withholds_data() { |
| 628 | warn!("Steel is SEALED: no wallet master key has been supplied, so the \ |
| 629 | {} configured database(s) are shut and routes that need them will \ |
| 630 | answer 503. Static vhosts, redirects, proxy routes and certificate \ |
| 631 | renewal all serve normally. Sign in at /admin with an admin \ |
| 632 | passphrase to unseal.", |
| 633 | db_specs.len()); |
| 634 | // The alert that matters most. The websites are up, so nothing |
| 635 | // looks wrong from outside, and without this the operator learns |
| 636 | // that the databases are shut from a user complaint. Queued on the |
| 637 | // runtime handle because `raise` spawns, and the runtime is not |
| 638 | // current until `block_on` below. |
| 639 | if let Some(a) = admin_state.alerter().cloned() { |
| 640 | let n = db_specs.len(); |
| 641 | rt.spawn(async move { |
| 642 | a.raise(AlertEvent::SealedStart { db_count: n }); |
| 643 | }); |
| 644 | } |
| 645 | } else if admin_state.is_sealed() { |
| 646 | info!("Steel is sealed -- no wallet master key has been supplied -- but \ |
| 647 | no vhost has a database configured, so nothing is waiting on it and \ |
| 648 | everything serves normally."); |
| 649 | } |
| 650 | |
| 651 | for vh in &vhosts_cfg { |
| 652 | let public_dir = match res!(vh.get_public_dir(&root_path)) { |
| 653 | Some(p) => p, |
| 654 | None => PathBuf::new(), |
| 655 | }; |
| 656 | let static_routes = res!(vh.get_static_route_paths( |
| 657 | &root_path, |
| 658 | HashMap::new(), |
| 659 | )); |
| 660 | let default_index_files = res!(vh.get_default_index_files()); |
| 661 | |
| 662 | // Resolve {file:} placeholders in API route headers and in |
| 663 | // in-process handler config (e.g. a Stripe secret key). |
| 664 | let mut api_routes = vh.api_routes.clone(); |
| 665 | for route in &mut api_routes { |
| 666 | res!(route.resolve_headers(root_path.as_ref())); |
| 667 | res!(route.resolve_config(root_path.as_ref())); |
| 668 | } |
| 669 | if !api_routes.is_empty() { |
| 670 | info!("Vhost '{}': {} API route(s) configured.", |
| 671 | vh.primary_hostname(), api_routes.len()); |
| 672 | } |
| 673 | if !vh.proxy_routes.is_empty() { |
| 674 | info!("Vhost '{}': {} proxy route(s) configured.", |
| 675 | vh.primary_hostname(), vh.proxy_routes.len()); |
| 676 | } |
| 677 | for route in &vh.ws_routes { |
| 678 | info!("Vhost '{}': ws route {} -> ws://{}:{}{}", |
| 679 | vh.primary_hostname(), route.path, |
| 680 | route.upstream_host, route.upstream_port, route.upstream_path); |
| 681 | } |
| 682 | |
| 683 | // Resolve {file:} placeholders in webhook route config. |
| 684 | let mut webhook_routes = vh.webhook_routes.clone(); |
| 685 | for route in &mut webhook_routes { |
| 686 | res!(route.resolve_config(root_path.as_ref())); |
| 687 | } |
| 688 | if !webhook_routes.is_empty() { |
| 689 | info!("Vhost '{}': {} webhook route(s) configured.", |
| 690 | vh.primary_hostname(), webhook_routes.len()); |
| 691 | } |
| 692 | |
| 693 | // Resolve {file:}/{env:} placeholders in the publish module's destination credentials, so a |
| 694 | // token reaches the sender resolved and never sits in the config in the clear. |
| 695 | let mut publish = vh.publish.clone(); |
| 696 | if let Some(p) = &mut publish { |
| 697 | res!(p.resolve_secrets(root_path.as_ref())); |
| 698 | } |
| 699 | |
| 700 | let web_handler = AppWebHandler::new( |
| 701 | server_cfg.clone(), |
| 702 | public_dir, |
| 703 | static_routes, |
| 704 | default_index_files, |
| 705 | dev_mode, |
| 706 | api_routes, |
| 707 | webhook_routes, |
| 708 | self.webhook_registry.clone(), |
| 709 | self.api_handler_registry.clone(), |
| 710 | tls_client.clone(), |
| 711 | Some(admin_state.clone()), |
| 712 | Some(traffic.clone()), |
| 713 | publish.map(Arc::new), |
| 714 | mail_sender.clone(), |
| 715 | Arc::new(vh.site_admins.clone()), |
| 716 | ); |
| 717 | |
| 718 | // One terminal manager per vhost, built only where the config asks |
| 719 | // for it. It is shared two ways: attached to the WS syntax handler |
| 720 | // for the `term_*` management commands, and held on the runtime so |
| 721 | // the `/term/` bridge router can see whether the feature exists at |
| 722 | // all -- the presence of this is the config gate that decides |
| 723 | // whether an upgrade to `/term/<name>` is dispatched or refused. |
| 724 | let term_manager = vh.term_config.as_ref().map(|tc| Arc::new( |
| 725 | crate::srv::ws::term::TerminalManager::new( |
| 726 | &tc.session_prefix, |
| 727 | &tc.launch_command, |
| 728 | ) |
| 729 | )); |
| 730 | // Every build is opened now, so a missing or unreadable archive stops start-up |
| 731 | // rather than surfacing as failed tiles. None of this names a tile. |
| 732 | let tiles = match &vh.tiles { |
| 733 | Some(tc) => { |
| 734 | let svc = res!(TileService::new( |
| 735 | tc, |
| 736 | vh.primary_hostname(), |
| 737 | server_cfg.hsts_max_age_secs as u64, |
| 738 | TileSource::open, |
| 739 | )); |
| 740 | info!("Vhost '{}': tiles under {} from {} build(s), current '{}'; \ |
| 741 | access log off.", vh.primary_hostname(), svc.prefix(), |
| 742 | tc.builds.len(), svc.current()); |
| 743 | Some(Arc::new(svc)) |
| 744 | } |
| 745 | None => None, |
| 746 | }; |
| 747 | let runtime = Arc::new(VhostRuntime { |
| 748 | hostnames: vh.hostnames.clone(), |
| 749 | web_handler, |
| 750 | ws_handler: match &term_manager { |
| 751 | Some(tm) => ws_handler.clone().with_term_manager(tm.clone()), |
| 752 | None => ws_handler.clone(), |
| 753 | }, |
| 754 | ws_syntax: ws_syntax.clone(), |
| 755 | redirects: vh.redirects.clone(), |
| 756 | proxy_routes: vh.proxy_routes.clone(), |
| 757 | ws_routes: vh.ws_routes.clone(), |
| 758 | term_manager: term_manager.clone(), |
| 759 | uses_sessions: vh.uses_sessions(), |
| 760 | permissions_policy: vh.permissions_policy.clone(), |
| 761 | tiles, |
| 762 | access_log: vh.access_log, |
| 763 | }); |
| 764 | |
| 765 | let primary_lc = vh.primary_hostname().to_lowercase(); |
| 766 | if default_vhost_key.is_none() { |
| 767 | default_vhost_key = Some(primary_lc.clone()); |
| 768 | } |
| 769 | for h in &vh.hostnames { |
| 770 | vhost_map.insert(h.to_lowercase(), runtime.clone()); |
| 771 | } |
| 772 | } |
| 773 | |
| 774 | let default_vhost = match default_vhost_key { |
| 775 | Some(k) => k, |
| 776 | None => return Ok(Evaluation::Error(fmt!( |
| 777 | "No vhosts configured -- at least one vhost is required."))), |
| 778 | }; |
| 779 | |
| 780 | let protocol = Protocol::Web { |
| 781 | vhosts: Arc::new(vhost_map), |
| 782 | default_vhost, |
| 783 | dev_mode, |
| 784 | }; |
| 785 | |
| 786 | let server_context = ServerContext::new( |
| 787 | server_cfg, |
| 788 | root_path.clone(), |
| 789 | vhost_dbs.clone(), |
| 790 | db_specs.clone(), |
| 791 | protocol, |
| 792 | Some(traffic.clone()), |
| 793 | Some(admin_state.clone()), |
| 794 | ); |
| 795 | |
| 796 | let server = Server::new(server_context); |
| 797 | |
| 798 | // From here until the stores are shut, a stop request is answered |
| 799 | // by the wind-up below rather than by the listener ending the |
| 800 | // process itself. Set before the opener is spawned, not after the |
| 801 | // listeners bind, so that a signal arriving while a store is being |
| 802 | // opened still finds somebody who will close it. |
| 803 | stop::serving(true); |
| 804 | |
| 805 | // The database opener waits for the master key and then opens |
| 806 | // every configured Ozone instance. Spawned on the runtime handle |
| 807 | // *before* `block_on` drives it, so it is already waiting on the |
| 808 | // unseal signal by the time the first request can arrive. |
| 809 | // |
| 810 | // When the operator has already supplied a passphrase (via the |
| 811 | // shell's `unseal`, or `STEEL_ADMIN_PASS`) the key is present and |
| 812 | // this returns without waiting -- so the familiar start-up path is |
| 813 | // unchanged apart from the databases now opening in parallel with |
| 814 | // the listeners binding, rather than before them. |
| 815 | // Cloned, not moved: the same map is what the wind-up reads to |
| 816 | // find the stores it has to close. |
| 817 | rt.spawn(open_dbs_on_unseal( |
| 818 | admin_state.clone(), |
| 819 | vhost_dbs.clone(), |
| 820 | db_specs, |
| 821 | uid, |
| 822 | )); |
| 823 | |
| 824 | info!("Starting server..."); |
| 825 | for line in srv_const::SPLASH.lines() { |
| 826 | info!("{}", line); |
| 827 | } |
| 828 | |
| 829 | match rt.block_on(server.start()) { |
| 830 | Ok(()) => info!("Server stopped gracefully."), |
| 831 | Err(e) => error!(Error::Upstream(Arc::new(e), ErrMsg { |
| 832 | tags: &[ErrTag::IO, ErrTag::Thread], |
| 833 | msg: fmt!("Result of the attempt to execute the server within the Tokio runtime."), |
| 834 | })), |
| 835 | } |
| 836 | |
| 837 | // The runtime is shut down, not dropped. Dropping a runtime waits |
| 838 | // for its blocking pool, and in dev mode that pool holds the file |
| 839 | // watcher -- `notify`'s own event loop, which by design never |
| 840 | // returns. A stop therefore came all the way home, closed |
| 841 | // everything, said so, and then hung for ever in the drop, one |
| 842 | // line before the end of the function. Nothing here needs those |
| 843 | // threads: what mattered has already finished inside `block_on`, |
| 844 | // and asynchronous tasks are dropped at their await points. |
| 845 | rt.shutdown_timeout(Duration::from_secs(RUNTIME_STOP_SECS)); |
| 846 | |
| 847 | // The point of all of it. Whatever brought the server home -- |
| 848 | // a stop request or a failure -- every store it holds is shut |
| 849 | // before the process ends, rather than being killed open. After |
| 850 | // the runtime has gone, so the periodic persist task cannot be |
| 851 | // writing into a database while it is being closed. |
| 852 | let closed = close_vhost_dbs(&vhost_dbs); |
| 853 | info!("{} database(s) closed.", closed); |
| 854 | stop::serving(false); |
| 855 | |
| 856 | log_finish_wait!(); |
| 857 | |
| 858 | Ok(Evaluation::Exit) |
| 859 | } |
| 860 | |
| 861 | } |
| 862 | |
| 863 | pub fn build_outbound_tls_client() |
| 864 | -> Outcome<Arc<tokio_rustls::rustls::ClientConfig>> |
| 865 | { |
| 866 | use tokio_rustls::rustls::{ |
| 867 | ClientConfig, |
| 868 | RootCertStore, |
| 869 | pki_types::CertificateDer, |
| 870 | }; |
| 871 | |
| 872 | // Common system CA bundle paths. |
| 873 | let ca_paths = [ |
| 874 | "/etc/ssl/certs/ca-certificates.crt", // Debian/Ubuntu |
| 875 | "/etc/pki/tls/certs/ca-bundle.crt", // Fedora/RHEL |
| 876 | "/etc/ssl/cert.pem", // Alpine/macOS |
| 877 | ]; |
| 878 | let ca_file = match ca_paths.iter().find(|p| Path::new(p).exists()) { |
| 879 | Some(p) => *p, |
| 880 | None => return Err(err!( |
| 881 | "No system CA bundle found. Tried: {:?}", ca_paths; |
| 882 | Init, Missing, File)), |
| 883 | }; |
| 884 | |
| 885 | info!("Loading system CA certificates from '{}'...", ca_file); |
| 886 | let pem_data = match std::fs::read(ca_file) { |
| 887 | Ok(d) => d, |
| 888 | Err(e) => return Err(err!(e, |
| 889 | "Failed to read CA bundle '{}'.", ca_file; |
| 890 | IO, File, Read)), |
| 891 | }; |
| 892 | |
| 893 | let mut store = RootCertStore::empty(); |
| 894 | let mut count = 0u32; |
| 895 | // Parse PEM-encoded certificates. |
| 896 | let mut cursor = &pem_data[..]; |
| 897 | loop { |
| 898 | match rustls_pemfile::read_one(&mut cursor) { |
| 899 | Ok(Some(rustls_pemfile::Item::X509Certificate(cert))) => { |
| 900 | let der = CertificateDer::from(cert); |
| 901 | match store.add(der) { |
| 902 | Ok(()) => count += 1, |
| 903 | Err(_) => (), // Skip malformed certs silently. |
| 904 | } |
| 905 | } |
| 906 | Ok(Some(_)) => continue, // Skip non-certificate items. |
| 907 | Ok(None) => break, // End of file. |
| 908 | Err(_) => break, // Parse error; stop. |
| 909 | } |
| 910 | } |
| 911 | if count == 0 { |
| 912 | return Err(err!( |
| 913 | "CA bundle '{}' contained no usable certificates.", ca_file; |
| 914 | Init, Invalid, File)); |
| 915 | } |
| 916 | info!("Loaded {} CA certificate(s) for outbound HTTPS.", count); |
| 917 | |
| 918 | let mut config = ClientConfig::builder() |
| 919 | .with_root_certificates(store) |
| 920 | .with_no_client_auth(); |
| 921 | // Advertise HTTP/1.1 via ALPN so CDN-fronted servers (e.g. |
| 922 | // Fireworks.ai behind Cloudflare) don't close the connection |
| 923 | // after the TLS handshake when no protocol is negotiated. |
| 924 | config.alpn_protocols = vec![b"http/1.1".to_vec()]; |
| 925 | Ok(Arc::new(config)) |
| 926 | } |
| 927 | |
| 928 | |
| 929 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 930 | // │ DATABASE OPENER AND CLOSER │ |
| 931 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 932 | |
| 933 | type VhostDb = O3db< |
| 934 | { id::UID_LEN }, |
| 935 | id::Uid, |
| 936 | EncryptionScheme, |
| 937 | HashScheme, |
| 938 | HashScheme, |
| 939 | ChecksumScheme, |
| 940 | >; |
| 941 | |
| 942 | fn close_vhost_dbs( |
| 943 | vhost_dbs: &VhostDbs<{ id::UID_LEN }, id::Uid, VhostDb>, |
| 944 | ) -> usize { |
| 945 | let dbs = lock_read_or_recover!(vhost_dbs, |
| 946 | "The per-vhost database map is poisoned. Closing what it holds anyway: \ |
| 947 | an open store is the thing worth minding here."); |
| 948 | if dbs.is_empty() { |
| 949 | info!("No database is open, so there is none to close."); |
| 950 | return 0; |
| 951 | } |
| 952 | let mut closed = 0; |
| 953 | for (vhost, (db, _uid)) in dbs.iter() { |
| 954 | info!("Closing the database for vhost '{}'...", vhost); |
| 955 | let db = lock_read_or_recover!(db, |
| 956 | "The database handle for vhost '{}' is poisoned. Closing it anyway.", |
| 957 | vhost); |
| 958 | match db.close() { |
| 959 | Ok(()) => { |
| 960 | info!("Closed the database for vhost '{}'.", vhost); |
| 961 | closed += 1; |
| 962 | }, |
| 963 | Err(e) => error!(e, |
| 964 | "Closing the database for vhost '{}'. Its files may need a \ |
| 965 | repair pass on the next start.", vhost), |
| 966 | } |
| 967 | } |
| 968 | closed |
| 969 | } |
| 970 | |
| 971 | async fn open_dbs_on_unseal( |
| 972 | admin_state: Arc<AdminState>, |
| 973 | vhost_dbs: VhostDbs<{ id::UID_LEN }, id::Uid, O3db< |
| 974 | { id::UID_LEN }, |
| 975 | id::Uid, |
| 976 | EncryptionScheme, |
| 977 | HashScheme, |
| 978 | HashScheme, |
| 979 | ChecksumScheme, |
| 980 | >>, |
| 981 | db_specs: Vec<VhostDbSpec>, |
| 982 | uid: id::Uid, |
| 983 | ) { |
| 984 | if db_specs.is_empty() { |
| 985 | return; |
| 986 | } |
| 987 | |
| 988 | let enc_key = match admin_state.await_master_key().await { |
| 989 | Ok(k) => k, |
| 990 | Err(e) => { |
| 991 | error!(e, "Waiting for the wallet master key. No database will \ |
| 992 | be opened; Steel remains sealed."); |
| 993 | return; |
| 994 | } |
| 995 | }; |
| 996 | |
| 997 | let outcome = tokio::task::spawn_blocking(move || -> Outcome<usize> { |
| 998 | let mut opened = 0; |
| 999 | for spec in &db_specs { |
| 1000 | info!("Starting database for vhost '{}' at {:?}...", |
| 1001 | spec.vhost_key, spec.db_dir); |
| 1002 | let mut db = res!(new_db(&spec.db_dir, &enc_key)); |
| 1003 | // Label is the vhost's canonical name so log output |
| 1004 | // disambiguates multi-vhost deployments. |
| 1005 | let label = fmt!("db_{}", spec.vhost_key); |
| 1006 | res!(db.start(&label)); |
| 1007 | res!(ok!(db.updated_api()).activate_gc(true)); |
| 1008 | |
| 1009 | std::thread::sleep(Duration::from_millis(200)); |
| 1010 | |
| 1011 | let (start, msgs) = res!(db.api().ping_bots(app_const::GET_DATA_WAIT)); |
| 1012 | info!("Vhost '{}': {} ping replies in {:?}.", |
| 1013 | spec.vhost_key, msgs.len(), start.elapsed()); |
| 1014 | |
| 1015 | // Publish as soon as each database is up, rather than |
| 1016 | // batching at the end: a vhost whose database is ready has |
| 1017 | // no reason to keep answering 503 while a later one starts. |
| 1018 | let mut guard = lock_write!(vhost_dbs, |
| 1019 | "Attaching the database for vhost '{}'.", spec.vhost_key); |
| 1020 | guard.insert( |
| 1021 | spec.vhost_key.clone(), |
| 1022 | (Arc::new(RwLock::new(db)), uid), |
| 1023 | ); |
| 1024 | opened += 1; |
| 1025 | } |
| 1026 | Ok(opened) |
| 1027 | }).await; |
| 1028 | |
| 1029 | match outcome { |
| 1030 | Ok(Ok(n)) => info!("Unsealed: {} database(s) open and attached.", n), |
| 1031 | Ok(Err(e)) => error!(e, "Opening the per-vhost databases after unseal."), |
| 1032 | Err(e) => error!(err!(e, |
| 1033 | "The database opener task failed to join."; |
| 1034 | Thread, Panic), |
| 1035 | "Opening the per-vhost databases after unseal."), |
| 1036 | } |
| 1037 | } |
| 1038 | |
| 1039 | |
| 1040 | |
| 1041 | /// The site's own newsletter sender: the DKIM identities and an outbound SMTP client, built from |
| 1042 | /// the `mail` block whether or not the mail server is enabled, and shared by every vhost. |
| 1043 | /// |
| 1044 | /// `None` where there is no sending identity -- a hostname to greet with and at least one DKIM |
| 1045 | /// key -- or where one cannot be built. Newsletter signup then answers "not set up" rather than |
| 1046 | /// recording a pending subscriber it could never confirm. A failure warns and disables the |
| 1047 | /// newsletter rather than refusing to start, because the websites do not depend on it: on |
| 1048 | /// 2026-09-21 a disabled `mail` block naming an RSA key that was never generated stopped birch, |
| 1049 | /// a forge proxy that sends no mail, from starting. The keys are still loaded when mail is |
| 1050 | /// disabled, since that is what lets a site send its newsletter signed without becoming an MX. |
| 1051 | fn newsletter_sender( |
| 1052 | mail: Option<crate::srv::cfg::MailConfig>, |
| 1053 | root: &oxedyne_fe2o3_core::path::NormPathBuf, |
| 1054 | ) |
| 1055 | -> Option<Arc<crate::srv::publish::send::MailSender>> |
| 1056 | { |
| 1057 | let mail_cfg = match mail { |
| 1058 | Some(m) if !m.hostname.is_empty() |
| 1059 | && (!m.dkim_key_file.is_empty() || !m.dkim_rsa_key_file.is_empty()) => m, |
| 1060 | _ => { |
| 1061 | info!("No mail configured, so the newsletter is unavailable; signup says so."); |
| 1062 | return None; |
| 1063 | }, |
| 1064 | }; |
| 1065 | let dkim = match crate::srv::server::load_dkim_signers(&mail_cfg, root) { |
| 1066 | Ok(d) => d, |
| 1067 | Err(e) => { |
| 1068 | warn!("The newsletter is disabled, and signup will say so: its DKIM key(s) could \ |
| 1069 | not be loaded ({}). The websites are unaffected.", e); |
| 1070 | return None; |
| 1071 | }, |
| 1072 | }; |
| 1073 | let domain = if mail_cfg.dkim_domain.is_empty() { |
| 1074 | mail_cfg.hostname.clone() |
| 1075 | } else { |
| 1076 | mail_cfg.dkim_domain.clone() |
| 1077 | }; |
| 1078 | let default_from = fmt!("news@{}", domain); |
| 1079 | match crate::srv::publish::send::MailSender::new( |
| 1080 | mail_cfg.hostname.clone(), dkim.clone(), default_from.clone()) |
| 1081 | { |
| 1082 | Ok(s) => { |
| 1083 | info!("Newsletter sender ready (default from {}, {} DKIM key(s)).", |
| 1084 | default_from, dkim.len()); |
| 1085 | Some(Arc::new(s)) |
| 1086 | }, |
| 1087 | Err(e) => { |
| 1088 | warn!("Building the newsletter sender failed ({}); newsletter \ |
| 1089 | signup will report itself unavailable.", e); |
| 1090 | None |
| 1091 | }, |
| 1092 | } |
| 1093 | } |
| 1094 | |
| 1095 | |
| 1096 | #[cfg(test)] |
| 1097 | mod tests { |
| 1098 | use super::*; |
| 1099 | use crate::srv::cfg::MailConfig; |
| 1100 | use oxedyne_fe2o3_core::path::NormalPath; |
| 1101 | |
| 1102 | /// A fresh scratch directory, removed by the caller. |
| 1103 | fn scratch(tag: &str) -> Outcome<std::path::PathBuf> { |
| 1104 | let nanos = std::time::SystemTime::now() |
| 1105 | .duration_since(std::time::UNIX_EPOCH) |
| 1106 | .map(|d| d.as_nanos()) |
| 1107 | .unwrap_or(0); |
| 1108 | let dir = std::env::temp_dir().join(fmt!( |
| 1109 | "fe2o3_steel_newsletter_{}_{}_{}", tag, std::process::id(), nanos)); |
| 1110 | res!(std::fs::create_dir_all(&dir), IO, File); |
| 1111 | Ok(dir) |
| 1112 | } |
| 1113 | |
| 1114 | /// THE BIRCH START-UP FAILURE: a disabled `mail` block naming an RSA key that was never |
| 1115 | /// generated disables the newsletter and lets the server start, where it used to refuse. |
| 1116 | #[test] |
| 1117 | fn a_missing_newsletter_key_disables_the_newsletter_and_does_not_stop_start_up() -> Outcome<()> { |
| 1118 | let dir = res!(scratch("missing")); |
| 1119 | let root = Path::new(&dir).normalise().absolute(); |
| 1120 | let mut cfg = MailConfig::default(); |
| 1121 | cfg.enabled = false; |
| 1122 | cfg.hostname = fmt!("birch.example.test"); |
| 1123 | cfg.dkim_rsa_key_file = fmt!("mail/dkim_rsa.key"); |
| 1124 | assert!(newsletter_sender(Some(cfg), &root).is_none(), |
| 1125 | "a key that is not there cannot sign, so there is no newsletter sender"); |
| 1126 | let _ = std::fs::remove_dir_all(&dir); |
| 1127 | Ok(()) |
| 1128 | } |
| 1129 | |
| 1130 | /// Mail disabled is not signing disabled: a site with a key sends its newsletter signed, |
| 1131 | /// which is why the key is not skipped wholesale when `enabled` is false. |
| 1132 | #[test] |
| 1133 | fn a_send_only_newsletter_is_still_signed_when_mail_is_disabled() -> Outcome<()> { |
| 1134 | let dir = res!(scratch("signed")); |
| 1135 | let root = Path::new(&dir).normalise().absolute(); |
| 1136 | let mut cfg = MailConfig::default(); |
| 1137 | cfg.enabled = false; |
| 1138 | cfg.hostname = fmt!("site.example.test"); |
| 1139 | cfg.dkim_key_file = fmt!("mail/dkim.key"); |
| 1140 | let sender = match newsletter_sender(Some(cfg), &root) { |
| 1141 | Some(s) => s, |
| 1142 | None => return Err(err!("A send-only site lost its newsletter sender."; Test)), |
| 1143 | }; |
| 1144 | let shown = fmt!("{:?}", sender); |
| 1145 | assert!(shown.contains("news@site.example.test") && shown.contains("dkim: 1"), |
| 1146 | "the sender must sign with its one key: {}", shown); |
| 1147 | let _ = std::fs::remove_dir_all(&dir); |
| 1148 | Ok(()) |
| 1149 | } |
| 1150 | |
| 1151 | #[test] |
| 1152 | fn no_mail_block_means_no_newsletter() { |
| 1153 | let root = Path::new("/nonexistent").normalise().absolute(); |
| 1154 | assert!(newsletter_sender(None, &root).is_none()); |
| 1155 | let mut bare = MailConfig::default(); |
| 1156 | bare.hostname = fmt!("site.example.test"); |
| 1157 | assert!(newsletter_sender(Some(bare), &root).is_none(), "no key, no sending identity"); |
| 1158 | } |
| 1159 | } |