Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_steel/src/srv/https.rs

67.4 KiB, 245 runs

created by r1870400018:987, 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::traffic::{
3 self,
4 RequestRecord,
5 },
6 cfg::{
7 ProxyRoute,
8 RedirectMatch,
9 RedirectRule,
10 WsRoute,
11 },
12 constant,
13 context::ServerContext,
14 tiles::TileRequest,
15 wsproxy,
16};
17
18use oxedyne_fe2o3_core::{
19 prelude::*,
20 error::ErrTag,
21 id::ParseId,
22 rand::RanDef,
23};
24use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
25use oxedyne_fe2o3_iop_db::api::Database;
26use oxedyne_fe2o3_iop_hash::api::Hasher;
27use oxedyne_fe2o3_jdat::{
28 prelude::*,
29 id::{
30 IdDat,
31 NumIdDat,
32 },
33};
34use oxedyne_fe2o3_net::{
35 conc::AsyncReadIterator,
36 http::{
37 encoding,
38 fields::{
39 Cookie,
40 HeaderFieldValue,
41 HeaderName,
42 },
43 fwd::{
44 self,
45 ForwardedPolicy,
46 },
47 handler::WebHandler,
48 header::{
49 HttpHeadline,
50 HttpMethod,
51 },
52 msg::{
53 HttpMessageReader,
54 HttpMessage,
55 ReadLimits,
56 },
57 status::HttpStatus,
58 },
59 id::Sid,
60 ws::handler::WebSocketHandler,
61};
62
63use std::{
64 net::SocketAddr,
65 pin::Pin,
66 sync::Arc,
67 time::{
68 Duration,
69 Instant,
70 },
71};
72
73use tokio::{
74 net::TcpStream,
75 io::{
76 AsyncRead,
77 AsyncReadExt,
78 AsyncWrite,
79 AsyncWriteExt,
80 },
81};
82use tokio_rustls::server::TlsStream;
83
84
85// A log line about one connection or request, written only where the vhost
86// keeps an access log.
87macro_rules! alog {
88 ($on:expr, $($arg:tt)+) => {
89 if $on {
90 log!($($arg)+);
91 }
92 };
93}
94
95impl<
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 ServerContext<UIDL, UID, ENC, KH, DB, WH, WSH>
105{
106 /// The token-gated health body, or `None` when the request is not for the
107 /// health path (so the caller falls through to normal dispatch).
108 ///
109 /// When the path *is* the health path this always answers -- a `200` with the
110 /// integer body to a caller presenting the right token, a `404` to anyone
111 /// else -- and never falls through, so the path never reaches a file router
112 /// that might serve something at the same location. A missing token and a
113 /// wrong token take the identical `404` path, compared in constant time, so
114 /// the route leaks no timing oracle and its existence stays invisible. The
115 /// body is served even while the box is sealed (the watcher reads a `503` as
116 /// down), and the caller answers this before it logs the request, so the
117 /// token in the header is never written to a log.
118 fn health_response(&self, request: &HttpMessage) -> Option<HttpMessage> {
119 // Both a path and a token must be configured; a path without a token is
120 // a body with no gate, which the config contract disables.
121 if self.cfg.health_path.is_empty() || self.cfg.health_token.is_empty() {
122 return None;
123 }
124 let path = match &request.header.headline {
125 HttpHeadline::Request { loc, .. } => loc.path.as_string().to_string(),
126 _ => return None,
127 };
128 if path != self.cfg.health_path {
129 return None;
130 }
131 // A 404 that is byte-for-byte the unauthorised answer, reused for both
132 // the no-admin-state and wrong-token cases so neither is distinguishable.
133 let not_found = || HttpMessage::respond_with_text(
134 HttpStatus::NotFound, "Not Found");
135 let admin = match self.admin_state.as_ref() {
136 Some(a) => a,
137 None => return Some(not_found()),
138 };
139 let presented = match request.header.fields.get_one(
140 &HeaderName::NonStandard("x-steel-health-token".to_string()))
141 {
142 Some(v) => fmt!("{}", v),
143 None => String::new(),
144 };
145 // Constant-time, so a missing and a wrong token are one path.
146 if !oxedyne_fe2o3_core::byte::ct_eq(
147 presented.as_bytes(), self.cfg.health_token.as_bytes())
148 {
149 return Some(not_found());
150 }
151 // Assembled per request, stamp ages included: each stamp is read now rather than
152 // sampled, so a job whose timer has died shows an age that keeps growing.
153 let body = admin.health_body();
154 Some(HttpMessage::new_response(HttpStatus::OK)
155 .with_field(
156 HeaderName::ContentType,
157 HeaderFieldValue::Generic("application/json; charset=utf-8".to_string()),
158 )
159 .with_field(
160 HeaderName::CacheControl,
161 HeaderFieldValue::Generic("no-store".to_string()),
162 )
163 .with_body(body.to_json().into_bytes()))
164 }
165
166 pub async fn handle_https(
167 self,
168 mut stream: TlsStream<TcpStream>,
169 sni: Option<String>,
170 src_addr: SocketAddr,
171 )
172 -> Outcome<()>
173 {
174 let id = fmt!("Https|Cx:{}", IdDat::<4, u32>::randef()); // Cx = Connection id.
175
176 // Resolve the vhost once per connection from the SNI. All requests
177 // on a single TLS connection are considered to target the same vhost,
178 // which matches how every HTTP/1.1 and HTTP/2 client behaves.
179 let vhost = self.vhost_for(sni.as_deref());
180 let log_level = res!(self.cfg.log_level());
181 // A vhost with its access log off writes no line naming a peer or a
182 // request, and gives the traffic recorder nothing.
183 let logged = vhost.access_log;
184 alog!(logged, log_level, "{}: connection from {:?}, sni={:?}, vhost='{}'.",
185 id, src_addr, sni, vhost.primary_hostname());
186
187 let (mut read_stream, mut write_stream) = tokio::io::split(&mut stream);
188
189 // Build per-connection read limits from ServerConfig so the
190 // reader enforces the configured header / body bounds and the
191 // slowloris read deadline. A zero value in the config means
192 // "disabled" and maps to `None` in `ReadLimits`.
193 let limits = ReadLimits {
194 max_header_bytes: if self.cfg.http_max_header_bytes == 0 {
195 None
196 } else {
197 Some(self.cfg.http_max_header_bytes as usize)
198 },
199 max_body_bytes: if self.cfg.http_max_body_bytes == 0 {
200 None
201 } else {
202 Some(self.cfg.http_max_body_bytes as usize)
203 },
204 header_read_timeout: if self.cfg.http_header_read_timeout_ms == 0 {
205 None
206 } else {
207 Some(Duration::from_millis(
208 self.cfg.http_header_read_timeout_ms,
209 ))
210 },
211 };
212
213 let mut reader: HttpMessageReader<
214 '_,
215 { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE },
216 { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE },
217 _,
218 > = HttpMessageReader::with_limits(Pin::new(&mut read_stream), limits);
219
220 loop {
221 let result = reader.next().await;
222 // Capture request start time + method/path before the
223 // request is consumed by the dispatch chain. Used at
224 // the bottom of the loop to emit a TrafficRecord that
225 // covers the full handle-to-write duration.
226 let req_started_at = Instant::now();
227 match result {
228 Some(Ok(request)) => {
229 // Health route, answered before anything is logged so the
230 // token in the header is never written to a log, and before
231 // vhost/Host dispatch so it needs no vhost of its own. When
232 // it answers, the connection is done.
233 if let Some(mut resp) = self.health_response(&request) {
234 resp.set_connection_close(true);
235 if let Err(e) = resp.write_all(&mut write_stream).await {
236 warn!("{}: could not send health response: {}", id, e);
237 }
238 break;
239 }
240 // Validate the Host header against the vhost hostnames.
241 // A mismatch means an SNI/Host disagreement, which is a
242 // misdirected client; we return 421.
243 if let Some(HeaderFieldValue::Generic(host_hdr)) =
244 request.header.fields.get_one(&HeaderName::Host)
245 {
246 if !vhost.accepts_host(host_hdr) {
247 warn!("{}: Host header '{}' does not match vhost '{}' \
248 (hostnames={:?}); returning 421 Misdirected Request.",
249 id, host_hdr, vhost.primary_hostname(), vhost.hostnames);
250 let mut resp = HttpMessage::respond_with_text(
251 HttpStatus::MisdirectedRequest,
252 "Misdirected request: Host header does not match SNI.",
253 );
254 resp.set_connection_close(true);
255 match resp.write_all(&mut write_stream).await {
256 Ok(()) => (),
257 Err(e) => return Err(err!(e,
258 "{}: Could not send 421 response.", id;
259 IO, Network, Wire, Write)),
260 }
261 break;
262 }
263 }
264
265 // Tiles, answered before the request is logged, before the
266 // session cookie is read and before the traffic recorder
267 // sees it, since a tile request says where its viewer
268 // looked. See `srv::tiles`. The connection stays open for
269 // the next tile unless the client asked to close it.
270 if let Some(tiles) = vhost.tiles.as_ref() {
271 if let Some(treq) = TileRequest::of(&request) {
272 if let Some(mut resp) = tiles.respond(treq).await {
273 let close = request.get_connection_close();
274 if close {
275 resp.set_connection_close(true);
276 }
277 match resp.write_all(&mut write_stream).await {
278 Ok(()) => (),
279 Err(e) => return Err(err!(e,
280 "{}: Could not send a tile response.", id;
281 IO, Network, Wire, Write)),
282 }
283 if close {
284 break;
285 }
286 continue;
287 }
288 }
289 }
290
291 alog!(logged, log_level, "{}: Incoming from {:?}:", id, src_addr);
292 if logged {
293 request.log(log_get_level!());
294 }
295
296 // Pull method+path out for traffic recording
297 // before the request is moved into the dispatch
298 // chain. Both are cheap clones.
299 let (rec_method, rec_path) = match &request.header.headline {
300 HttpHeadline::Request { method, loc } => (
301 fmt!("{}", method),
302 loc.path.as_string().to_string(),
303 ),
304 _ => (String::new(), String::new()),
305 };
306
307 // What content codings the client will take, kept as a
308 // string because the request itself is moved into the
309 // dispatch chain long before the response is encoded.
310 let accept_encoding = encoding::accept_encoding(
311 &request.header.fields,
312 );
313
314 // Resolve (or issue) the session identifier for this
315 // request. If the client already carries a session
316 // cookie, parse it. Otherwise, when anonymous sessions
317 // are enabled *and this vhost has somewhere to keep one*,
318 // mint a fresh `Sid`, remember it as a pending
319 // `Set-Cookie` header to attach to the response, and use
320 // it to scope session commands on this request.
321 //
322 // A vhost with no database is excluded because a session
323 // identifier is a key prefix into that database and
324 // nothing else. Issuing one there set a cookie on every
325 // response a static site made -- every stylesheet, every
326 // image, every module -- which no shared cache will store,
327 // and which the operator would then have to account for to
328 // anyone asking what it was for. It was for nothing.
329 let raw_sid_str = request.header.fields.get_session_id();
330 let mut issued_cookie: Option<Cookie> = None;
331 let (sid_opt, sid_str) = match raw_sid_str {
332 Some(ref s) => {
333 let parsed = Sid::parse_id(s).ok();
334 (parsed, raw_sid_str.clone())
335 }
336 None => {
337 if self.cfg.allow_anonymous_sessions && vhost.uses_sessions {
338 let new_sid: Sid = Sid::randef();
339 let s = fmt!("{}", new_sid);
340 issued_cookie = Some(
341 self.cfg.session_cookie_default(s.clone()),
342 );
343 alog!(logged, log_level,
344 "{}: issuing anonymous session {}.", id, s);
345 (Some(new_sid), Some(s))
346 } else {
347 (None, None)
348 }
349 }
350 };
351
352 if request.is_websocket_upgrade() {
353 // ── WebSocket routes ───────────────────────
354 // Checked first: a route naming one exact path
355 // is more specific than a proxy prefix that
356 // happens to contain it, and more specific
357 // than Steel's own WS handler, which answers
358 // every other path.
359 if !vhost.ws_routes.is_empty() {
360 if let HttpHeadline::Request { ref loc, .. } =
361 request.header.headline
362 {
363 let ws_path = loc.path.as_string().to_string();
364 if let Some(ws_route) = vhost.ws_routes.iter()
365 .find(|r| r.matches(&ws_path))
366 {
367 alog!(logged, log_level,
368 "{}: ws route {} -> ws://{}:{}{}",
369 id, ws_path,
370 ws_route.upstream_host,
371 ws_route.upstream_port,
372 ws_route.upstream_path,
373 );
374 let reunited = read_stream.unsplit(write_stream);
375 return self.handle_ws_route(
376 reunited,
377 request,
378 ws_route,
379 src_addr,
380 &id,
381 ).await;
382 }
383 }
384 }
385 // Check proxy routes next — if a proxy route
386 // matches, tunnel the WebSocket to the upstream
387 // instead of handling it with Steel's own WS
388 // handler. This allows proxied applications
389 // that use WebSocket (e.g. web terminals) to
390 // work through the reverse proxy.
391 if !vhost.proxy_routes.is_empty() {
392 if let HttpHeadline::Request { ref loc, .. } =
393 request.header.headline
394 {
395 let proxy_path = loc.path.as_string().to_string();
396 if let Some(proxy_route) = vhost.proxy_routes.iter()
397 .filter(|r| proxy_path.starts_with(&r.path_prefix))
398 .max_by_key(|r| r.path_prefix.len())
399 {
400 alog!(logged, log_level,
401 "{}: proxy ws {} -> {}:{}{}",
402 id, proxy_path,
403 proxy_route.upstream_host,
404 proxy_route.upstream_port,
405 if proxy_route.upstream_tls { " (tls)" } else { "" },
406 );
407 let reunited = read_stream.unsplit(write_stream);
408 return self.handle_proxy_websocket(
409 reunited,
410 request,
411 proxy_route,
412 src_addr,
413 &id,
414 ).await;
415 }
416 }
417 }
418 // ── Red chat WS dispatch ───────────────────
419 // If the request path is /chat and the vhost
420 // ── Terminal WS dispatch ──────────────────
421 // If the request path starts with /term/,
422 // route to the terminal I/O bridge instead
423 // of the normal text-protocol WS handler.
424 // This allows binary terminal data to flow
425 // over a separate WS channel.
426 if let HttpHeadline::Request { ref loc, .. } =
427 request.header.headline
428 {
429 let req_path = loc.path.as_string();
430 if req_path.starts_with("/term/") {
431 let session_name = req_path
432 .strip_prefix("/term/")
433 .unwrap_or("")
434 .to_string();
435 if session_name.is_empty() {
436 alog!(logged, log_level,
437 "{}: terminal WS missing session name.", id);
438 let mut resp = HttpMessage::respond_with_text(
439 HttpStatus::BadRequest,
440 "Missing terminal session name.",
441 );
442 resp.set_connection_close(true);
443 let _ = resp.write_all(&mut write_stream).await;
444 break;
445 }
446 // ── Gate: fail closed on config and auth ──
447 // The bridge attaches a raw pty to a shell, so the upgrade is
448 // dispatched only when two things both hold, and each fails
449 // closed. The vhost must have enabled terminals -- a vhost with
450 // no `term_config` leaves `term_manager` unset -- and the
451 // request must carry an authenticated operator session, the
452 // same principal the admin routes read from the dashboard
453 // cookie. Absent either, the answer is a refusal, never a
454 // shell. The decision lives in `terminal_gate` so it is
455 // testable apart from this connection machinery.
456 let principal = self.admin_state.as_ref()
457 .and_then(|st| crate::srv::admin::handler::extract_principal(
458 st, &request.header.fields));
459 match terminal_gate(vhost.term_manager.is_some(), principal.as_ref()) {
460 TerminalGate::NotEnabled => {
461 warn!("{}: terminal WS refused for '{}': vhost '{}' has no term_config.",
462 id, session_name, vhost.primary_hostname());
463 let mut resp = HttpMessage::respond_with_text(
464 HttpStatus::Forbidden,
465 "Terminal features are not enabled on this host.",
466 );
467 resp.set_connection_close(true);
468 let _ = resp.write_all(&mut write_stream).await;
469 break;
470 }
471 TerminalGate::NotAuthenticated => {
472 warn!("{}: terminal WS refused for '{}': no authenticated operator session.",
473 id, session_name);
474 let mut resp = HttpMessage::respond_with_text(
475 HttpStatus::Forbidden,
476 "Terminal access requires an authenticated operator session.",
477 );
478 resp.set_connection_close(true);
479 let _ = resp.write_all(&mut write_stream).await;
480 break;
481 }
482 TerminalGate::Dispatch => (),
483 }
484 alog!(logged, log_level,
485 "{}: terminal ws -> '{}'", id, session_name);
486 let reunited = read_stream.unsplit(write_stream);
487 return crate::srv::ws::term::handle_terminal_websocket::<
488 UIDL, UID, ENC, KH, DB, _,
489 >(
490 reunited,
491 session_name,
492 request,
493 &id,
494 ).await;
495 }
496 }
497 alog!(logged, log_level, "Connection upgrading to websocket...");
498 // The raw sid string is enough for the WS handler:
499 // it only needs a stable per-client key prefix, not
500 // the typed numeric identifier.
501 // Resolve the operator principal for this request, if
502 // any, so the handler's `term_*` management commands --
503 // which spawn, list, rename and kill terminal sessions --
504 // are gated on an authenticated operator exactly as the
505 // `/term/` bridge above is. Absent an operator session
506 // the flag is false and those commands refuse.
507 let ws_operator = self.admin_state.as_ref()
508 .and_then(|st| crate::srv::admin::handler
509 ::extract_principal(st, &request.header.fields))
510 .is_some();
511 let ws_handler = vhost.ws_handler.clone()
512 .attach_sid(sid_str.clone())
513 .with_operator_authed(ws_operator);
514 let reunited_stream = read_stream.unsplit(write_stream);
515 let vhost_db = self.db_for_vhost(vhost.primary_hostname());
516 return self.handle_websocket(
517 reunited_stream,
518 ws_handler,
519 vhost.ws_syntax.clone(),
520 vhost_db,
521 request,
522 &id,
523 ).await;
524 }
525
526 let mut response = None;
527 let close_requested = request.get_connection_close();
528 if close_requested {
529 let mut msg = HttpMessage::new_response(HttpStatus::OK);
530 msg.set_connection_close(true);
531 response = Some(msg);
532 }
533
534 // Per-route rate limit: sensitive URL prefixes
535 // (login forms, admin login) go through a
536 // dedicated, tighter guard so a brute-force
537 // password hammer gets kicked off the login
538 // path faster than a normal browsing session.
539 if let HttpHeadline::Request { ref loc, .. } =
540 request.header.headline
541 {
542 let path_str = loc.path.as_string();
543 if self.cfg.auth_path_prefixes.iter()
544 .any(|p| path_str.starts_with(p.as_str()))
545 {
546 if let Some(admin) = self.admin_state.as_ref() {
547 match admin.auth_guard.check(&src_addr.ip()) {
548 Ok(d) if d.should_drop() => {
549 warn!("{}: auth guard dropping {} from {}: {:?}",
550 id, path_str, src_addr, d);
551 admin.r429.incr();
552 let mut resp = HttpMessage::respond_with_text(
553 HttpStatus::TooManyRequests,
554 "Too many authentication attempts. \
555 Please wait and try again.",
556 );
557 resp.set_connection_close(true);
558 match resp.write_all(&mut write_stream).await {
559 Ok(()) => (),
560 Err(we) => warn!(
561 "{}: failed to emit 429: {}",
562 id, we),
563 }
564 break;
565 }
566 Ok(_) => (),
567 Err(e) => warn!(
568 "{}: auth guard error for {}: {}",
569 id, src_addr, e),
570 }
571 }
572 }
573 }
574
575 match request.header.headline.clone() {
576 HttpHeadline::Request { method, loc } => {
577 // Redirect rules fire before the file router.
578 let request_uri = loc.path.as_string().to_string();
579 if let Some(rule) = Self::match_redirect(
580 &vhost.redirects,
581 &request_uri,
582 ) {
583 let target = rule.resolve_target(&request_uri);
584 alog!(logged, log_level,
585 "{}: redirect {} {} -> {} ({})",
586 id, rule.status, request_uri, target,
587 match rule.match_kind {
588 RedirectMatch::Exact => "exact",
589 RedirectMatch::Prefix => "prefix",
590 RedirectMatch::All => "all",
591 });
592 let status = match rule.status {
593 301 => HttpStatus::MovedPermanently,
594 302 => HttpStatus::Found,
595 303 => HttpStatus::SeeOther,
596 307 => HttpStatus::TemporaryRedirect,
597 308 => HttpStatus::PermanentRedirect,
598 _ => HttpStatus::MovedPermanently,
599 };
600 let resp = HttpMessage::new_response(status)
601 .with_field(
602 HeaderName::Location,
603 HeaderFieldValue::Generic(target),
604 );
605 response = Some(resp);
606 } else {
607 // ── Reverse proxy routes ──────────────
608 // Checked after redirects, before static
609 // files and API routes. Longest matching
610 // prefix wins. WebSocket upgrades are
611 // tunnelled; regular HTTP is streamed.
612 if !vhost.proxy_routes.is_empty() {
613 let proxy_path = loc.path.as_string().to_string();
614 if let Some(proxy_route) = vhost.proxy_routes.iter()
615 .filter(|r| proxy_path.starts_with(&r.path_prefix))
616 .max_by_key(|r| r.path_prefix.len())
617 {
618 alog!(logged, log_level,
619 "{}: proxy {} -> {}:{}{}",
620 id, proxy_path,
621 proxy_route.upstream_host,
622 proxy_route.upstream_port,
623 if proxy_route.upstream_tls { " (tls)" } else { "" },
624 );
625
626 if request.is_websocket_upgrade() {
627 let reunited = read_stream.unsplit(write_stream);
628 return self.handle_proxy_websocket(
629 reunited,
630 request,
631 proxy_route,
632 src_addr,
633 &id,
634 ).await;
635 }
636
637 let proxy_result = self.handle_proxy_http(
638 request,
639 proxy_route,
640 &mut write_stream,
641 src_addr,
642 &id,
643 ).await;
644
645 let (proxy_status, proxy_bytes) = match proxy_result {
646 Ok(sb) => sb,
647 Err(e) => {
648 warn!("{}: proxy error: {}", id, e);
649 let mut resp = HttpMessage::respond_with_text(
650 HttpStatus::BadGateway,
651 "Bad Gateway: upstream proxy error.",
652 );
653 resp.set_connection_close(true);
654 match resp.write_all(&mut write_stream).await {
655 Ok(()) => (),
656 Err(we) => warn!(
657 "{}: failed to emit 502: {}",
658 id, we),
659 }
660 break;
661 }
662 };
663
664 // Record traffic for the proxied request.
665 if let Some(recorder) = self.traffic.as_ref().filter(|_| logged) {
666 let dur_us = req_started_at
667 .elapsed().as_micros() as u64;
668 let record = RequestRecord {
669 when_ns: traffic::now_ns(),
670 vhost: vhost.primary_hostname()
671 .to_string(),
672 method: rec_method.clone(),
673 path: rec_path.clone(),
674 status: proxy_status,
675 peer: fmt!("{}", src_addr),
676 bytes: proxy_bytes,
677 duration_us: dur_us,
678 };
679 if let Err(e) = recorder.record(record) {
680 warn!("{}: traffic recorder rejected entry: {}",
681 id, e);
682 }
683 }
684
685 // Proxied responses are written
686 // directly to the stream. Close
687 // the connection because the
688 // upstream uses Connection: close
689 // and we cannot guarantee
690 // keep-alive semantics.
691 break;
692 }
693 }
694
695 // Wrap the incoming header fields in an
696 // `Arc` so the downstream handler (and
697 // any API / webhook handler it dispatches
698 // to) can read request headers without
699 // another copy. Cloning here rather than
700 // moving because the POST branch still
701 // needs `request.header.fields` a few
702 // lines down.
703 let req_headers = Arc::new(
704 request.header.fields.clone(),
705 );
706 let body = request.body;
707 match method {
708 // A `HEAD` asks what the `GET` would answer, minus the
709 // answer. RFC 9110 9.3.2 says the fields must be the ones
710 // the `GET` would carry, so the only honest way to produce
711 // them is to do the `GET` and withhold the body at the
712 // wire. Handled here rather than left to fall through:
713 // unhandled, it reached no branch at all, no response was
714 // ever built, and the caller sat until the read timed out
715 // -- so `curl -I`, and every uptime monitor that speaks
716 // `HEAD` first, saw a 408 from a server that was fine.
717 HttpMethod::GET | HttpMethod::HEAD => {
718 let head_only = method == HttpMethod::HEAD;
719 // Told to the handler as well as applied at the wire,
720 // so a `GET` path that keeps a tally can decline to
721 // count a request that asked for no prose.
722 let mut loc = loc;
723 if head_only {
724 loc.data.insert(
725 dat!("head_only"),
726 dat!(true),
727 );
728 }
729 let vhost_db = self.db_for_vhost(
730 vhost.primary_hostname(),
731 );
732 let result = vhost.web_handler.handle_get(
733 loc,
734 response,
735 body,
736 req_headers.clone(),
737 vhost_db,
738 &sid_opt,
739 src_addr,
740 &id,
741 ).await;
742 response = res!(result);
743 if head_only {
744 response = response.map(|r| r.head_only());
745 }
746 }
747 HttpMethod::POST => {
748 // Carry Content-Type from the incoming
749 // request into loc.data so the handler
750 // can forward it to the upstream.
751 let mut loc = loc;
752 if let Some((vals, _)) = request.header.fields.get_all(
753 &HeaderName::ContentType,
754 ) {
755 if let Some(v) = vals.first() {
756 loc.data.insert(
757 dat!("content_type"),
758 dat!(v.to_string()),
759 );
760 }
761 }
762 let vhost_db = self.db_for_vhost(
763 vhost.primary_hostname(),
764 );
765 let result = vhost.web_handler.handle_post(
766 loc,
767 response,
768 body,
769 req_headers.clone(),
770 vhost_db,
771 &sid_opt,
772 src_addr,
773 &id,
774 ).await;
775 response = res!(result);
776 }
777 _ => if logged {
778 fault!("{}: Unsupported HTTP request method '{}'.",
779 id, method);
780 },
781 }
782 }
783 }
784 _ => if logged {
785 fault!("{}: Unsupported HTTP '{:?}'.", id, request.header.headline);
786 },
787 }
788
789 alog!(logged, log_level, "Outgoing HTTPS message:");
790 let mut rec_status: u16 = 0;
791 let mut rec_bytes: Option<u64> = None;
792 match response {
793 Some(mut msg) => {
794 // Attach the freshly-issued session cookie to
795 // the outgoing response, if any. Works for both
796 // file responses and redirect responses.
797 if let Some(cookie) = issued_cookie.take() {
798 msg = msg.set_cookie(cookie);
799 }
800 // Inject HSTS header if configured, so browsers
801 // remember to use HTTPS for subsequent visits.
802 if self.cfg.hsts_max_age_secs > 0 {
803 msg.header.fields.insert(
804 HeaderName::StrictTransportSecurity,
805 HeaderFieldValue::Generic(fmt!(
806 "max-age={}; includeSubDomains",
807 self.cfg.hsts_max_age_secs,
808 )),
809 None,
810 );
811 }
812 // Baseline security response headers: cheap
813 // defence in depth against content sniffing,
814 // clickjacking and referrer leakage. Each one
815 // is a hard-coded conservative default; tighten
816 // via per-deployment patches if the defaults
817 // bite an integration.
818 if self.cfg.security_headers_enabled {
819 msg.header.fields.insert(
820 HeaderName::XContentTypeOptions,
821 HeaderFieldValue::Generic(
822 "nosniff".to_string()),
823 None,
824 );
825 msg.header.fields.insert(
826 HeaderName::XFrameOptions,
827 HeaderFieldValue::Generic(
828 "SAMEORIGIN".to_string()),
829 None,
830 );
831 msg.header.fields.insert(
832 HeaderName::ReferrerPolicy,
833 HeaderFieldValue::Generic(
834 "strict-origin-when-cross-origin"
835 .to_string()),
836 None,
837 );
838 // Permissions-Policy: deny every sensor
839 // feature by default. A vhost may replace the
840 // whole policy through its `permissions_policy`
841 // config -- a capture ceremony sets
842 // `camera=(self)`, say -- so a browser-capability
843 // grant is explicit and per-site, never a code
844 // default that loosens every deployment at once.
845 let permissions_policy = vhost.permissions_policy
846 .as_deref()
847 .unwrap_or("accelerometer=(), camera=(), \
848 geolocation=(), gyroscope=(), \
849 magnetometer=(), microphone=(), \
850 payment=(), usb=()");
851 msg.header.fields.insert(
852 HeaderName::PermissionsPolicy,
853 HeaderFieldValue::Generic(
854 permissions_policy.to_string()),
855 None,
856 );
857 }
858 // Content-Security-Policy: only emit if the
859 // operator configured a value. An empty CSP
860 // string is not injected at all because the
861 // absence of the header is different from
862 // `default-src 'none'` -- the former is a
863 // no-op, the latter blocks the whole page.
864 if !self.cfg.content_security_policy.is_empty() {
865 msg.header.fields.insert(
866 HeaderName::ContentSecurityPolicy,
867 HeaderFieldValue::Generic(
868 self.cfg.content_security_policy.clone()),
869 None,
870 );
871 }
872 // Pull the status code out for the
873 // traffic record before the message is
874 // consumed by write_all. HttpStatus is
875 // repr(u16), so a direct cast yields
876 // the wire code (200, 404, etc.).
877 if let HttpHeadline::Response { status } =
878 &msg.header.headline
879 {
880 rec_status = *status as u16;
881 }
882 // Encode the body, if the type gains by it and the
883 // client said it would take one. Done here, at the
884 // last point before the wire, so every response the
885 // server produces is covered by one rule -- a
886 // static file, a rendered page, a JSON answer --
887 // rather than each producer having to remember.
888 //
889 // `Content-Length` is taken from the body when the
890 // message is written, so it describes the encoded
891 // body by construction. Getting that wrong is not a
892 // cosmetic fault: a length that does not match
893 // leaves the client waiting for bytes that never
894 // arrive, and every message after it on a
895 // kept-alive connection is read at the wrong offset.
896 if self.cfg.compression_enabled {
897 let content_type = msg.header.fields
898 .get_one(&HeaderName::ContentType)
899 .map(|v| fmt!("{}", v))
900 .unwrap_or_default();
901 if encoding::is_compressible(&content_type) {
902 let coding = encoding::choose_for(
903 accept_encoding.as_deref(),
904 &content_type,
905 msg.body_len(),
906 self.cfg.compression_min_bytes as usize,
907 );
908 // On the blocking pool: encoding a
909 // megabyte of markup is processor work,
910 // and a single-core host has one async
911 // worker to starve. The handle is asked
912 // for rather than assumed: `current`
913 // panics where there is no runtime, and a
914 // panic in a pool thread is an error the
915 // caller cannot read.
916 msg = match tokio::task::spawn_blocking(move || {
917 let rt = match tokio::runtime::Handle::try_current() {
918 Ok(rt) => rt,
919 Err(e) => return Err(err!(e,
920 "The response encoder found no runtime \
921 to wait on."; Init, Missing)),
922 };
923 rt.block_on(encoding::encode(msg, coding))
924 }).await {
925 Ok(Ok(encoded)) => encoded,
926 Ok(Err(e)) => return Err(err!(e,
927 "{}: Could not encode the response body.", id;
928 IO, Encode)),
929 Err(e) => return Err(err!(e,
930 "{}: The response encoder did not finish.", id;
931 IO, Encode)),
932 };
933 }
934 }
935
936 // The length the response actually carries, which for a
937 // body sent from a file is the window rather than the
938 // empty buffer beside it.
939 rec_bytes = Some(msg.body_len() as u64);
940 match msg.write_all(&mut write_stream).await {
941 Ok(()) => (),
942 Err(e) => return Err(err!(e,
943 "{}: Could not send response.", id;
944 IO, Network, Wire, Write)),
945 }
946 }
947 None => alog!(logged, log_level, " None"),
948 }
949
950 // Emit a traffic record for this request now that
951 // the response has been fully written. Recording
952 // is a bounded short critical section; failures
953 // are logged but never propagated, since the
954 // request itself succeeded and we do not want
955 // the dashboard to break the data path.
956 if let Some(recorder) = self.traffic.as_ref().filter(|_| logged) {
957 let dur_us = req_started_at.elapsed().as_micros() as u64;
958 let record = RequestRecord {
959 when_ns: traffic::now_ns(),
960 vhost: vhost.primary_hostname().to_string(),
961 method: rec_method,
962 path: rec_path,
963 status: rec_status,
964 peer: fmt!("{}", src_addr),
965 bytes: rec_bytes,
966 duration_us: dur_us,
967 };
968 if let Err(e) = recorder.record(record) {
969 warn!("{}: traffic recorder rejected entry: {}",
970 id, e);
971 }
972 }
973 }
974 Some(Err(e)) => {
975 // A reader error is often a configured limit
976 // breach (oversized body, oversized header,
977 // slowloris timeout). Translate those into a
978 // proper HTTP status so the client sees a
979 // deliberate rejection instead of a silent
980 // connection drop, then close the connection.
981 let tags = e.tags();
982 let (status, msg) = if tags.contains(&ErrTag::TooBig) {
983 (HttpStatus::ContentTooLarge,
984 "Request exceeds the configured size limit.")
985 } else if tags.contains(&ErrTag::Timeout) {
986 (HttpStatus::RequestTimeout,
987 "Request timed out while reading headers.")
988 } else {
989 // Returned, the error is logged by the accept loop, so a vhost with its
990 // access log off closes the connection quietly instead.
991 if !logged {
992 break;
993 }
994 warn!("{}: HTTP read error: {}", id, e);
995 return Err(e);
996 };
997 alog!(logged, LogLevel::Warn, "{}: dropping connection ({}): {}", id, status, e);
998 let mut resp = HttpMessage::respond_with_text(status, msg);
999 resp.set_connection_close(true);
1000 match resp.write_all(&mut write_stream).await {
1001 Ok(()) => (),
1002 Err(we) => alog!(logged, LogLevel::Warn,
1003 "{}: failed to emit {} response: {}", id, status, we),
1004 }
1005 break;
1006 }
1007 None => {
1008 break;
1009 }
1010 }
1011 }
1012
1013 // Gracefully close the TLS connection.
1014 let reunited_stream = read_stream.unsplit(write_stream);
1015 let result = reunited_stream.shutdown().await;
1016 // A peer that has already gone makes this fail on most connections, and on a vhost
1017 // with its access log off even a bare line per connection is a record of traffic.
1018 if let Err(e) = result {
1019 if logged {
1020 error!(e.into());
1021 }
1022 }
1023 alog!(logged, log_level, "{}: Connection with {:?} closed.", id, src_addr);
1024
1025 Ok(())
1026 }
1027
1028 fn match_redirect<'a>(
1029 rules: &'a [RedirectRule],
1030 request_path: &str,
1031 )
1032 -> Option<&'a RedirectRule>
1033 {
1034 for rule in rules {
1035 if rule.matches(request_path) {
1036 return Some(rule);
1037 }
1038 }
1039 None
1040 }
1041
1042 async fn handle_proxy_http<W>(
1043 &self,
1044 request: HttpMessage,
1045 route: &ProxyRoute,
1046 client_w: &mut W,
1047 src_addr: SocketAddr,
1048 id: &str,
1049 )
1050 -> Outcome<(u16, Option<u64>)>
1051 where W: AsyncWriteExt + Unpin,
1052 {
1053 // Extract method, request path and the raw query. The query must ride
1054 // through verbatim: an upstream that dispatches on a query parameter
1055 // (e.g. `?view=`) never sees it otherwise, and silently gets the
1056 // default.
1057 let (method, path, query) = match &request.header.headline {
1058 HttpHeadline::Request { method, loc } => {
1059 (fmt!("{}", method), loc.path.as_string().to_string(), loc.query.clone())
1060 }
1061 _ => return Err(err!(
1062 "{}: proxy: request is not an HTTP request.", id;
1063 Invalid, Bug)),
1064 };
1065
1066 let upstream_path = match query.is_empty() {
1067 true => route.upstream_path_for(&path),
1068 false => fmt!("{}?{}", route.upstream_path_for(&path), query),
1069 };
1070
1071 // Who, if anybody, is entitled to have spoken the forwarding headers before this hop.
1072 let policy = res!(ForwardedPolicy::new(&self.cfg.trusted_proxies));
1073
1074 // Connect to the upstream.
1075 let mut upstream = match TcpStream::connect(
1076 (route.upstream_host.as_str(), route.upstream_port),
1077 ).await {
1078 Ok(s) => s,
1079 Err(e) => return Err(err!(e,
1080 "{}: proxy: failed to connect to {}:{}.",
1081 id, route.upstream_host, route.upstream_port;
1082 IO, Network, Init)),
1083 };
1084
1085 // Build the request bytes to send to the upstream. The caller's headers ride through
1086 // except the ones this hop owns, and the forwarding headers this hop appends go last --
1087 // see `fe2o3_net::http::fwd`, which both proxy paths share so a fix to one fixes both.
1088 let req = fwd::build_proxy_request_head(
1089 &method,
1090 &upstream_path,
1091 &route.upstream_host,
1092 &request,
1093 &src_addr,
1094 &policy,
1095 request.body.len(),
1096 );
1097
1098 // Write request to upstream.
1099 match upstream.write_all(req.as_bytes()).await {
1100 Ok(()) => (),
1101 Err(e) => return Err(err!(e,
1102 "{}: proxy: failed to write request to upstream.", id;
1103 IO, Network, Wire, Write)),
1104 }
1105 if !request.body.is_empty() {
1106 match upstream.write_all(&request.body).await {
1107 Ok(()) => (),
1108 Err(e) => return Err(err!(e,
1109 "{}: proxy: failed to write request body to upstream.", id;
1110 IO, Network, Wire, Write)),
1111 }
1112 }
1113 match upstream.flush().await {
1114 Ok(()) => (),
1115 Err(e) => return Err(err!(e,
1116 "{}: proxy: failed to flush upstream.", id;
1117 IO, Network, Wire, Write)),
1118 }
1119
1120 // Stream the response from upstream to client.
1121 // Read into a buffer, find the header/body boundary,
1122 // parse the status code, then stream everything.
1123 let mut buf = vec![0u8; 16384];
1124 let mut accum: Vec<u8> = Vec::new();
1125 let mut status_code: u16 = 0;
1126 let mut total_body_bytes: u64 = 0;
1127 let mut headers_forwarded = false;
1128
1129 loop {
1130 let n = match upstream.read(&mut buf).await {
1131 Ok(0) => break,
1132 Ok(n) => n,
1133 Err(e) => return Err(err!(e,
1134 "{}: proxy: error reading upstream response.", id;
1135 IO, Network, Wire, Read)),
1136 };
1137
1138 if !headers_forwarded {
1139 accum.extend_from_slice(&buf[..n]);
1140 // Look for end-of-headers marker.
1141 if let Some(pos) = accum.windows(4).position(|w| w == b"\r\n\r\n") {
1142 let header_end = pos + 4;
1143 let header_bytes = &accum[..header_end];
1144 let body_start = &accum[header_end..];
1145
1146 // Parse status code from the first line.
1147 if let Some(line_end) = header_bytes.iter().position(|&b| b == b'\r') {
1148 let status_line = String::from_utf8_lossy(&header_bytes[..line_end]);
1149 // Format: "HTTP/1.1 200 OK"
1150 let parts: Vec<&str> = status_line.splitn(3, ' ').collect();
1151 if parts.len() >= 2 {
1152 if let Ok(code) = parts[1].parse::<u16>() {
1153 status_code = code;
1154 }
1155 }
1156 }
1157
1158 // Forward the response headers to the client.
1159 match client_w.write_all(header_bytes).await {
1160 Ok(()) => (),
1161 Err(e) => return Err(err!(e,
1162 "{}: proxy: failed to write response headers to client.", id;
1163 IO, Network, Wire, Write)),
1164 }
1165
1166 // Forward any body bytes that arrived with the headers.
1167 if !body_start.is_empty() {
1168 match client_w.write_all(body_start).await {
1169 Ok(()) => (),
1170 Err(e) => return Err(err!(e,
1171 "{}: proxy: failed to write initial body to client.", id;
1172 IO, Network, Wire, Write)),
1173 }
1174 total_body_bytes += body_start.len() as u64;
1175 }
1176
1177 headers_forwarded = true;
1178 } else if accum.len() > 65536 {
1179 return Err(err!(
1180 "{}: proxy: upstream response headers exceed 64 KiB.", id;
1181 IO, Network, Input, TooBig));
1182 }
1183 } else {
1184 // Stream body chunks directly.
1185 match client_w.write_all(&buf[..n]).await {
1186 Ok(()) => (),
1187 Err(e) => return Err(err!(e,
1188 "{}: proxy: failed to stream body to client.", id;
1189 IO, Network, Wire, Write)),
1190 }
1191 total_body_bytes += n as u64;
1192 }
1193 }
1194
1195 match client_w.flush().await {
1196 Ok(()) => (),
1197 Err(e) => return Err(err!(e,
1198 "{}: proxy: failed to flush client stream.", id;
1199 IO, Network, Wire, Write)),
1200 }
1201
1202 if status_code == 0 {
1203 status_code = 200; // Fallback if parsing failed.
1204 }
1205
1206 log!(log_get_level!(),
1207 "{}: proxy: {} {} -> {} ({} body bytes)",
1208 id, method, path, status_code, total_body_bytes);
1209
1210 Ok((status_code, Some(total_body_bytes)))
1211 }
1212
1213 async fn handle_proxy_websocket<S>(
1214 self,
1215 client: &mut S,
1216 request: HttpMessage,
1217 route: &ProxyRoute,
1218 src_addr: SocketAddr,
1219 id: &str,
1220 )
1221 -> Outcome<()>
1222 where S: AsyncRead + AsyncWrite + Unpin,
1223 {
1224 let path = match &request.header.headline {
1225 HttpHeadline::Request { loc, .. } => loc.path.as_string().to_string(),
1226 _ => return Err(err!(
1227 "{}: proxy ws: request is not an HTTP request.", id;
1228 Invalid, Bug)),
1229 };
1230 let upstream_path = res!(wsproxy::upstream_target(
1231 &request,
1232 &route.upstream_path_for(&path),
1233 ));
1234 let policy = res!(ForwardedPolicy::new(&self.cfg.trusted_proxies));
1235 wsproxy::tunnel_upgrade(
1236 client,
1237 &request,
1238 &route.upstream_host,
1239 route.upstream_port,
1240 &upstream_path,
1241 src_addr,
1242 &policy,
1243 id,
1244 ).await
1245 }
1246
1247 async fn handle_ws_route<S>(
1248 self,
1249 client: &mut S,
1250 request: HttpMessage,
1251 route: &WsRoute,
1252 src_addr: SocketAddr,
1253 id: &str,
1254 )
1255 -> Outcome<()>
1256 where S: AsyncRead + AsyncWrite + Unpin,
1257 {
1258 let upstream_path = res!(wsproxy::upstream_target(&request, &route.upstream_path));
1259 let policy = res!(ForwardedPolicy::new(&self.cfg.trusted_proxies));
1260 wsproxy::tunnel_upgrade(
1261 client,
1262 &request,
1263 &route.upstream_host,
1264 route.upstream_port,
1265 &upstream_path,
1266 src_addr,
1267 &policy,
1268 id,
1269 ).await
1270 }
1271}
1272
1273
1274// ┌───────────────────────────────────────────────────────────────────────────┐
1275// │ TERMINAL DISPATCH GATE │
1276// └───────────────────────────────────────────────────────────────────────────┘
1277
1278/// The outcome of deciding whether a `/term/<name>` upgrade may be dispatched to
1279/// the terminal bridge. Both conditions fail closed: a vhost that never enabled
1280/// terminals is [`NotEnabled`](Self::NotEnabled), and a request with no operator
1281/// session is [`NotAuthenticated`](Self::NotAuthenticated). Only when the vhost
1282/// has terminals and the caller is an authenticated operator is the answer
1283/// [`Dispatch`](Self::Dispatch).
1284///
1285/// This exists so the decision is testable apart from the connection machinery
1286/// in `handle_https`, which cannot be stood up without a TLS stream.
1287#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1288pub enum TerminalGate {
1289 Dispatch,
1290 NotEnabled,
1291 NotAuthenticated,
1292}
1293
1294/// Decides whether to dispatch an upgrade to the terminal bridge.
1295///
1296/// `term_enabled` is whether the vhost carries a `term_config` (its
1297/// `term_manager` is set); `principal` is the operator resolved from the
1298/// request's dashboard cookie, if any. The gate is deliberately conjunctive and
1299/// order-sensitive only for the sake of a precise refusal: configuration is
1300/// checked first, so a host that never asked for terminals never even reports
1301/// whether the caller was authenticated.
1302pub fn terminal_gate(
1303 term_enabled: bool,
1304 principal: Option<&crate::srv::admin::AdminPrincipal>,
1305)
1306 -> TerminalGate
1307{
1308 if !term_enabled {
1309 return TerminalGate::NotEnabled;
1310 }
1311 if principal.is_none() {
1312 return TerminalGate::NotAuthenticated;
1313 }
1314 TerminalGate::Dispatch
1315}
1316
1317
1318#[cfg(test)]
1319mod terminal_gate_tests {
1320 use super::{
1321 terminal_gate,
1322 TerminalGate,
1323 };
1324 use crate::srv::admin::AdminPrincipal;
1325
1326 fn operator() -> AdminPrincipal {
1327 AdminPrincipal {
1328 name: "op".to_string(),
1329 scopes: vec!["dashboard:view".to_string()],
1330 expires_at: u64::MAX,
1331 }
1332 }
1333
1334 // An unauthenticated upgrade to `/term/` is refused even where the feature
1335 // is enabled: a terminal-enabled vhost with no operator session must not
1336 // dispatch. This is the hole that shipped -- the router dispatched every
1337 // `/term/` upgrade regardless of the caller.
1338 #[test]
1339 fn unauthenticated_term_upgrade_is_refused() {
1340 assert_eq!(
1341 terminal_gate(true, None),
1342 TerminalGate::NotAuthenticated,
1343 "an unauthenticated /term/ upgrade must be refused, not dispatched",
1344 );
1345 }
1346
1347 // A vhost with no `term_config` refuses regardless of who is asking: even a
1348 // fully authenticated operator does not reach a shell on a host that never
1349 // enabled terminals. This is the config gate the router never consulted --
1350 // `None` must mean off, failing closed.
1351 #[test]
1352 fn term_without_config_is_refused_even_for_an_operator() {
1353 let op = operator();
1354 assert_eq!(
1355 terminal_gate(false, Some(&op)),
1356 TerminalGate::NotEnabled,
1357 "a vhost with no term_config must refuse /term/ for anyone",
1358 );
1359 // And with no config and no session, still refused, config first.
1360 assert_eq!(terminal_gate(false, None), TerminalGate::NotEnabled);
1361 }
1362
1363 // The one path that dispatches: terminals enabled and an operator present.
1364 #[test]
1365 fn enabled_and_authenticated_dispatches() {
1366 let op = operator();
1367 assert_eq!(terminal_gate(true, Some(&op)), TerminalGate::Dispatch);
1368 }
1369}