From e10d5da08fa2a42e355e5df39e8c0b2eb6539a1b Mon Sep 17 00:00:00 2001 From: Matt Kantor Date: Sat, 30 May 2026 06:59:06 -0400 Subject: [PATCH 1/3] Update a comment. --- src/content/handlebars_helpers/get.rs | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/src/content/handlebars_helpers/get.rs b/src/content/handlebars_helpers/get.rs index 0d18654..8b2ce72 100644 --- a/src/content/handlebars_helpers/get.rs +++ b/src/content/handlebars_helpers/get.rs @@ -176,11 +176,10 @@ where Ok(all_bytes) }, ); - // `FileBody`/`ProcessBody` rely on `actix_web::web::block`, which needs - // a Tokio runtime. When called from an HTTP handler one is already - // active, so `futures::executor::block_on` can drive the future to - // completion. From sync contexts (e.g. unit tests, the `eval`/`get` - // CLI subcommands) we spin up a fresh runtime to wait for the future. + // `block_on`, needs an active future runtime. When this is called from + // an HTTP handler one's already active, but in sync contexts (e.g. unit + // tests, the `eval`/`get` CLI subcommands) we spin up a fresh runtime + // to wait for the future. let bytes = match tokio::runtime::Handle::try_current() { Ok(_) => executor::block_on(bytes_future), Err(_) => actix_rt::System::new().block_on(bytes_future), From ed1763fbb45651dd6db7b2f5da23f99e4992e2bc Mon Sep 17 00:00:00 2001 From: Matt Kantor Date: Sat, 30 May 2026 08:04:26 -0400 Subject: [PATCH 2/3] Make `ByteStream` response bodies `Send`. I'm working on #5 which will involve moving rendering off of the main request-handling thread. This means `Media` needs to be `Send` so that it can cross thread boundaries. --- src/content/body.rs | 8 ++++---- src/content/mod.rs | 4 ++-- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/src/content/body.rs b/src/content/body.rs index fc5e2e8..d470961 100644 --- a/src/content/body.rs +++ b/src/content/body.rs @@ -6,7 +6,7 @@ use super::StreamError; use crate::bug_message; use bytes::Bytes; use futures::Stream; -use futures::future::{Future, FutureExt, LocalBoxFuture}; +use futures::future::{BoxFuture, Future, FutureExt}; use std::cmp; use std::fs::File; use std::io::ErrorKind::Interrupted; @@ -21,7 +21,7 @@ use actix_web::error::BlockingError; use actix_web::web; type ChunkResult = Result, BlockingError>; -type ChunkOperation<'a, T> = LocalBoxFuture<'a, ChunkResult>; +type ChunkOperation<'a, T> = BoxFuture<'a, ChunkResult>; fn flatten_chunk_result(result: ChunkResult) -> Result { match result { @@ -106,7 +106,7 @@ impl Stream for FileBody { file.by_ref().take(max_bytes).read_to_end(&mut buffer)?; Ok((file, Bytes::from(buffer))) }) - .boxed_local(), + .boxed(), ); self.poll_next(context) } @@ -200,7 +200,7 @@ impl Stream for ProcessBody { } } }) - .boxed_local(); + .boxed(); self.next = Some(next); self.poll_next(context) diff --git a/src/content/mod.rs b/src/content/mod.rs index 6645271..4231d25 100644 --- a/src/content/mod.rs +++ b/src/content/mod.rs @@ -29,12 +29,12 @@ pub use content_registry::{ContentRepresentations, RegisteredContent}; pub use route::Route; // This is just a trait alias to help make type signatures a bit saner. -pub trait ByteStream: Stream> +pub trait ByteStream: Stream> + Send where Self: Unpin, { } -impl ByteStream for T where T: Stream> + Unpin {} +impl ByteStream for T where T: Stream> + Unpin + Send {} /// A piece of rendered content along with its media type. pub struct Media { From 384da5c0902d56abc00de4277a26879380a36acc Mon Sep 17 00:00:00 2001 From: Matt Kantor Date: Sat, 30 May 2026 08:16:12 -0400 Subject: [PATCH 3/3] Render content off the request handler's thread. Fixes #5. --- src/http.rs | 418 ++++++++++++++++++++++++++++++++++++---------------- 1 file changed, 294 insertions(+), 124 deletions(-) diff --git a/src/http.rs b/src/http.rs index fd4bc5e..253cd84 100644 --- a/src/http.rs +++ b/src/http.rs @@ -127,11 +127,6 @@ where let path = request.uri().path(); - let content_engine = app_data - .shared_content_engine - .read() - .expect("RwLock for ContentEngine has been poisoned"); - let (route, media_range_from_url) = { let media_range_from_url = MimeGuess::from_path(path).first(); let path_without_extension = if media_range_from_url.is_some() { @@ -148,7 +143,7 @@ where return error_response( http::StatusCode::BAD_REQUEST, format!("HTTP request path `{path}` could not be parsed into a Route: {error}"), - &*content_engine, + &app_data.shared_content_engine, RequestData { route: None, query_parameters: HashMap::new(), @@ -157,7 +152,8 @@ where &app_data.error_handler_route, vec![&mime::TEXT_PLAIN], HeaderMap::new(), - ); + ) + .await; } Ok(request_route) => { if request_route.as_ref() == "/" { @@ -182,7 +178,7 @@ where return error_response( http::StatusCode::BAD_REQUEST, format!("Malformed query string `{query_string}`: {error}"), - &*content_engine, + &app_data.shared_content_engine, RequestData { route: Some(route), query_parameters: HashMap::new(), @@ -191,7 +187,8 @@ where &app_data.error_handler_route, vec![&mime::TEXT_PLAIN], HeaderMap::new(), - ); + ) + .await; } }; @@ -201,7 +198,7 @@ where return error_response( http::StatusCode::BAD_REQUEST, format!("Failed to handle request headers: {error}"), - &*content_engine, + &app_data.shared_content_engine, RequestData { route: Some(route), query_parameters, @@ -210,7 +207,8 @@ where &app_data.error_handler_route, vec![&mime::TEXT_PLAIN], HeaderMap::new(), - ); + ) + .await; } }; @@ -229,7 +227,7 @@ where request.headers().get(header::ACCEPT), error ), - &*content_engine, + &app_data.shared_content_engine, RequestData { route: Some(route), query_parameters, @@ -238,25 +236,52 @@ where &app_data.error_handler_route, vec![&mime::TEXT_PLAIN], HeaderMap::new(), - ); + ) + .await; } }, }; - let render_result = content_engine.get(&route).map(|content| { - let render_context = content_engine.render_context( - Some(route.clone()), - query_parameters.clone(), - request_headers.clone(), - ); - content.render(render_context, acceptable_media_ranges.clone()) - }); + // Rendering returns a lazy `ByteStream`, but the `render` call itself can + // do synchronous work (e.g. templates are rendered entirely synchronously + // and may even block on I/O via the `get` helper). We don't want to block + // the actix worker thread, so it's offloaded to a separate thread pool via + // `web::block`. + // + // TODO: Render concurrency is currently implicitly bounded by the size of + // tokio's blocking pool, and when it's full incoming requests will pile up. + // Instead this could be gated by an `Arc` with a + // bounded wait queue. When the queue is full and more requests come in, + // HTTP 503 responses could immediately be sent in order to shed load. + let render_result = { + let shared_content_engine = app_data.shared_content_engine.clone(); + let route = route.clone(); + let query_parameters = query_parameters.clone(); + let request_headers = request_headers.clone(); + let acceptable_media_ranges = acceptable_media_ranges + .iter() + .map(|&media_range| media_range.clone()) + .collect::>(); + web::block(move || { + let content_engine = shared_content_engine + .read() + .expect("RwLock for ContentEngine has been poisoned"); + content_engine.get(&route).map(|content| { + let render_context = + content_engine.render_context(Some(route), query_parameters, request_headers); + content.render(render_context, acceptable_media_ranges.iter()) + }) + }) + .await + }; + // The outer `Result` models errors from `web::block`, while the inner one + // models render errors. match render_result { - Some(Ok(Media { + Ok(Some(Ok(Media { content, media_type, - })) => { + }))) => { log::info!( "Responding with {}, body from {} as {}", http::StatusCode::OK, @@ -292,45 +317,70 @@ where }), ) } - Some(Err(error @ RenderError::CannotProvideAcceptableMediaType)) => error_response( - http::StatusCode::NOT_ACCEPTABLE, - format!("Cannot provide an acceptable response: {error}"), - &*content_engine, - RequestData { - route: Some(route), - query_parameters, - request_headers, - }, - &app_data.error_handler_route, - acceptable_media_ranges, - HeaderMap::new(), - ), - Some(Err(error)) => error_response( - http::StatusCode::INTERNAL_SERVER_ERROR, - format!("Failed to render content: {error}"), - &*content_engine, - RequestData { - route: Some(route), - query_parameters, - request_headers, - }, - &app_data.error_handler_route, - acceptable_media_ranges, - HeaderMap::new(), - ), - None => error_response( - http::StatusCode::NOT_FOUND, - "No content found at route", - &*content_engine, - RequestData { - route: Some(route), - query_parameters, - request_headers, - }, - &app_data.error_handler_route, - acceptable_media_ranges, - HeaderMap::new(), - ), + Ok(Some(Err(error @ RenderError::CannotProvideAcceptableMediaType))) => { + error_response( + http::StatusCode::NOT_ACCEPTABLE, + format!("Cannot provide an acceptable response: {error}"), + &app_data.shared_content_engine, + RequestData { + route: Some(route), + query_parameters, + request_headers, + }, + &app_data.error_handler_route, + acceptable_media_ranges, + HeaderMap::new(), + ) + .await + } + Ok(Some(Err(error))) => { + error_response( + http::StatusCode::INTERNAL_SERVER_ERROR, + format!("Failed to render content: {error}"), + &app_data.shared_content_engine, + RequestData { + route: Some(route), + query_parameters, + request_headers, + }, + &app_data.error_handler_route, + acceptable_media_ranges, + HeaderMap::new(), + ) + .await + } + Ok(None) => { + error_response( + http::StatusCode::NOT_FOUND, + "No content found at route", + &app_data.shared_content_engine, + RequestData { + route: Some(route), + query_parameters, + request_headers, + }, + &app_data.error_handler_route, + acceptable_media_ranges, + HeaderMap::new(), + ) + .await + } + Err(blocking_error) => { + error_response( + http::StatusCode::INTERNAL_SERVER_ERROR, + format!("Failed to render content: {blocking_error}"), + &app_data.shared_content_engine, + RequestData { + route: Some(route), + query_parameters, + request_headers, + }, + &app_data.error_handler_route, + acceptable_media_ranges, + HeaderMap::new(), + ) + .await + } } } @@ -357,11 +407,6 @@ where .app_data::>() .expect("App data was not of the expected type!"); - let content_engine = app_data - .shared_content_engine - .read() - .expect("RwLock for ContentEngine has been poisoned"); - let mut response_headers = HeaderMap::with_capacity(1); response_headers.insert( http::header::ALLOW, @@ -371,7 +416,7 @@ where error_response( http::StatusCode::METHOD_NOT_ALLOWED, format!("The {} request method is not supported", request.method()), - &*content_engine, + &app_data.shared_content_engine, RequestData { route: None, query_parameters: HashMap::new(), @@ -381,6 +426,7 @@ where vec![&mime::TEXT_PLAIN], response_headers, ) + .await } fn log_request(request: &HttpRequest) { @@ -406,10 +452,10 @@ fn log_request(request: &HttpRequest) { ); } -fn error_response( +async fn error_response( status_code: http::StatusCode, details: Details, - content_engine: &Engine, + shared_content_engine: &Arc>, request_data: RequestData, error_handler_route: &Option, acceptable_media_ranges: Vec<&MediaRange>, @@ -436,61 +482,83 @@ where response_builder.insert_header((header_name.clone(), header_value.clone())); } - error_handler_route - .as_ref() - .and_then(|route| { - content_engine.get(route).and_then(|content| { - let error_context = content_engine - .render_context( - request_data.route.clone(), - request_data.query_parameters, - request_data.request_headers, - ) - .into_error_context(status_code.as_u16()); - match content.render(error_context, acceptable_media_ranges) { - Ok(rendered_content) => Some((route, rendered_content)), - Err(rendering_error) => { - log::error!( - "Error occurred while rendering error handler: {rendering_error}" - ); - None - } - } + // Render the error handler (if one is configured) off the request-handler + // thread, for the same reason as the main content render in `get`: + // rendering is synchronous and may block the request-handling worker. + let rendered_error_handler = match error_handler_route { + None => None, + Some(error_handler_route) => { + let shared_content_engine = shared_content_engine.clone(); + let request_route = request_data.route.clone(); + let error_handler_route = error_handler_route.clone(); + let status_code_value = status_code.as_u16(); + let acceptable_media_ranges = acceptable_media_ranges + .iter() + .map(|&media_range| media_range.clone()) + .collect::>(); + web::block(move || { + let content_engine = shared_content_engine + .read() + .expect("RwLock for ContentEngine has been poisoned"); + content_engine + .get(&error_handler_route) + .and_then(|content| { + let error_context = content_engine + .render_context(request_route, request_data.query_parameters, request_data.request_headers) + .into_error_context(status_code_value); + match content.render(error_context, acceptable_media_ranges.iter()) { + Ok(rendered_content) => { + Some((error_handler_route, rendered_content)) + } + Err(rendering_error) => { + log::error!( + "Error occurred while rendering error handler: {rendering_error}" + ); + None + } + } + }) }) - }) - .map( - |( - error_handler_route, - Media { - media_type, - content, - }, - )| { - match request_data.route.clone() { - Some(request_route) => log::warn!( - "Responding with {} for {}, body from {} as {}: {}", - status_code, - request_route, - error_handler_route, - media_type, - details.as_ref() - ), - None => log::warn!( - "Responding with {}, body from {} as {}: {}", - status_code, - error_handler_route, - media_type, - details.as_ref() - ), - }; - response_builder - .content_type(media_type.to_string()) - .streaming(content.inspect_err(|error| { - log::error!("An error occurred while streaming a response body: {error}",); - })) + .await + .unwrap_or_else(|blocking_error| { + log::error!("Error handler render task failed: {blocking_error}"); + None + }) + } + }; + + match rendered_error_handler { + Some(( + error_handler_route, + Media { + media_type, + content, }, - ) - .unwrap_or_else(|| { + )) => { + match request_data.route.clone() { + Some(request_route) => log::warn!( + "Responding with {} for {}, body from {} as {}: {}", + status_code, + request_route, + error_handler_route, + media_type, + details.as_ref() + ), + None => log::warn!( + "Responding with {}, body from {} as {}: {}", + status_code, + error_handler_route, + media_type, + details.as_ref() + ), + }; + response_builder + .content_type(media_type.to_string()) + .streaming(content.inspect_err(|error| { + log::error!("An error occurred while streaming a response body: {error}",); + })) + } + None => { // Send a default error response if the error handler failed or was // not specified. let media_type = "text/plain"; @@ -514,7 +582,8 @@ where .canonical_reason() .unwrap_or("Something Went Wrong"), ) - }) + } + } } fn acceptable_media_ranges_from_accept_header<'a>( @@ -570,9 +639,13 @@ mod tests { use actix_web::http::StatusCode; use actix_web::http::header::{HeaderName, HeaderValue}; use actix_web::test::TestRequest; + use awc::Client as HttpClient; use bytes::Bytes; use maplit::hashmap; + use std::fs; + use std::os::unix::fs::PermissionsExt; use std::path::Path; + use std::time::{Duration, Instant}; use test_log::test; type TestContentEngine<'a> = FilesystemBasedContentEngine<'a, ServerInfo>; @@ -1400,4 +1473,101 @@ mod tests { assert_eq!(response.status(), StatusCode::BAD_REQUEST); } + + #[actix_web::test] + async fn slow_rendering_does_not_block_the_request_handler_thread() { + // Create a content directory supporting three routes: + // 1. /fast (a trivial static file) + // 2. /slow (an executable that takes a while to terminate) + // 3. /blocking (a template which `get`s the slow executable) + let temp_dir = tempfile::tempdir().expect("Failed to create temp dir"); + let temp_dir_path = temp_dir.path(); + { + fs::write(temp_dir_path.join("fast.txt"), "i am fast") + .expect("Failed to write fast.txt"); + fs::write(temp_dir_path.join("blocking.txt.hbs"), "{{get \"/slow\"}}") + .expect("Failed to write blocking.txt.hbs"); + let slow_executable = temp_dir_path.join("slow.txt.sh"); + fs::write(&slow_executable, "#!/bin/sh\nsleep 3\necho i am slow\n") + .expect("Failed to write slow.txt.sh"); + fs::set_permissions(&slow_executable, fs::Permissions::from_mode(0o755)) + .expect("Failed to set permissions on slow.txt.sh"); + }; + + let content_directory = + ContentDirectory::from_root(&temp_dir_path).expect("Failed to read content directory"); + let shared_content_engine = FilesystemBasedContentEngine::from_content_directory( + content_directory, + ServerInfo { + version: ServerVersion(""), + operator_path: PathBuf::new(), + socket_address: None, + }, + ) + .expect("Content engine could not be created"); + + // Start the server with a single worker (so a blocked request handler + // thread stalls all other requests). + let server = HttpServer::new(move || { + App::new() + .app_data(AppData { + shared_content_engine: shared_content_engine.clone(), + index_route: None, + error_handler_route: None, + }) + .default_service(web::to(dispatch::)) + }) + .keep_alive(KeepAlive::Disabled) + .workers(1) // Just one worker! + .bind(("127.0.0.1", 0)) + .expect("Failed to bind test server"); + let address = server + .addrs() + .into_iter() + .next() + .expect("Test server didn't bind to any addresses"); + let server = server.run(); + actix_rt::spawn(server); + + // Issue HTTP GET and report how long it took alongside the response. + async fn timed_get(client: &HttpClient, url: String) -> (Duration, StatusCode, Bytes) { + let started = Instant::now(); + let mut response = client + .get(url) + .insert_header((header::ACCEPT, "text/plain")) + .send() + .await + .expect("Failed to send request"); + let body = response.body().await.expect("Failed to read response body"); + (started.elapsed(), response.status(), body) + } + + let client = HttpClient::new(); + + // Request /blocking and give it time to start rendering. + let slow_request = actix_rt::spawn({ + let client = client.clone(); + let url = format!("http://{address}/blocking"); + async move { timed_get(&client, url).await } + }); + actix_rt::time::sleep(Duration::from_millis(500)).await; + + let (fast_elapsed, fast_status, fast_body) = + timed_get(&client, format!("http://{address}/fast")).await; + + assert_eq!(fast_status, StatusCode::OK, "Fast request did not succeed"); + assert_eq!(fast_body, "i am fast", "Fast response body was incorrect"); + assert!( + fast_elapsed < Duration::from_millis(1500), + "Fast request took {} seconds, which means the worker thread was blocked by \ + the in-flight slow render", + (fast_elapsed.as_millis() as f64) / 1000f64 + ); + + let (_slow_elapsed, slow_status, slow_body) = + slow_request.await.expect("Slow request task panicked"); + + assert_eq!(slow_status, StatusCode::OK, "Slow request did not succeed"); + assert_eq!(slow_body, "i am slow\n", "Slow response body was incorrect"); + } }