oxedyne/fe2o3/fe2o3_steel/src/srv/context.rs
13.4 KiB, 157 runs
created by r1870400018:975, 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::srv::{ |
| 2 | admin::{ |
| 3 | state::AdminState, |
| 4 | traffic::TrafficRecorder, |
| 5 | }, |
| 6 | cfg::{ |
| 7 | ProxyRoute, |
| 8 | RedirectRule, |
| 9 | ServerConfig, |
| 10 | WsRoute, |
| 11 | }, |
| 12 | id, |
| 13 | tiles::{ |
| 14 | TileService, |
| 15 | TileSource, |
| 16 | }, |
| 17 | }; |
| 18 | |
| 19 | use oxedyne_fe2o3_core::{ |
| 20 | prelude::*, |
| 21 | id::ParseId, |
| 22 | path::NormPathBuf, |
| 23 | rand::Rand, |
| 24 | }; |
| 25 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 26 | use oxedyne_fe2o3_hash::{ |
| 27 | csum::ChecksumScheme, |
| 28 | hash::HashScheme, |
| 29 | }; |
| 30 | use oxedyne_fe2o3_iop_crypto::enc::Encrypter; |
| 31 | use oxedyne_fe2o3_iop_db::api::Database; |
| 32 | use oxedyne_fe2o3_iop_hash::api::Hasher; |
| 33 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 34 | use oxedyne_fe2o3_net::{ |
| 35 | http::{ |
| 36 | handler::WebHandler, |
| 37 | msg::HttpMessage, |
| 38 | }, |
| 39 | id::Sid, |
| 40 | ws::{ |
| 41 | WebSocket, |
| 42 | handler::WebSocketHandler, |
| 43 | }, |
| 44 | }; |
| 45 | use oxedyne_fe2o3_o3db_sync::{ |
| 46 | O3db, |
| 47 | base::cfg::OzoneConfig, |
| 48 | data::core::RestSchemesInput, |
| 49 | }; |
| 50 | use oxedyne_fe2o3_syntax::core::SyntaxRef; |
| 51 | |
| 52 | use std::{ |
| 53 | collections::{ |
| 54 | BTreeMap, |
| 55 | HashMap, |
| 56 | }, |
| 57 | marker::PhantomData, |
| 58 | net::SocketAddr, |
| 59 | path::Path, |
| 60 | sync::{ |
| 61 | Arc, |
| 62 | RwLock, |
| 63 | }, |
| 64 | }; |
| 65 | |
| 66 | use tokio::io::{ |
| 67 | AsyncRead, |
| 68 | AsyncWrite, |
| 69 | }; |
| 70 | |
| 71 | |
| 72 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 73 | // │ VHOST RUNTIME │ |
| 74 | // │ │ |
| 75 | // │ One per configured vhost. Carries everything the request path needs to │ |
| 76 | // │ serve that specific site: handlers, hostnames for validation, redirects. │ |
| 77 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 78 | |
| 79 | #[derive(Clone, Debug)] |
| 80 | pub struct VhostRuntime< |
| 81 | WH: WebHandler, |
| 82 | WSH: WebSocketHandler, |
| 83 | > { |
| 84 | pub hostnames: Vec<String>, |
| 85 | pub web_handler: WH, |
| 86 | pub ws_handler: WSH, |
| 87 | pub ws_syntax: SyntaxRef, |
| 88 | pub redirects: Vec<RedirectRule>, |
| 89 | pub proxy_routes: Vec<ProxyRoute>, |
| 90 | pub ws_routes: Vec<WsRoute>, |
| 91 | pub term_manager: Option<Arc<crate::srv::ws::term::TerminalManager>>, |
| 92 | pub uses_sessions: bool, |
| 93 | pub permissions_policy: Option<String>, // replaces the default Permissions-Policy header when set |
| 94 | pub tiles: Option<Arc<TileService<TileSource>>>, |
| 95 | pub access_log: bool, // `false` keeps this vhost's requests out of the log and the recorder |
| 96 | } |
| 97 | |
| 98 | impl< |
| 99 | WH: WebHandler, |
| 100 | WSH: WebSocketHandler, |
| 101 | > |
| 102 | VhostRuntime<WH, WSH> |
| 103 | { |
| 104 | pub fn primary_hostname(&self) -> &str { |
| 105 | self.hostnames.first().map(|s| s.as_str()).unwrap_or("") |
| 106 | } |
| 107 | |
| 108 | pub fn accepts_host(&self, host: &str) -> bool { |
| 109 | let host_lc = host.to_lowercase(); |
| 110 | // Strip any :port suffix. |
| 111 | let host_lc = match host_lc.find(':') { |
| 112 | Some(i) => host_lc[..i].to_string(), |
| 113 | None => host_lc, |
| 114 | }; |
| 115 | self.hostnames.iter().any(|h| h.to_lowercase() == host_lc) |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | |
| 120 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 121 | // │ PROTOCOL │ |
| 122 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 123 | |
| 124 | #[derive(Clone, Debug)] |
| 125 | pub enum Protocol< |
| 126 | WH: WebHandler, |
| 127 | WSH: WebSocketHandler, |
| 128 | > { |
| 129 | Web { |
| 130 | vhosts: Arc<HashMap<String, Arc<VhostRuntime<WH, WSH>>>>, |
| 131 | default_vhost: String, |
| 132 | dev_mode: bool, |
| 133 | }, |
| 134 | } |
| 135 | |
| 136 | |
| 137 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 138 | // │ SERVER CONTEXT │ |
| 139 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 140 | |
| 141 | pub type VhostDbs<const UIDL: usize, UID, DB> = |
| 142 | Arc<RwLock<HashMap<String, (Arc<RwLock<DB>>, UID)>>>; |
| 143 | |
| 144 | #[derive(Clone, Debug)] |
| 145 | pub struct VhostDbSpec { |
| 146 | pub vhost_key: String, |
| 147 | pub db_dir: std::path::PathBuf, |
| 148 | } |
| 149 | |
| 150 | pub struct ServerContext< |
| 151 | const UIDL: usize, |
| 152 | UID: NumIdDat<UIDL> + 'static, |
| 153 | ENC: Encrypter, // Symmetric encryption of database. |
| 154 | KH: Hasher, // Hashes database keys. |
| 155 | DB: Database<UIDL, UID, ENC, KH>, |
| 156 | WH: WebHandler, |
| 157 | WSH: WebSocketHandler, |
| 158 | > { |
| 159 | pub cfg: ServerConfig, |
| 160 | pub root: NormPathBuf, |
| 161 | pub vhost_dbs: VhostDbs<UIDL, UID, DB>, |
| 162 | pub protocol: Protocol<WH, WSH>, |
| 163 | pub traffic: Option<Arc<TrafficRecorder>>, |
| 164 | pub admin_state: Option<Arc<AdminState>>, |
| 165 | pub db_specs: Vec<VhostDbSpec>, |
| 166 | phantom3: PhantomData<ENC>, |
| 167 | phantom4: PhantomData<KH>, |
| 168 | } |
| 169 | |
| 170 | impl< |
| 171 | const UIDL: usize, |
| 172 | UID: NumIdDat<UIDL> + 'static, |
| 173 | ENC: Encrypter + 'static, |
| 174 | KH: Hasher + 'static, |
| 175 | DB: Database<UIDL, UID, ENC, KH> + 'static, |
| 176 | WH: WebHandler + 'static, |
| 177 | WSH: WebSocketHandler + 'static, |
| 178 | > |
| 179 | Clone for ServerContext<UIDL, UID, ENC, KH, DB, WH, WSH> |
| 180 | { |
| 181 | fn clone(&self) -> Self { |
| 182 | Self { |
| 183 | cfg: self.cfg.clone(), |
| 184 | root: self.root.clone(), |
| 185 | vhost_dbs: self.vhost_dbs.clone(), |
| 186 | protocol: self.protocol.clone(), |
| 187 | traffic: self.traffic.clone(), |
| 188 | admin_state: self.admin_state.clone(), |
| 189 | db_specs: self.db_specs.clone(), |
| 190 | phantom3: PhantomData, |
| 191 | phantom4: PhantomData, |
| 192 | } |
| 193 | } |
| 194 | } |
| 195 | |
| 196 | impl< |
| 197 | const UIDL: usize, |
| 198 | UID: NumIdDat<UIDL> + 'static, |
| 199 | ENC: Encrypter + 'static, |
| 200 | KH: Hasher + 'static, |
| 201 | DB: Database<UIDL, UID, ENC, KH> + 'static, |
| 202 | WH: WebHandler + 'static, |
| 203 | WSH: WebSocketHandler + 'static, |
| 204 | > |
| 205 | ServerContext<UIDL, UID, ENC, KH, DB, WH, WSH> |
| 206 | { |
| 207 | pub fn new( |
| 208 | cfg: ServerConfig, |
| 209 | root: NormPathBuf, |
| 210 | vhost_dbs: VhostDbs<UIDL, UID, DB>, |
| 211 | db_specs: Vec<VhostDbSpec>, |
| 212 | protocol: Protocol<WH, WSH>, |
| 213 | traffic: Option<Arc<TrafficRecorder>>, |
| 214 | admin_state: Option<Arc<AdminState>>, |
| 215 | ) |
| 216 | -> Self |
| 217 | { |
| 218 | Self { |
| 219 | cfg, |
| 220 | root, |
| 221 | vhost_dbs, |
| 222 | db_specs, |
| 223 | protocol, |
| 224 | traffic, |
| 225 | admin_state, |
| 226 | phantom3: PhantomData, |
| 227 | phantom4: PhantomData, |
| 228 | } |
| 229 | } |
| 230 | |
| 231 | pub fn db_for_vhost( |
| 232 | &self, |
| 233 | vhost_key: &str, |
| 234 | ) |
| 235 | -> Option<(Arc<RwLock<DB>>, UID)> |
| 236 | { |
| 237 | let guard = match self.vhost_dbs.read() { |
| 238 | Ok(g) => g, |
| 239 | Err(_) => { |
| 240 | // This returns `Option`, not `Outcome`, so a poisoned lock |
| 241 | // cannot be propagated. Report it and answer as though the |
| 242 | // vhost has no database: the caller degrades to a 404 or a |
| 243 | // 503 rather than serving from a map nobody can read. |
| 244 | fault!("The vhost database map lock is poisoned; treating \ |
| 245 | '{}' as having no database.", vhost_key); |
| 246 | return None; |
| 247 | } |
| 248 | }; |
| 249 | guard.get(&vhost_key.to_lowercase()).cloned() |
| 250 | } |
| 251 | |
| 252 | pub fn is_sealed(&self) -> bool { |
| 253 | match &self.admin_state { |
| 254 | Some(state) => state.is_sealed(), |
| 255 | None => false, |
| 256 | } |
| 257 | } |
| 258 | |
| 259 | pub fn db_pending_for_vhost(&self, vhost_key: &str) -> bool { |
| 260 | if !self.is_sealed() { |
| 261 | return false; |
| 262 | } |
| 263 | let key = vhost_key.to_lowercase(); |
| 264 | self.db_specs.iter().any(|spec| spec.vhost_key == key) |
| 265 | } |
| 266 | |
| 267 | pub fn clone_self(&self) -> Self { |
| 268 | self.clone() |
| 269 | } |
| 270 | |
| 271 | pub fn vhost_for(&self, sni: Option<&str>) -> Arc<VhostRuntime<WH, WSH>> { |
| 272 | match &self.protocol { |
| 273 | Protocol::Web { vhosts, default_vhost, .. } => { |
| 274 | if let Some(name) = sni { |
| 275 | if let Some(vh) = vhosts.get(&name.to_lowercase()) { |
| 276 | return vh.clone(); |
| 277 | } |
| 278 | } |
| 279 | // Fall through to default. |
| 280 | match vhosts.get(&default_vhost.to_lowercase()) { |
| 281 | Some(vh) => vh.clone(), |
| 282 | None => { |
| 283 | // Should not happen if startup validated properly. |
| 284 | // Return the first entry if any; otherwise panic is |
| 285 | // impossible here because start-up would have failed. |
| 286 | vhosts.values().next().cloned().expect( |
| 287 | "ServerContext::vhost_for: no vhosts configured \ |
| 288 | -- this should have been rejected at start-up.", |
| 289 | ) |
| 290 | } |
| 291 | } |
| 292 | } |
| 293 | } |
| 294 | } |
| 295 | |
| 296 | pub fn err_id() -> String { |
| 297 | Rand::generate_random_string(6, "abcdefghikmnpqrstuvw0123456789") |
| 298 | } |
| 299 | |
| 300 | pub fn get_session_id( |
| 301 | msg: &HttpMessage, |
| 302 | src_addr: &SocketAddr, |
| 303 | ) |
| 304 | -> Option<Sid> |
| 305 | { |
| 306 | match msg.header.fields.get_session_id() { |
| 307 | Some(sid_string) => match Sid::parse_id(&sid_string) { |
| 308 | Ok(n) => Some(n), |
| 309 | Err(e) => { |
| 310 | error!(e, "The session cookie string '{}' in a message from \ |
| 311 | {:?} cannot be decoded to a {}.", |
| 312 | sid_string, src_addr, std::any::type_name::<Sid>()); |
| 313 | None |
| 314 | }, |
| 315 | }, |
| 316 | None => None, |
| 317 | } |
| 318 | } |
| 319 | } |
| 320 | |
| 321 | pub fn new_db( |
| 322 | db_root: &Path, |
| 323 | enc_key: &[u8], |
| 324 | ) |
| 325 | -> Outcome<O3db< |
| 326 | { id::UID_LEN }, |
| 327 | id::Uid, |
| 328 | EncryptionScheme, |
| 329 | HashScheme, |
| 330 | HashScheme, |
| 331 | ChecksumScheme, |
| 332 | >> |
| 333 | { |
| 334 | // Start from the library's own production configuration and state only what a Steel server |
| 335 | // deliberately wants different. This was previously a full struct literal, and it carried |
| 336 | // the values from `o3db_sync`'s *test* setup -- a 1.5 KB chunking threshold and a 64-byte |
| 337 | // chunk size -- so every Steel store split any value past 1.5 KB into 64-byte pieces. An |
| 338 | // exhaustive literal is what let that happen and what kept it: it has to restate every field, |
| 339 | // so a wrong one is invisible among the right ones, and a field added upstream cannot reach a |
| 340 | // caller who has already spelled them all out. Naming only the deviations makes each one a |
| 341 | // decision, and everything else tracks the library. |
| 342 | let cfg = OzoneConfig { |
| 343 | // A Steel server may hold several stores, so it takes a tenth of the library's cache. |
| 344 | cache_size_limit_bytes: 100_000_000, |
| 345 | // One writer per zone: the server's stores are small and write-light. |
| 346 | num_wbots_per_zone: 1, |
| 347 | // Zone state is reported more often than the library default, so the dashboard's view of |
| 348 | // a store is near-live rather than five seconds stale. |
| 349 | zone_state_update_secs: 1, |
| 350 | // No per-zone size caps: a Steel store grows with the application that owns it. |
| 351 | zone_overrides: BTreeMap::new(), |
| 352 | // Everything else -- the chunking above all -- is the library's own production default. |
| 353 | ..Default::default() |
| 354 | }; |
| 355 | |
| 356 | let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(enc_key)); |
| 357 | let crc32 = ChecksumScheme::new_crc32(); |
| 358 | let schms_input = RestSchemesInput::new( |
| 359 | Some(aes_gcm.clone()), |
| 360 | None::<HashScheme>, |
| 361 | None::<HashScheme>, |
| 362 | Some(crc32.clone()), |
| 363 | ); |
| 364 | O3db::new( |
| 365 | &db_root, |
| 366 | Some(cfg), |
| 367 | schms_input, |
| 368 | id::Uid::default(), |
| 369 | ) |
| 370 | } |
| 371 | |
| 372 | pub fn no_db() |
| 373 | -> Outcome<HashMap<String, (Arc<RwLock<O3db< |
| 374 | { id::UID_LEN }, |
| 375 | id::Uid, |
| 376 | EncryptionScheme, |
| 377 | HashScheme, |
| 378 | HashScheme, |
| 379 | ChecksumScheme, |
| 380 | >>>, |
| 381 | id::Uid, |
| 382 | )>> |
| 383 | { |
| 384 | Ok(HashMap::new()) |
| 385 | } |
| 386 | |
| 387 | pub fn new_ws_no_db< |
| 388 | 'a, |
| 389 | S: AsyncRead + AsyncWrite + Unpin, |
| 390 | WSH: WebSocketHandler, |
| 391 | >( |
| 392 | stream: &'a mut S, |
| 393 | ws_handler: WSH, |
| 394 | ) |
| 395 | -> Outcome<WebSocket< |
| 396 | 'a, |
| 397 | { id::UID_LEN }, |
| 398 | id::Uid, |
| 399 | EncryptionScheme, |
| 400 | HashScheme, |
| 401 | O3db< |
| 402 | { id::UID_LEN }, |
| 403 | id::Uid, |
| 404 | EncryptionScheme, |
| 405 | HashScheme, |
| 406 | HashScheme, |
| 407 | ChecksumScheme, |
| 408 | >, |
| 409 | S, |
| 410 | WSH, |
| 411 | >> |
| 412 | { |
| 413 | Ok(WebSocket::< |
| 414 | '_, |
| 415 | { id::UID_LEN }, |
| 416 | id::Uid, |
| 417 | EncryptionScheme, |
| 418 | HashScheme, |
| 419 | O3db< |
| 420 | { id::UID_LEN }, |
| 421 | id::Uid, |
| 422 | EncryptionScheme, |
| 423 | HashScheme, |
| 424 | HashScheme, |
| 425 | ChecksumScheme, |
| 426 | >, |
| 427 | S, |
| 428 | WSH, |
| 429 | >::new_client( |
| 430 | stream, |
| 431 | ws_handler, |
| 432 | 10, |
| 433 | 20, |
| 434 | )) |
| 435 | } |