Oregami
Repositories/oxedyne/fe2o3

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

1use 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
19use oxedyne_fe2o3_core::{
20 prelude::*,
21 id::ParseId,
22 path::NormPathBuf,
23 rand::Rand,
24};
25use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
26use oxedyne_fe2o3_hash::{
27 csum::ChecksumScheme,
28 hash::HashScheme,
29};
30use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
31use oxedyne_fe2o3_iop_db::api::Database;
32use oxedyne_fe2o3_iop_hash::api::Hasher;
33use oxedyne_fe2o3_jdat::id::NumIdDat;
34use oxedyne_fe2o3_net::{
35 http::{
36 handler::WebHandler,
37 msg::HttpMessage,
38 },
39 id::Sid,
40 ws::{
41 WebSocket,
42 handler::WebSocketHandler,
43 },
44};
45use oxedyne_fe2o3_o3db_sync::{
46 O3db,
47 base::cfg::OzoneConfig,
48 data::core::RestSchemesInput,
49};
50use oxedyne_fe2o3_syntax::core::SyntaxRef;
51
52use 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
66use 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)]
80pub 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
98impl<
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)]
125pub 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
141pub type VhostDbs<const UIDL: usize, UID, DB> =
142 Arc<RwLock<HashMap<String, (Arc<RwLock<DB>>, UID)>>>;
143
144#[derive(Clone, Debug)]
145pub struct VhostDbSpec {
146 pub vhost_key: String,
147 pub db_dir: std::path::PathBuf,
148}
149
150pub 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
170impl<
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
196impl<
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
321pub 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
372pub 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
387pub 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}