Oregami
Repositories/oxedyne/fe2o3

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

68.7 KiB, 422 runs

created by r1870400018:947, 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 handler as admin_handler,
4 ozone_view as admin_ozone_view,
5 state::AdminState,
6 traffic::TrafficRecorder,
7 },
8 api::{
9 self,
10 ApiHandlerRegistry,
11 },
12 cache,
13 cfg::{
14 ApiRoute,
15 ServerConfig,
16 WebhookRoute,
17 },
18 dev::refresh::HtmlModifier,
19 console as site_console,
20 publish::{
21 self,
22 PublishConfig,
23 Subscription,
24 ai as publish_ai,
25 comment as publish_comment,
26 declare as publish_declare,
27 page as publish_page,
28 send::MailSender,
29 store as publish_store,
30 subscribe as publish_subscribe,
31 },
32 webhook::{
33 self,
34 WebhookRegistry,
35 },
36};
37
38use oxedyne_fe2o3_core::{
39 prelude::*,
40 file::{
41 OsPath,
42 },
43 map::MapMut,
44 path::NormalPath,
45 rand::Rand,
46};
47use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
48use oxedyne_fe2o3_iop_db::api::Database;
49use oxedyne_fe2o3_iop_hash::api::Hasher;
50use oxedyne_fe2o3_jdat::{
51 prelude::*,
52 id::NumIdDat,
53};
54use oxedyne_fe2o3_net::{
55 file::RequestPath,
56 http::{
57 client::{
58 http_request,
59 https_request,
60 },
61 encoding,
62 fields::{
63 HeaderFields,
64 HeaderFieldValue,
65 HeaderName,
66 },
67 handler::WebHandler,
68 header::HttpMethod,
69 loc::HttpLocator,
70 msg::{
71 FileWindow,
72 HttpMessage,
73 },
74 range::{
75 self,
76 RangeOutcome,
77 },
78 status::HttpStatus,
79 },
80};
81
82use std::{
83 collections::BTreeMap,
84 fmt::Debug,
85 net::SocketAddr,
86 path::{
87 Path,
88 PathBuf,
89 },
90 sync::{
91 Arc,
92 RwLock,
93 },
94};
95
96use tokio::{
97 self,
98 io::AsyncReadExt,
99};
100use tokio_rustls::rustls::ClientConfig;
101
102
103#[derive(Clone, Debug)]
104pub struct AppWebHandler<
105 M: MapMut<String, OsPath> + Clone + Debug + Send + Sync,
106>{
107 // Config
108 pub cfg: ServerConfig,
109 // State
110 pub public_dir: PathBuf,
111 pub static_routes: M,
112 pub default_index_files: Vec<String>,
113 pub dev_mode: bool,
114 pub api_routes: Vec<ApiRoute>,
115 pub webhook_routes: Vec<WebhookRoute>,
116 pub webhook_registry: Arc<WebhookRegistry>,
117 pub api_handler_registry: Arc<ApiHandlerRegistry>,
118 pub tls_client: Option<Arc<ClientConfig>>,
119 pub admin_state: Option<Arc<AdminState>>,
120 pub traffic: Option<Arc<TrafficRecorder>>,
121 pub publish: Option<Arc<PublishConfig>>,
122 pub mail: Option<Arc<MailSender>>,
123 pub site_admins: Arc<Vec<String>>,
124}
125
126impl<
127 M: MapMut<String, OsPath> + Clone + Debug + Send + Sync,
128>
129 AppWebHandler<M>
130{
131 pub fn new(
132 cfg: ServerConfig,
133 public_dir: PathBuf,
134 static_routes: M,
135 default_index_files: Vec<String>,
136 dev_mode: bool,
137 api_routes: Vec<ApiRoute>,
138 webhook_routes: Vec<WebhookRoute>,
139 webhook_registry: Arc<WebhookRegistry>,
140 api_handler_registry: Arc<ApiHandlerRegistry>,
141 tls_client: Option<Arc<ClientConfig>>,
142 admin_state: Option<Arc<AdminState>>,
143 traffic: Option<Arc<TrafficRecorder>>,
144 publish: Option<Arc<PublishConfig>>,
145 mail: Option<Arc<MailSender>>,
146 site_admins: Arc<Vec<String>>,
147 )
148 -> Self
149 {
150 Self {
151 cfg,
152 public_dir,
153 static_routes,
154 default_index_files,
155 dev_mode,
156 api_routes,
157 webhook_routes,
158 webhook_registry,
159 api_handler_registry,
160 tls_client,
161 admin_state,
162 traffic,
163 publish,
164 mail,
165 site_admins,
166 }
167 }
168
169 pub fn err_id() -> String {
170 Rand::generate_random_string(6, "abcdefghikmnpqrstuvw0123456789")
171 }
172
173 async fn router(
174 &self,
175 loc: &HttpLocator,
176 id: &String,
177 )
178 -> Outcome<PathBuf>
179 {
180 let route = loc.path.as_string();
181
182 // First try configured static routes.
183 match self.static_routes.get(route) {
184 Some(os_path) => match os_path {
185 OsPath::Dir(path) => {
186 // Path is already normalised and absolute.
187 for filename in &self.default_index_files {
188 let candidate = path.clone().join(filename);
189 if candidate.exists() {
190 return Ok(candidate);
191 }
192 }
193 return Err(err!(
194 "{}: No default index files found in directory {:?}. \
195 Tried: {:?}", id, path, self.default_index_files;
196 File, NotFound));
197 }
198 OsPath::File(path) => return Ok(path.clone()),
199 }
200 None => {
201 // Fallback: try to serve directly from public directory.
202 let clean_path = if route.starts_with('/') {
203 &route[1..]
204 } else {
205 route
206 };
207
208 let path = Path::new(clean_path).normalise();
209 if path.escapes() {
210 return Err(err!(
211 "{}: Request path '{}' would escape the public directory.",
212 id, route;
213 Invalid, Path, Security));
214 }
215
216 let full_path = self.public_dir.clone().join(path);
217
218 // If it's a directory, try index files.
219 if full_path.is_dir() {
220 for filename in &self.default_index_files {
221 let candidate = full_path.join(filename);
222 if candidate.exists() {
223 return Ok(candidate);
224 }
225 }
226 }
227
228 return Ok(full_path);
229 }
230 }
231 }
232}
233
234impl<
235 M: MapMut<String, OsPath> + Clone + Debug + Send + Sync,
236>
237 WebHandler for AppWebHandler<M>
238{
239
240 fn handle_get<
241 const SIDL: usize,
242 const UIDL: usize,
243 SID: NumIdDat<SIDL> + 'static,
244 UID: NumIdDat<UIDL> + 'static,
245 ENC: Encrypter,
246 KH: Hasher,
247 DB: Database<UIDL, UID, ENC, KH>,
248 >(
249 &self,
250 loc: HttpLocator,
251 _response: Option<HttpMessage>,
252 _body: Vec<u8>,
253 req_headers: Arc<HeaderFields>,
254 db: Option<(Arc<RwLock<DB>>, UID)>,
255 _sid_opt: &Option<SID>,
256 peer: SocketAddr,
257 id: &String,
258 )
259 -> impl std::future::Future<Output = Outcome<Option<HttpMessage>>> + Send
260 {
261 let dev_mode = self.dev_mode;
262 let rpath = loc.path.clone();
263 let id = id.to_string();
264 let request_path = loc.path.as_string().to_string();
265 // The path and the query are parsed apart, so a handler wanting the query must be handed
266 // it: `request_path` never carries one, and splitting it on `?` finds nothing.
267 let request_query = loc.query.clone();
268 let api_routes = self.api_routes.clone();
269 let api_handler_registry = self.api_handler_registry.clone();
270 let tls_client = self.tls_client.clone();
271 let admin_state = self.admin_state.clone();
272 let publish = self.publish.clone();
273 let site_admins = self.site_admins.clone();
274 // A `HEAD` reached this branch because it asks what the `GET` would
275 // answer; the dispatch said so in `loc.data`. Read once here, since two
276 // things below turn on it: the read tally, and whether a `Range` is
277 // honoured at all.
278 let head_only = matches!(
279 loc.data.get(&dat!("head_only")),
280 Some(Dat::Bool(true)),
281 );
282
283 async move {
284 // The dashboard owns the entire `/admin` and `/admin/*`
285 // subtree on every vhost it is configured for. Dispatch
286 // before any API/webhook/static route lookups so app
287 // routes cannot accidentally shadow it.
288 //
289 // Two-stage dispatch: ozone-prefixed routes go to the
290 // generic ozone_view module (which needs the per-vhost
291 // db typed parameters in scope), all other dashboard
292 // routes go to the non-generic handler module.
293 if request_path == "/admin/database"
294 || request_path.starts_with("/admin/database/")
295 {
296 if let Some(state) = &admin_state {
297 let resp = res!(admin_ozone_view::handle_get(
298 state.as_ref(),
299 db.as_ref(),
300 &request_path,
301 &request_query,
302 &req_headers,
303 &id,
304 ).await);
305 return Ok(Some(resp));
306 }
307 return Ok(Some(cache::generated(HttpMessage::respond_with_text(
308 HttpStatus::NotFound,
309 "Not found.",
310 ))));
311 }
312 if request_path == "/admin"
313 || request_path.starts_with("/admin/")
314 {
315 if let Some(state) = &admin_state {
316 let resp = res!(admin_handler::handle_get(
317 state.as_ref(),
318 &request_path,
319 &req_headers,
320 peer,
321 &id,
322 ).await);
323 return Ok(Some(resp));
324 }
325 // Dashboard not configured. Pretend the route does
326 // not exist so we do not leak the existence of an
327 // admin endpoint.
328 return Ok(Some(cache::generated(HttpMessage::respond_with_text(
329 HttpStatus::NotFound,
330 "Not found.",
331 ))));
332 }
333
334 // The site console, at `/manage`, before static files. A site that
335 // has content to manage claims the prefix -- so a signed-in member
336 // can learn their id and ask to be an admin even before anyone is
337 // one, which is the bootstrap. A site that neither publishes nor
338 // has admins skips this and `/manage` means what it did before. The
339 // gate is inside: a listed member manages, a signed-in member who
340 // is not listed is shown their id, an anonymous visitor is sent home.
341 if (publish.is_some() || !site_admins.is_empty())
342 && site_console::owns(&request_path)
343 {
344 let resp = res!(site_console::handle_get(
345 site_admins.as_ref(),
346 admin_state.as_deref(),
347 publish.as_deref(),
348 db.as_ref(),
349 &request_path,
350 &request_query,
351 &req_headers,
352 &id,
353 ).await);
354 return Ok(Some(resp));
355 }
356
357 // The published prose owns its prefix and everything under
358 // it, so a post named like a file on disk is still the
359 // post. Dispatched before API and static routes for the
360 // same reason the dashboard is: a prefix a vhost has
361 // claimed should not be shadowed by what happens to sit in
362 // its webroot. A vhost publishing nothing skips this
363 // entirely and the paths mean whatever they meant before.
364 if let Some(cfg) = &publish {
365 if cfg.owns(&request_path) {
366 // The newsletter's public endpoints sit under the same prefix and touch the
367 // subscriber store, not the posts, so they are answered before the posts are
368 // read. A confirm or unsubscribe is a GET, since it is followed from an email;
369 // the sign-up form is a GET too, and its POST is handled in `handle_post`.
370 match cfg.subscription_of(&request_path) {
371 Some(Subscription::Subscribe) => {
372 return Ok(Some(publish_subscribe::subscribe_form(cfg.as_ref())));
373 }
374 Some(Subscription::Confirm) => {
375 let resp = res!(publish_subscribe::handle_confirm(
376 cfg.as_ref(), db.as_ref(), &request_query, &id));
377 return Ok(Some(resp));
378 }
379 Some(Subscription::Unsubscribe) => {
380 let resp = res!(publish_subscribe::handle_unsubscribe(
381 cfg.as_ref(), db.as_ref(), &request_query, &id));
382 return Ok(Some(resp));
383 }
384 None => {}
385 }
386 // What the site declares about itself and about the things it shows that are
387 // not posts. Answered before the posts are read, for the same reason the
388 // picture below is: a front page drawing a mark beside a book would otherwise
389 // pull every post and its rendered prose to find one word.
390 if request_path == cfg.declare_path() {
391 let keys: Vec<String> =
392 cfg.declare.items.iter().map(|i| i.key.clone()).collect();
393 let levels = match db.as_ref() {
394 Some(dbh) => match publish_store::get_levels(dbh, &keys, &id) {
395 Ok(levels) => levels,
396 Err(e) => {
397 // The site's own declarations will not read. Say nothing rather
398 // than guess: an unread level must not become an undeclared
399 // work that then reads as a claim nobody made.
400 error!(e, "{}: publish: the declarations will not read", id);
401 BTreeMap::new()
402 }
403 },
404 None => BTreeMap::new(),
405 };
406 return Ok(Some(res!(publish_declare::serve_json(
407 &cfg.declare, &levels, &id))));
408 }
409 // A member's uploaded picture: bytes in the site's own database, asked for by
410 // the byline that points at it. Answered here, before the posts are read,
411 // because it is not a post and reading every post to serve a picture would be
412 // work for nothing.
413 if let Some(user) = request_path.strip_prefix(&cfg.avatar_prefix()) {
414 let found = match db.as_ref() {
415 Some(dbh) => match publish_store::get_avatar(dbh, user) {
416 Ok(found) => found,
417 Err(e) => {
418 error!(e, "{}: publish: the picture of '{}' will not read",
419 id, user);
420 None
421 }
422 },
423 None => None,
424 };
425 return Ok(Some(match found {
426 Some((kind, bytes)) => {
427 info!("{}: publish: picture of '{}', {} bytes",
428 id, user, bytes.len());
429 HttpMessage::new_response(HttpStatus::OK)
430 .with_field(
431 HeaderName::ContentType,
432 HeaderFieldValue::Generic(kind),
433 )
434 // A picture changes when its owner changes it and not otherwise,
435 // so it is worth an hour of a reader's cache -- long enough to
436 // matter over a page of bylines, short enough that a new one
437 // shows up the same day.
438 .with_field(
439 HeaderName::CacheControl,
440 HeaderFieldValue::Generic(fmt!("public, max-age=3600")),
441 )
442 // An SVG is a document, and a document served from this origin
443 // may carry script that runs as this site. The picture is drawn
444 // in an `<img>`, where no browser runs it, but a URL can be
445 // opened directly -- so the response is sandboxed and given
446 // nothing it may fetch or execute. The same headers cost a PNG
447 // nothing, so they are not conditional on the type: a rule that
448 // applies sometimes is a rule that will one day be missed.
449 .with_field(
450 HeaderName::ContentSecurityPolicy,
451 HeaderFieldValue::Generic(
452 fmt!("default-src 'none'; style-src 'unsafe-inline'; \
453 sandbox")),
454 )
455 // And it is what it says it is: no sniffing a mistyped upload
456 // into something with a script in it.
457 .with_field(
458 HeaderName::XContentTypeOptions,
459 HeaderFieldValue::Generic(fmt!("nosniff")),
460 )
461 .with_body(bytes)
462 }
463 // A miss is not held. RFC 9110 15.5.5 makes a `404`
464 // heuristically cacheable, so a browser that asked
465 // before the member uploaded anything would keep the
466 // miss and never show the picture they then chose.
467 None => cache::generated(HttpMessage::respond_with_text(
468 HttpStatus::NotFound,
469 "no such picture",
470 )),
471 }));
472 }
473 // The read is the only part that touches the database, so the generics stop here
474 // and the renderers take a slice of posts.
475 let posts = match publish::read(cfg.as_ref(), db.as_ref(), &id) {
476 Ok(posts) => posts,
477 Err(e) => {
478 // The site cannot read its own prose. That is the site's fault, and a
479 // reader should be told rather than shown an empty shelf that looks
480 // like the truth.
481 error!(e, "{}: publish: cannot read the posts", id);
482 return Ok(Some(HttpMessage::respond_with_text(
483 HttpStatus::InternalServerError,
484 "the posts cannot be read",
485 )));
486 }
487 };
488 // The distinct authors the posts name, resolved to a face -- a display name and
489 // an avatar -- for an author row. Two requests draw one: the index, which renders
490 // its own filter, and the JSON, which hands the same faces to a page that renders
491 // the filter itself. A post view, the feed and the filter script want none of it,
492 // and a read per author on every one of those would be work nobody asked for. A
493 // directory has a database for none of this, and a post from one names no author
494 // regardless.
495 let author_names: Vec<String> =
496 if request_path == cfg.path || request_path == cfg.json_path() {
497 // Everyone the list names, for the filter's author row, and after them
498 // whoever else may write here, for the block above it saying what the
499 // site is about. A blog whose first post is not written yet still has
500 // someone who can say what it will be about, and this is how their
501 // description reaches an empty index.
502 let mut names: Vec<String> = posts.iter()
503 .map(|p| p.author.clone())
504 .filter(|a| !a.is_empty())
505 .collect();
506 names.extend(site_admins.iter().cloned());
507 if let Some(dbh) = db.as_ref() {
508 match publish_store::admins_get(dbh, &id) {
509 Ok(granted) => names.extend(granted),
510 Err(e) => warn!(
511 "{}: publish: the granted admins will not read: {}", id, e),
512 }
513 }
514 names
515 } else {
516 // A post being read wants one face: whoever wrote it, for the byline and
517 // the note beneath it. The feed and the filter script want none.
518 publish_page::served_post(cfg.as_ref(), &posts, &request_path)
519 .map(|p| p.author.clone())
520 .filter(|a| !a.is_empty())
521 .into_iter()
522 .collect()
523 };
524 let authors = if author_names.is_empty() {
525 Vec::new()
526 } else {
527 match db.as_ref() {
528 Some(dbh) => publish_store::resolve_authors(dbh, &author_names),
529 None => Vec::new(),
530 }
531 };
532 // The conversation below the post, where the site has a database to keep one
533 // in. A comments read that fails costs the conversation and never the prose: a
534 // reader came for the post.
535 // The site's own switch, which the console sets; the config is only where it
536 // started. Closing stops new comments and does **not** hide the ones already
537 // published: a conversation that happened still happened, and taking it off the
538 // page would be a deletion nobody asked for.
539 let open = publish_comment::comments_open(db.as_ref(), cfg.comments);
540 let served = publish_page::served_post(cfg.as_ref(), &posts, &request_path);
541 let view = match (db.as_ref(), served) {
542 (Some(dbh), Some(post)) => {
543 match publish_comment::site_secret(dbh) {
544 Ok(secret) => {
545 // The order a reader asked for, and the page of it they are on.
546 let newest = publish_page::query_word(&request_query, "order")
547 .as_deref() == Some("newest");
548 let ranker = if newest {
549 publish_comment::Ranker::Recent
550 } else {
551 publish_comment::Ranker::Chronological
552 };
553 let all = match publish_comment::public_for_post(
554 dbh, &post.slug, ranker, &id,
555 ) {
556 Ok(items) => publish_comment::thread(items),
557 Err(e) => {
558 warn!("{}: publish: comments on '{}' will not read: {}",
559 id, post.slug, e);
560 Vec::new()
561 }
562 };
563 let count = publish_comment::count_threads(&all);
564 let want = publish_page::query_word(&request_query, "cpage")
565 .and_then(|p| p.parse::<usize>().ok())
566 .unwrap_or(1);
567 let (threads, at, pages) = publish_comment::page_of(all, want);
568 Some(publish_page::CommentsView {
569 threads,
570 count,
571 page: at,
572 pages,
573 order: if newest { "newest" } else { "oldest" },
574 path: cfg.path_of(&post.slug),
575 challenge: publish_comment::pow_challenge(&post.slug, &secret),
576 said: publish_page::said_of(&request_query),
577 open,
578 // Whoever holds the cookie for a comment on this page may
579 // still correct it, if the window stands.
580 editable: publish_page::edit_claim(&req_headers)
581 .filter(|(cid, token)| {
582 publish_comment::edit_token_ok(cid, &secret, token)
583 }),
584 })
585 }
586 Err(e) => {
587 warn!("{}: publish: no comment secret: {}", id, e);
588 None
589 }
590 }
591 }
592 _ => None,
593 };
594 let resp = res!(publish_page::handle_get(
595 cfg.as_ref(),
596 &posts,
597 &authors,
598 &request_path,
599 &request_query,
600 view.as_ref(),
601 &id,
602 ));
603 // The tally, where a post was actually served to somebody who is neither its
604 // author nor a machine. It is kept here because this is the last place holding
605 // the database: the renderers take a slice of posts on purpose.
606 //
607 // A tally that cannot be written costs the tally and never the page. A reader
608 // asked for prose, and a counter failing is not their problem; it is logged and
609 // the response goes out regardless.
610 if let (Some(db), Some(post)) = (
611 db.as_ref(),
612 publish_page::served_post(cfg.as_ref(), &posts, &request_path),
613 ) {
614 let ua = match req_headers.get_one(&HeaderName::UserAgent) {
615 Some(HeaderFieldValue::Generic(s)) => Some(s.as_str()),
616 _ => None,
617 };
618 let seen = publish::counts_as_read(
619 site_console::session::cookie_value(&req_headers).is_some(),
620 ua,
621 // Set by the dispatch where the request was a HEAD, which asked
622 // for no prose and so read none.
623 head_only,
624 );
625 if seen {
626 if let Err(e) = publish_store::reads_bump(db, &post.slug) {
627 warn!("{}: publish: the read of '{}' was not counted: {}",
628 id, post.slug, e);
629 }
630 }
631 }
632 return Ok(Some(resp));
633 }
634 }
635
636 // Check API routes before falling through to static file
637 // serving. Handler-mode routes dispatch in-process;
638 // proxy-mode routes forward to the upstream over TLS or
639 // plain HTTP depending on the scheme flag. A proxy-mode
640 // GET is useful for loopback app binaries that respond to
641 // plain GETs (e.g. dynamic JSON endpoints), and for
642 // third-party APIs that the app wants to mirror through
643 // Steel so browser calls pick up the operator-supplied
644 // headers (auth tokens, rate-limit keys).
645 if let Some(route) = api_routes.iter().find(|r| r.path == request_path) {
646 if route.handler.is_some() {
647 debug!("{}: GET {} -> api handler '{}'",
648 id, request_path,
649 route.handler.as_deref().unwrap_or("?"));
650 let resp = res!(api::dispatch(
651 &api_handler_registry,
652 route,
653 HttpMethod::GET,
654 &loc,
655 &[],
656 &req_headers,
657 &tls_client,
658 &id,
659 ).await);
660 return Ok(Some(resp));
661 }
662 if route.upstream_host.is_some() {
663 let resp = res!(forward_api_proxy(
664 route,
665 HttpMethod::GET,
666 &loc,
667 &[],
668 &req_headers,
669 &tls_client,
670 &id,
671 ).await);
672 return Ok(Some(resp));
673 }
674 }
675
676 let abs_path = match self.router(&loc, &id).await {
677 Ok(path) => path, // The path may not exist, but at least we have one.
678 Err(e) => {
679 // Tap out early if the route is definitely not known.
680 error!(e);
681 // A file that is not there today may be there after the next
682 // deploy, and a `404` a store was free to invent a lifetime
683 // for outlives the absence it describes.
684 return Ok(Some(cache::generated(
685 HttpMessage::respond_with_text(
686 HttpStatus::NotFound,
687 "File not found.",
688 ).with_field(
689 HeaderName::ContentType,
690 RequestPath::content_type(rpath.as_path()),
691 )
692 )));
693 }
694 };
695
696 let id_clone = id.clone();
697 let req_headers_clone = req_headers.clone();
698 let static_max_age_secs = self.cfg.static_max_age_secs;
699 let fingerprint_secs = self.cfg.fingerprint_max_age_secs;
700 let compression_min_bytes = if self.cfg.compression_enabled {
701 self.cfg.compression_min_bytes as usize
702 } else {
703 usize::MAX // No body reaches the floor, so none is encoded.
704 };
705 let result = tokio::task::spawn_blocking(move || {
706 // Asked for rather than assumed: `current` panics where there is no
707 // runtime, and a panic in a pool thread is an error the caller
708 // cannot read.
709 let rt = match tokio::runtime::Handle::try_current() {
710 Ok(rt) => rt,
711 Err(e) => return Err(err!(e,
712 "{}: The file reader found no runtime to wait on.", id_clone;
713 Init, Missing)),
714 };
715 rt.block_on(async {
716 Ok(match tokio::fs::File::open(&abs_path).await {
717 Ok(mut file) => {
718 let content_type = RequestPath::content_type(abs_path.as_path());
719 let content_type_str = content_type.to_string();
720
721 // The filesystem is asked before the file is read: the
722 // length answers a range, the modification time makes an
723 // entity tag, and a client that already holds the entity
724 // needs no body at all -- reading one would be the very
725 // cost the tag exists to avoid.
726 let meta = res!(file.metadata().await);
727 // A directory opens like a file on Linux and reads like
728 // nothing at all, so the route that reached one is a route
729 // to nothing.
730 if !meta.is_file() {
731 debug!("{}: {:?} is not a file.", id_clone, abs_path);
732 return Ok(cache::generated(
733 HttpMessage::respond_with_text(
734 HttpStatus::NotFound,
735 "File not found.",
736 )));
737 }
738 let total = meta.len();
739
740 // In development an entry document is rewritten on the
741 // way out to carry the refresh hook, so the file on disk
742 // is not the entity being sent. Such a document is never
743 // cached, never revalidated and never cut into windows --
744 // it is simply resent, whole, every time.
745 let verbatim = !(dev_mode && cache::is_document(&content_type_str));
746
747 // Which form of this file the client would end up
748 // holding, decided here so the conditional request
749 // below is answered about the copy it actually has.
750 // The `200` carries the plain tag and the encoder
751 // names the coding in it on the way out, so the
752 // naming happens exactly once.
753 let coding = encoding::choose(
754 &req_headers_clone,
755 &content_type_str,
756 total as usize,
757 compression_min_bytes,
758 );
759 let varies = encoding::is_compressible(&content_type_str);
760
761 let validators = if verbatim {
762 let etag = res!(cache::entity_tag(&meta));
763 let directive = cache::cache_control(
764 &content_type_str,
765 abs_path.as_path(),
766 static_max_age_secs,
767 fingerprint_secs,
768 );
769 Some((etag, directive))
770 } else {
771 None
772 };
773
774 if let Some((etag, directive)) = &validators {
775 // Either form of the tag settles it. A client
776 // holding the encoded copy sends the tag with
777 // the coding in it; one holding the plain copy
778 // sends the plain tag, and is equally entitled
779 // to be told it is current -- what it holds is
780 // what it would get.
781 let encoded_tag = encoding::tagged(etag, coding);
782 let held = if cache::is_current(&req_headers_clone, &encoded_tag) {
783 Some(encoded_tag)
784 } else if cache::is_current(&req_headers_clone, etag) {
785 Some(etag.clone())
786 } else {
787 None
788 };
789 if let Some(held) = held {
790 debug!("{}: {:?} is unchanged; 304.", id_clone, abs_path);
791 return Ok(res!(cache::not_modified(
792 held,
793 directive.clone(),
794 varies,
795 )));
796 }
797 }
798
799 if !verbatim {
800 // The development refresh hook, injected into a
801 // document that is read whole because it is rewritten
802 // whole.
803 let mut contents = Vec::new();
804 match file.read_to_end(&mut contents).await {
805 Ok(_n) => (),
806 Err(e) => {
807 let err_id = Self::err_id();
808 error!(e.into(),
809 "{}: While trying to serve file '{:?}' (err_id: {})",
810 id_clone, abs_path, err_id,
811 );
812 return Ok(HttpMessage::respond_with_text(
813 HttpStatus::InternalServerError,
814 fmt!("Problem during request processing \
815 (err_id: {}).", err_id),
816 ));
817 }
818 }
819 let contents_str = res!(String::from_utf8(contents));
820 let modified =
821 res!(HtmlModifier::inject_dev_refresh(&contents_str));
822 return Ok(HttpMessage::new_response(HttpStatus::OK)
823 .with_field(HeaderName::ContentType, content_type)
824 .with_body(modified.into_bytes()));
825 }
826
827 // The bytes on disk are the bytes sent, so a window of
828 // them can be asked for and answered. Every such response
829 // says so, because a browser will not offer a scrubber on
830 // a video the server has not advertised.
831 //
832 // Except to a `HEAD`. RFC 9110 14.2 defines range
833 // handling for `GET` alone and requires a server to
834 // ignore the field on any other method, so a `HEAD`
835 // carrying one is answered about the whole file: a `200`
836 // and the length of all of it, which is what the asker
837 // wanted to know. Answering `206` there would name a
838 // window nobody can read and understate the size of what
839 // they were asking about.
840 let outcome = if head_only {
841 RangeOutcome::Whole
842 } else {
843 range::resolve(&req_headers_clone, total)
844 };
845 if let RangeOutcome::NotSatisfiable = outcome {
846 debug!("{}: {:?} cannot answer the range asked of it; 416.",
847 id_clone, abs_path);
848 return Ok(range::not_satisfiable(total));
849 }
850
851 let (status, window) = match outcome {
852 RangeOutcome::Partial(w) => {
853 debug!("{}: {:?} {}.", id_clone, abs_path, w.content_range());
854 (HttpStatus::PartialContent, Some(w))
855 }
856 _ => (HttpStatus::OK, None),
857 };
858
859 let mut response = range::with_accept_ranges(
860 HttpMessage::new_response(status)
861 .with_field(HeaderName::ContentType, content_type)
862 );
863
864 if let Some((etag, directive)) = validators {
865 response = response
866 .with_field(HeaderName::ETag, res!(
867 HeaderFieldValue::new(&HeaderName::ETag, &etag)))
868 .with_field(HeaderName::CacheControl, res!(
869 HeaderFieldValue::new(
870 &HeaderName::CacheControl, &directive)));
871 }
872
873 // The body is named rather than read: a window of a two
874 // gigabyte recording weighs what the window weighs, and
875 // the file goes to the socket a chunk at a time when the
876 // message is written.
877 match window {
878 Some(w) => response
879 .with_field(
880 HeaderName::ContentRange,
881 range::content_range_field(&w),
882 )
883 .with_file_window(FileWindow::new(
884 abs_path.clone(), w.start, w.len())),
885 None => response.with_file_window(FileWindow::new(
886 abs_path.clone(), 0, total)),
887 }
888 }
889 Err(_e) => {
890 debug!("{}: File {:?} not found.", id_clone, abs_path);
891 cache::generated(HttpMessage::respond_with_text(
892 HttpStatus::NotFound,
893 "File not found.",
894 ).with_field(
895 HeaderName::ContentType,
896 RequestPath::content_type(abs_path.as_path()),
897 ))
898 }
899 })
900 })
901 });
902
903 match result.await {
904 Ok(response) => match response {
905 Ok(http_msg) => Ok(Some(http_msg)),
906 Err(e) => Err(e),
907 },
908 Err(e) => Err(err!(e,
909 "{}: Error while executing async file read.", id;
910 IO, File, Read)),
911 }
912 }
913
914 }
915
916 fn handle_post<
917 const SIDL: usize,
918 const UIDL: usize,
919 SID: NumIdDat<SIDL> + 'static,
920 UID: NumIdDat<UIDL> + 'static,
921 ENC: Encrypter,
922 KH: Hasher,
923 DB: Database<UIDL, UID, ENC, KH>,
924 >(
925 &self,
926 loc: HttpLocator,
927 _response: Option<HttpMessage>,
928 body: Vec<u8>,
929 req_headers: Arc<HeaderFields>,
930 db: Option<(Arc<RwLock<DB>>, UID)>,
931 _sid_opt: &Option<SID>,
932 peer: SocketAddr,
933 id: &String,
934 )
935 -> impl std::future::Future<Output = Outcome<Option<HttpMessage>>> + Send
936 {
937 let request_path = loc.path.as_string().to_string();
938 // A one-click unsubscribe arrives as a POST carrying its token in the query, so this
939 // handler needs the query as well as the path.
940 let request_query = loc.query.clone();
941 let id = id.to_string();
942 let api_routes = self.api_routes.clone();
943 let webhook_routes = self.webhook_routes.clone();
944 let webhook_registry = self.webhook_registry.clone();
945 let api_handler_registry = self.api_handler_registry.clone();
946 let tls_client = self.tls_client.clone();
947 let admin_state = self.admin_state.clone();
948 let publish = self.publish.clone();
949 let mail = self.mail.clone();
950 let site_admins = self.site_admins.clone();
951
952 async move {
953 // The newsletter sign-up: a public POST under the published prefix,
954 // answered before the console and the API routes. It touches the
955 // subscriber store and the DKIM mail sender, not the posts, and
956 // always answers the same themed page whether or not the address was
957 // already known -- so the form is never an oracle for the list.
958 if let Some(cfg) = &publish {
959 if let Some(Subscription::Subscribe) = cfg.subscription_of(&request_path) {
960 let peer_ip = peer.ip().to_string();
961 let resp = res!(publish_subscribe::handle_subscribe(
962 cfg.as_ref(), db.as_ref(), &mail, &body, Some(&peer_ip), &id).await);
963 return Ok(Some(resp));
964 }
965 // An unsubscribe by POST, which is what one-click means: the newsletter carries
966 // `List-Unsubscribe-Post`, and a mail client that honours it posts to the same URL
967 // rather than opening a browser. The token rides in the query exactly as it does
968 // for the link, so the two paths remove the same subscriber by the same means.
969 if let Some(Subscription::Unsubscribe) = cfg.subscription_of(&request_path) {
970 let resp = res!(publish_subscribe::handle_unsubscribe(
971 cfg.as_ref(), db.as_ref(), &request_query, &id));
972 return Ok(Some(resp));
973 }
974 // A comment on a post: a public POST under the post's own path, so which post is
975 // being commented on is carried by the URL and cannot be swapped in the body. The
976 // reader is answered with a redirect back to the post, carrying what to tell them --
977 // so a reload does not post the comment twice.
978 // An edit: only from whoever holds the token this comment was answered with, and
979 // only while the window stands.
980 if let Some(slug) = cfg.comment_edit_slug(&request_path)
981 .filter(|_| publish_comment::comments_open(db.as_ref(), cfg.comments))
982 {
983 if let Some(dbh) = db.as_ref() {
984 let secret = res!(publish_comment::site_secret(dbh));
985 let cid = site_console::form_field(&body, "id").unwrap_or_default();
986 let token = site_console::form_field(&body, "token").unwrap_or_default();
987 let source = site_console::form_field(&body, "body").unwrap_or_default();
988 let now = std::time::SystemTime::now()
989 .duration_since(std::time::UNIX_EPOCH)
990 .map(|d| d.as_secs())
991 .unwrap_or(0);
992 let ok = match res!(publish_comment::get(dbh, slug, &cid)) {
993 Some(c) => {
994 publish_comment::edit_token_ok(&cid, &secret, &token)
995 && publish_comment::editable(&c, now)
996 }
997 None => false,
998 };
999 // One answer whether the token was wrong, the window had closed or the
1000 // comment was never there: none of those is anybody's business to learn by
1001 // asking.
1002 let said = if ok
1003 && res!(publish_comment::edit(dbh, slug, &cid, &source))
1004 {
1005 "edited"
1006 } else {
1007 "noedit"
1008 };
1009 return Ok(Some(publish_page::comment_posted(cfg.as_ref(), slug, said)));
1010 }
1011 }
1012 // A preview: renders and stores nothing. Answered before the comment route, since
1013 // its path is the comment path with a suffix.
1014 if let Some(_slug) = cfg.comment_preview_slug(&request_path)
1015 .filter(|_| publish_comment::comments_open(db.as_ref(), cfg.comments))
1016 {
1017 if let Some(dbh) = db.as_ref() {
1018 let secret = res!(publish_comment::site_secret(dbh));
1019 let source = site_console::form_field(&body, "body").unwrap_or_default();
1020 let html = res!(publish_comment::preview(
1021 dbh,
1022 &source,
1023 Some(&peer.ip().to_string()),
1024 &secret,
1025 cfg.comment_rate_secs,
1026 ));
1027 return Ok(Some(publish_page::comment_preview(html)));
1028 }
1029 }
1030 if let Some(slug) = cfg.comment_slug(&request_path)
1031 .filter(|_| publish_comment::comments_open(db.as_ref(), cfg.comments))
1032 {
1033 // The slug must name a post a reader can actually see. Without this the endpoint
1034 // writes a record under any name at all -- an unauthenticated write to storage
1035 // keyed on a string the sender chose, which is a way to fill a disk rather than a
1036 // way to comment. Measured against a live site before it was fixed: a POST to
1037 // /readme/anything/comment answered 303 and stored a comment.
1038 let posts = match publish::read(cfg.as_ref(), db.as_ref(), &id) {
1039 Ok(p) => p,
1040 Err(e) => {
1041 error!(e, "{}: publish: cannot read the posts to place a comment", id);
1042 Vec::new()
1043 }
1044 };
1045 if !posts.iter().any(|p| p.slug == slug) {
1046 info!("{}: publish: a comment named no post of ours ('{}')", id, slug);
1047 return Ok(Some(HttpMessage::respond_with_text(
1048 HttpStatus::NotFound, "No such post.")));
1049 }
1050 let mut edit_cookie: Option<(String, String)> = None;
1051 let said = match db.as_ref() {
1052 Some(dbh) => {
1053 let secret = res!(publish_comment::site_secret(dbh));
1054 let f = |k: &str| site_console::form_field(&body, k).unwrap_or_default();
1055 let sub = publish_comment::Submission {
1056 slug,
1057 parent: site_console::form_field(&body, "parent"),
1058 name: f("name"),
1059 email: site_console::form_field(&body, "email"),
1060 body: f("body"),
1061 honeypot: f("website"),
1062 challenge: f("challenge"),
1063 nonce: f("nonce"),
1064 from: Some(peer.ip().to_string()),
1065 now: publish_comment::now_stamp(),
1066 };
1067 // The site's AI, where it has set one up, so a stranger's first comment
1068 // can be judged rather than always made to wait. A settings read that
1069 // fails is no AI, not a failed comment: the comment falls back to the
1070 // queue, which is where it would have gone anyway.
1071 let ai_settings = publish_ai::get_settings(dbh).ok();
1072 let got = res!(publish_comment::receive(
1073 dbh,
1074 &publish_comment::Moderator::default(),
1075 ai_settings.as_ref(),
1076 &tls_client,
1077 (cfg.comment_rate_secs, cfg.comment_rate_hourly),
1078 sub,
1079 &secret,
1080 &secret,
1081 &id,
1082 ).await);
1083 // The token that lets its author correct it, handed back once, in a
1084 // cookie that expires with the window. It names one comment and proves
1085 // nothing else, so it is safe to hold in a browser.
1086 if let Some(cid) = &got.1 {
1087 edit_cookie = Some((
1088 cid.clone(),
1089 publish_comment::edit_token(cid, &secret),
1090 ));
1091 }
1092
1093 // A comment that ended up waiting for a person, and operators who asked to
1094 // be told: send each of them the alert. Spawned, not awaited, because it is
1095 // the operator's business and not the reader's -- the reader has posted and
1096 // should not wait on an SMTP round-trip to hear so. Best-effort: a mail
1097 // that will not send is logged, and the comment waits in the queue either
1098 // way, which is where the operator will find it regardless.
1099 if matches!(got.0, publish_comment::Received::Held) {
1100 if let (Some(settings), Some(mail_arc)) =
1101 (ai_settings.as_ref(), mail.as_ref())
1102 {
1103 if !settings.alert_emails.is_empty() {
1104 let from = cfg.newsletter_from(mail_arc);
1105 let site = cfg.site_name.clone();
1106 let slug_s = slug.to_string();
1107 let addrs = settings.alert_emails.clone();
1108 let mailer = mail_arc.clone();
1109 let idc = id.clone();
1110 tokio::spawn(async move {
1111 for to in &addrs {
1112 if let Err(e) = mailer.send_moderation_alert(
1113 &from, to, &site, &slug_s).await
1114 {
1115 warn!("{}: publish: a moderation alert to {} \
1116 could not be sent: {}", idc, to, e);
1117 }
1118 }
1119 });
1120 }
1121 }
1122 }
1123 got.0.tell_reader().to_string()
1124 }
1125 None => "shut".to_string(),
1126 };
1127 let mut resp = publish_page::comment_posted(cfg.as_ref(), slug, &said);
1128 if let Some((cid, token)) = edit_cookie {
1129 resp = publish_page::with_edit_cookie(resp, &cid, &token);
1130 }
1131 return Ok(Some(resp));
1132 }
1133 }
1134
1135 // The site console's writes, gated on a site admin's session and
1136 // guarded against cross-site forgery. It answers only the paths it
1137 // writes to and hands the rest back. Dispatched on the same terms as
1138 // its pages; the gate inside denies a site with no admins, since a
1139 // write needs a listed member and an empty list names none.
1140 if (publish.is_some() || !site_admins.is_empty())
1141 && site_console::owns(&request_path)
1142 {
1143 if let Some(resp) = res!(site_console::handle_post(
1144 site_admins.as_ref(),
1145 admin_state.as_deref(),
1146 publish.as_deref(),
1147 db.as_ref(),
1148 &tls_client,
1149 &mail,
1150 &request_path,
1151 &req_headers,
1152 &body,
1153 peer,
1154 &id,
1155 ).await) {
1156 return Ok(Some(resp));
1157 }
1158 }
1159
1160 // Dashboard `/admin/*` POST handlers (login form POST,
1161 // logout, future mutations). Same precedence rule as
1162 // GET: dashboard owns the subtree and dispatches before
1163 // any other lookup.
1164 if request_path == "/admin"
1165 || request_path.starts_with("/admin/")
1166 {
1167 if let Some(state) = &admin_state {
1168 let resp = res!(admin_handler::handle_post(
1169 state.as_ref(),
1170 &request_path,
1171 &body,
1172 &req_headers,
1173 peer,
1174 &id,
1175 ).await);
1176 return Ok(Some(resp));
1177 }
1178 return Ok(Some(HttpMessage::respond_with_text(
1179 HttpStatus::NotFound,
1180 "Not found.",
1181 )));
1182 }
1183
1184 // Check webhook routes first. Two dispatch branches
1185 // depending on the mode the route was configured in:
1186 //
1187 // - in-process `handler` -- the route names a registered
1188 // `WebhookHandler` in the webhook registry; dispatch
1189 // runs inside the Steel process as before.
1190 // - forwarded `upstream` -- the route carries an upstream
1191 // URL; the raw body and most incoming headers are
1192 // forwarded verbatim to the upstream so downstream
1193 // signature verification (e.g. `Stripe-Signature`)
1194 // still sees an unmodified payload.
1195 if let Some(wh) = webhook_routes.iter().find(|r| r.path == request_path) {
1196 if wh.is_upstream() {
1197 debug!("{}: POST {} -> webhook upstream {:?}:{:?}",
1198 id, request_path, wh.upstream_host, wh.upstream_port);
1199 let resp = res!(forward_webhook(
1200 wh,
1201 &body,
1202 &req_headers,
1203 &tls_client,
1204 &id,
1205 ).await);
1206 return Ok(Some(resp));
1207 }
1208 debug!("{}: POST {} -> webhook handler '{}'",
1209 id, request_path,
1210 wh.handler.as_deref().unwrap_or("?"));
1211 return webhook::dispatch(
1212 &webhook_registry, wh, &body, &req_headers, &tls_client, &id,
1213 ).await;
1214 }
1215
1216 // Find a matching API route.
1217 let route = match api_routes.iter().find(|r| r.path == request_path) {
1218 Some(r) => r,
1219 None => {
1220 debug!("{}: POST {} -- no matching API route.", id, request_path);
1221 return Ok(Some(HttpMessage::respond_with_text(
1222 HttpStatus::NotFound,
1223 "No API route matches this path.",
1224 )));
1225 }
1226 };
1227
1228 // In-process handler path: dispatch to the registered ApiHandler.
1229 if route.handler.is_some() {
1230 debug!("{}: POST {} -> api handler '{}'",
1231 id, request_path,
1232 route.handler.as_deref().unwrap_or("?"));
1233 let resp = res!(api::dispatch(
1234 &api_handler_registry,
1235 route,
1236 HttpMethod::POST,
1237 &loc,
1238 &body,
1239 &req_headers,
1240 &tls_client,
1241 &id,
1242 ).await);
1243 return Ok(Some(resp));
1244 }
1245
1246 // Proxy path: forward to the upstream via the shared helper.
1247 let resp = res!(forward_api_proxy(
1248 route,
1249 HttpMethod::POST,
1250 &loc,
1251 &body,
1252 &req_headers,
1253 &tls_client,
1254 &id,
1255 ).await);
1256 Ok(Some(resp))
1257 }
1258 }
1259}
1260
1261async fn forward_api_proxy(
1262 route: &ApiRoute,
1263 method: HttpMethod,
1264 loc: &HttpLocator,
1265 body: &[u8],
1266 req_headers: &HeaderFields,
1267 tls_client: &Option<Arc<ClientConfig>>,
1268 id: &str,
1269)
1270 -> Outcome<HttpMessage>
1271{
1272 let upstream_host = match &route.upstream_host {
1273 Some(h) => h.as_str(),
1274 None => return Err(err!(
1275 "{}: API route '{}' is in proxy mode but has no upstream_host.",
1276 id, route.path;
1277 Init, Missing, Bug)),
1278 };
1279 let upstream_port = route.upstream_port.unwrap_or(
1280 if route.upstream_tls { 443 } else { 80 });
1281 let upstream_path = route.upstream_path
1282 .as_deref().unwrap_or("/");
1283
1284 // Owned strings for headers built here so their `&str` refs live
1285 // until the outbound client has finished formatting the request.
1286 let mut owned: Vec<(String, String)> = Vec::new();
1287
1288 // Route-configured headers (secret tokens, fixed auth). These
1289 // win against any incoming client header with the same name
1290 // below -- the operator's declared headers are authoritative.
1291 for (name, value) in &route.headers {
1292 owned.push((name.clone(), value.clone()));
1293 }
1294
1295 // Propagate client headers that in-process handlers used to
1296 // have direct access to. The name list covers everything an
1297 // elearnity handler currently inspects plus the small set of
1298 // conventional "pass-through" headers browsers and CLIs send.
1299 // Hop-by-hop headers (`Host`, `Connection`, `Content-Length`)
1300 // are skipped because the outbound client regenerates them.
1301 let propagate: &[HeaderName] = &[
1302 HeaderName::Accept,
1303 HeaderName::AcceptLanguage,
1304 HeaderName::AcceptEncoding,
1305 HeaderName::ContentType,
1306 HeaderName::UserAgent,
1307 HeaderName::Authorization,
1308 HeaderName::Origin,
1309 HeaderName::Referer,
1310 ];
1311 for name in propagate {
1312 // Skip duplicates: a route-configured header with the same
1313 // name already sits in `owned`, so the operator's choice
1314 // wins over the client's.
1315 let name_str = fmt!("{}", name);
1316 if owned.iter().any(|(n, _)| n.eq_ignore_ascii_case(&name_str)) {
1317 continue;
1318 }
1319 if let Some(HeaderFieldValue::Generic(v)) = req_headers.get_one(name) {
1320 owned.push((name_str, v.clone()));
1321 }
1322 }
1323
1324 // Also let the dispatcher-resolved Content-Type override ride
1325 // through, because the POST branch already stamped it into
1326 // `loc.data` before reaching us and we want to respect the
1327 // original byte-level content-type.
1328 if let Some(ct) = loc.data.get(&dat!("content_type")) {
1329 if let Dat::Str(s) = ct {
1330 if !owned.iter().any(|(n, _)| n.eq_ignore_ascii_case("Content-Type")) {
1331 owned.push(("Content-Type".to_string(), s.clone()));
1332 }
1333 }
1334 }
1335
1336 let hdrs: Vec<(&str, &str)> = owned.iter()
1337 .map(|(n, v)| (n.as_str(), v.as_str()))
1338 .collect();
1339
1340 debug!("{}: {} {} -> {}{}:{}{} ({} headers)",
1341 id, method, route.path,
1342 if route.upstream_tls { "https://" } else { "http://" },
1343 upstream_host, upstream_port, upstream_path, hdrs.len());
1344
1345 if route.upstream_tls {
1346 let tls_cfg = match tls_client {
1347 Some(cfg) => cfg.clone(),
1348 None => return Err(err!(
1349 "{}: API route '{}' configured with https:// upstream but \
1350 no TLS client is available.", id, route.path;
1351 Init, Missing)),
1352 };
1353 https_request(
1354 upstream_host,
1355 upstream_port,
1356 method,
1357 upstream_path,
1358 &hdrs,
1359 body,
1360 tls_cfg,
1361 ).await
1362 } else {
1363 http_request(
1364 upstream_host,
1365 upstream_port,
1366 method,
1367 upstream_path,
1368 &hdrs,
1369 body,
1370 ).await
1371 }
1372}
1373
1374async fn forward_webhook(
1375 route: &WebhookRoute,
1376 body: &[u8],
1377 req_headers: &HeaderFields,
1378 tls_client: &Option<Arc<ClientConfig>>,
1379 id: &str,
1380)
1381 -> Outcome<HttpMessage>
1382{
1383 let upstream_host = match &route.upstream_host {
1384 Some(h) => h.as_str(),
1385 None => return Err(err!(
1386 "{}: forward_webhook called on a route with no upstream_host.", id;
1387 Bug, Missing)),
1388 };
1389 let upstream_port = route.upstream_port.unwrap_or(
1390 if route.upstream_tls { 443 } else { 80 });
1391 let upstream_path = route.upstream_path
1392 .as_deref().unwrap_or("/");
1393
1394 // Propagate the headers a typical webhook provider expects to see
1395 // verbatim. The list is deliberately short: Content-Type for the
1396 // JSON/form encoding, any Stripe-Signature / X-Hub-Signature style
1397 // header for signature verification, and a selection of common
1398 // auxiliaries (User-Agent, Request-Id, Idempotency-Key). Callers
1399 // that need a broader set can extend this list without touching
1400 // the dispatch shape.
1401 let mut hdrs: Vec<(String, String)> = Vec::new();
1402 let propagate_names: &[HeaderName] = &[
1403 HeaderName::ContentType,
1404 HeaderName::UserAgent,
1405 ];
1406 for name in propagate_names {
1407 if let Some(HeaderFieldValue::Generic(v)) =
1408 req_headers.get_one(name)
1409 {
1410 hdrs.push((fmt!("{}", name), v.clone()));
1411 }
1412 }
1413 // Propagate any non-standard header whose name begins with a
1414 // canonical signature prefix. This catches Stripe-Signature,
1415 // X-Hub-Signature, X-Signature, X-Hmac-Signature, and similar
1416 // without hardcoding the provider.
1417 for (name, values) in req_headers.iter() {
1418 if let HeaderName::NonStandard(n) = name {
1419 let lower = n.to_lowercase();
1420 let is_signature_like = lower.contains("signature")
1421 || lower.contains("idempotency")
1422 || lower == "x-request-id";
1423 if !is_signature_like {
1424 continue;
1425 }
1426 if let Some(HeaderFieldValue::Generic(v)) = values.first() {
1427 hdrs.push((n.clone(), v.clone()));
1428 }
1429 }
1430 }
1431 let hdr_refs: Vec<(&str, &str)> = hdrs.iter()
1432 .map(|(n, v)| (n.as_str(), v.as_str()))
1433 .collect();
1434
1435 debug!("{}: forwarding webhook to {}{}:{}{} (body {} bytes, {} headers)",
1436 id,
1437 if route.upstream_tls { "https://" } else { "http://" },
1438 upstream_host, upstream_port, upstream_path,
1439 body.len(), hdr_refs.len());
1440
1441 if route.upstream_tls {
1442 let tls_cfg = match tls_client {
1443 Some(cfg) => cfg.clone(),
1444 None => return Err(err!(
1445 "{}: webhook route '{}' configured with https:// upstream \
1446 but no TLS client is available.", id, route.path;
1447 Init, Missing)),
1448 };
1449 https_request(
1450 upstream_host,
1451 upstream_port,
1452 HttpMethod::POST,
1453 upstream_path,
1454 &hdr_refs,
1455 body,
1456 tls_cfg,
1457 ).await
1458 } else {
1459 http_request(
1460 upstream_host,
1461 upstream_port,
1462 HttpMethod::POST,
1463 upstream_path,
1464 &hdr_refs,
1465 body,
1466 ).await
1467 }
1468}