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
| 1 | use 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 | |
| 38 | use oxedyne_fe2o3_core::{ |
| 39 | prelude::*, |
| 40 | file::{ |
| 41 | OsPath, |
| 42 | }, |
| 43 | map::MapMut, |
| 44 | path::NormalPath, |
| 45 | rand::Rand, |
| 46 | }; |
| 47 | use oxedyne_fe2o3_iop_crypto::enc::Encrypter; |
| 48 | use oxedyne_fe2o3_iop_db::api::Database; |
| 49 | use oxedyne_fe2o3_iop_hash::api::Hasher; |
| 50 | use oxedyne_fe2o3_jdat::{ |
| 51 | prelude::*, |
| 52 | id::NumIdDat, |
| 53 | }; |
| 54 | use 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 | |
| 82 | use 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 | |
| 96 | use tokio::{ |
| 97 | self, |
| 98 | io::AsyncReadExt, |
| 99 | }; |
| 100 | use tokio_rustls::rustls::ClientConfig; |
| 101 | |
| 102 | |
| 103 | #[derive(Clone, Debug)] |
| 104 | pub 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 | |
| 126 | impl< |
| 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 | |
| 234 | impl< |
| 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 | |
| 1261 | async 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 | |
| 1374 | async 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 | } |