From 45a46f681830253f0f676c49e2a96895ab17af81 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 21:25:21 -0300 Subject: [PATCH 1/4] fix(cloudfront): hold a function to a CPU-time budget, not wall-clock time TestFunction and TestConnectionFunction ran the handler on a worker thread and failed it with "function execution exceeded the 250ms time limit" when the thread had not answered within 250ms of wall-clock time, counting thread spawn, runtime setup and any time the host spent not running the thread. On a loaded CI runner a trivial echo handler hit that limit: the cloudfront_test_function e2e tests needed a retry in 4 of the last 26 E2E runs, and nextest's retried-failure output (#2520) captured the error on test_function_echoes_request. Real CloudFront Functions are bounded by compute (~1ms of CPU per request), not elapsed time. The worker now measures its own CPU time (CLOCK_THREAD_CPUTIME_ID; elapsed time where the platform has no thread CPU clock) and a run that used more than 250ms of it fails with the same error, keeping whatever it logged. ComputeUtilization is derived from that CPU time too. The caller's recv_timeout becomes a 10s wall-clock safety net for a run that never returns. boa's loop and recursion caps still stop `while(1){}`. --- Cargo.lock | 1 + crates/fakecloud-cloudfront/Cargo.toml | 3 + crates/fakecloud-cloudfront/src/js_runtime.rs | 228 ++++++++++++++---- website/content/docs/services/cloudfront.md | 2 +- 4 files changed, 185 insertions(+), 49 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7c713629b..ab5263b5c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4958,6 +4958,7 @@ dependencies = [ "http-body-util", "hyper 1.9.0", "hyper-util", + "libc", "parking_lot", "quick-xml", "reqwest", diff --git a/crates/fakecloud-cloudfront/Cargo.toml b/crates/fakecloud-cloudfront/Cargo.toml index a7cf673c5..1f7de057f 100644 --- a/crates/fakecloud-cloudfront/Cargo.toml +++ b/crates/fakecloud-cloudfront/Cargo.toml @@ -30,3 +30,6 @@ thiserror = { workspace = true } tokio = { workspace = true } tracing = { workspace = true } uuid = { workspace = true } + +[target.'cfg(unix)'.dependencies] +libc = "0.2" diff --git a/crates/fakecloud-cloudfront/src/js_runtime.rs b/crates/fakecloud-cloudfront/src/js_runtime.rs index b7c7a08bc..c4ee1b194 100644 --- a/crates/fakecloud-cloudfront/src/js_runtime.rs +++ b/crates/fakecloud-cloudfront/src/js_runtime.rs @@ -9,21 +9,26 @@ //! - JSON.stringify the return value; //! - capture `console.log/error` output as execution log lines. //! -//! Limits are enforced in two layers: +//! Limits are enforced in three layers: //! //! 1. boa's loop iteration + recursion caps trip on hot loops so the //! interpreter eventually returns control even under adversarial //! user JS. -//! 2. A wall-clock timeout: the actual execution runs on a dedicated -//! OS thread; the calling thread waits on a `mpsc::sync_channel` via -//! `recv_timeout`. If the JS doesn't finish in time we abandon the -//! worker thread (best-effort — boa's iteration limit will eventually -//! let it die) and return a timeout error. +//! 2. A compute budget: real CloudFront Functions are bounded by the CPU +//! a request consumes (~1ms, and 2MB of memory), not by elapsed time. +//! The execution runs on a dedicated OS thread that measures its own +//! CPU time, and a run that used more than [`EXECUTION_TIMEOUT`] of it +//! fails with the time-limit error. Measuring CPU rather than wall +//! time means a host that descheduled the thread -- a loaded CI runner +//! -- does not fail a handler that did almost no work. +//! 3. A wall-clock safety net: the calling thread waits on a +//! `mpsc::sync_channel` via `recv_timeout` for [`WALL_CLOCK_LIMIT`]. If +//! the JS still hasn't finished we abandon the worker thread +//! (best-effort -- boa's iteration limit will eventually let it die) +//! and return the same time-limit error. //! -//! Real CloudFront Functions are bounded at ~1ms of CPU per request and -//! 2MB of memory. We mirror that with a 250ms wall-clock budget — -//! looser than AWS to absorb cold-start jitter on shared CI runners, -//! still tight enough that `while(1){}` is killed in tests. +//! The 250ms budget is looser than AWS's so ordinary handlers never trip +//! it in a debug build, while `while(1){}` is still stopped in tests. use std::cell::RefCell; use std::rc::Rc; @@ -34,13 +39,18 @@ use boa_engine::object::ObjectInitializer; use boa_engine::property::Attribute; use boa_engine::{js_string, Context, JsValue, NativeFunction, Source}; -/// Hard wall-clock cap on a single TestFunction / TestConnectionFunction -/// invocation. AWS bounds production traffic at ~1ms of CPU; we set -/// 250ms so CI runners with noisy neighbours don't false-alarm on -/// well-formed handlers, while `while(1){}` is still killed well -/// inside any reasonable test timeout. +/// Compute budget for a single TestFunction / TestConnectionFunction +/// invocation, measured as CPU time on the executing thread. AWS bounds +/// production traffic at ~1ms of CPU; 250ms leaves ordinary handlers far +/// below it even unoptimized, while a CPU-bound handler still exceeds it. pub(crate) const EXECUTION_TIMEOUT: Duration = Duration::from_millis(250); +/// How long the caller waits for the worker thread before abandoning it. +/// A safety net for a run that never returns; the compute budget above is +/// the limit a handler is actually held to, so this is generous enough +/// that a descheduled thread on a loaded host still reports back. +const WALL_CLOCK_LIMIT: Duration = Duration::from_secs(10); + /// boa loop iteration cap. Tight enough that `while(1){}` exits the VM /// well within the wall-clock budget on any reasonable host, loose /// enough that small loops in real handlers run to completion. @@ -69,8 +79,18 @@ pub(crate) struct JsExecution { } /// Run `handler(event)` defined in `code` against `event_json` on a -/// dedicated worker thread, enforcing `EXECUTION_TIMEOUT`. +/// dedicated worker thread, holding it to the `EXECUTION_TIMEOUT` compute +/// budget. pub(crate) fn run_handler(code: &str, event_json: &[u8]) -> JsExecution { + run_handler_with_limits(code, event_json, EXECUTION_TIMEOUT, WALL_CLOCK_LIMIT) +} + +fn run_handler_with_limits( + code: &str, + event_json: &[u8], + compute_budget: Duration, + wall_limit: Duration, +) -> JsExecution { let code = code.to_owned(); let event = event_json.to_vec(); let (tx, rx) = mpsc::sync_channel::(1); @@ -79,20 +99,23 @@ pub(crate) fn run_handler(code: &str, event_json: &[u8]) -> JsExecution { // `Rc`s and is `!Send`. We can't pre-spawn a worker pool without // marshalling the script + event via channels anyway, so a fresh // thread per call is the simpler shape. - let started = Instant::now(); let _ = std::thread::Builder::new() .name("cloudfront-js".to_string()) // Boa's bytecode VM is recursive so we want a generous stack. .stack_size(8 * 1024 * 1024) .spawn(move || { - let result = run_handler_blocking(&code, &event, started); + let clock = ComputeClock::start(); + let mut result = run_handler_blocking(&code, &event, &clock); + if clock.elapsed() > compute_budget { + result = time_limit_exceeded(result.logs); + } // If the receiver has timed out and gone away the send // simply errors; we don't care — the worker is being // abandoned. let _ = tx.send(result); }); - match rx.recv_timeout(EXECUTION_TIMEOUT) { + match rx.recv_timeout(wall_limit) { Ok(mut exec) => { // Floor compute_utilization at 1% on success so callers // don't mistake a successful run for an unrun one. @@ -101,18 +124,7 @@ pub(crate) fn run_handler(code: &str, event_json: &[u8]) -> JsExecution { } exec } - Err(mpsc::RecvTimeoutError::Timeout) => { - let msg = format!( - "function execution exceeded the {}ms time limit", - EXECUTION_TIMEOUT.as_millis() - ); - JsExecution { - output: None, - error: Some(msg.clone()), - logs: vec![format!("ERROR: {msg}")], - compute_utilization: 101, - } - } + Err(mpsc::RecvTimeoutError::Timeout) => time_limit_exceeded(Vec::new()), Err(mpsc::RecvTimeoutError::Disconnected) => { // The worker thread either failed to spawn (`spawn()` errored // and the sender was dropped) or panicked partway through. @@ -129,7 +141,65 @@ pub(crate) fn run_handler(code: &str, event_json: &[u8]) -> JsExecution { } } -fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsExecution { +/// The error a run over its compute budget (or past the wall-clock safety +/// net) reports, keeping whatever it logged before the limit. +fn time_limit_exceeded(mut logs: Vec) -> JsExecution { + let msg = format!( + "function execution exceeded the {}ms time limit", + EXECUTION_TIMEOUT.as_millis() + ); + logs.push(format!("ERROR: {msg}")); + JsExecution { + output: None, + error: Some(msg), + logs, + compute_utilization: 101, + } +} + +/// Measures the compute a run consumed: CPU time on the current thread +/// where the platform reports it, elapsed time otherwise. +struct ComputeClock { + cpu_start: Option, + wall_start: Instant, +} + +impl ComputeClock { + fn start() -> Self { + Self { + cpu_start: thread_cpu_time(), + wall_start: Instant::now(), + } + } + + fn elapsed(&self) -> Duration { + match (self.cpu_start, thread_cpu_time()) { + (Some(start), Some(now)) => now.saturating_sub(start), + _ => self.wall_start.elapsed(), + } + } +} + +/// CPU time consumed so far by the calling thread. +#[cfg(unix)] +fn thread_cpu_time() -> Option { + let mut ts = libc::timespec { + tv_sec: 0, + tv_nsec: 0, + }; + // SAFETY: `ts` is a valid, writable timespec; CLOCK_THREAD_CPUTIME_ID + // is supported on Linux and macOS, and a failure is reported by the + // return code rather than by leaving `ts` partially written. + let rc = unsafe { libc::clock_gettime(libc::CLOCK_THREAD_CPUTIME_ID, &mut ts) }; + (rc == 0).then(|| Duration::new(ts.tv_sec as u64, ts.tv_nsec as u32)) +} + +#[cfg(not(unix))] +fn thread_cpu_time() -> Option { + None +} + +fn run_handler_blocking(code: &str, event_json: &[u8], clock: &ComputeClock) -> JsExecution { let logs: Rc>> = Rc::new(RefCell::new(Vec::new())); let mut ctx = Context::default(); ctx.runtime_limits_mut() @@ -138,17 +208,17 @@ fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsEx .set_recursion_limit(RECURSION_LIMIT); if let Err(err) = install_console(&mut ctx, &logs) { - return error_execution(format!("failed to install console: {err}"), &logs, started); + return error_execution(format!("failed to install console: {err}"), &logs, clock); } if let Err(err) = ctx.eval(Source::from_bytes(code.as_bytes())) { - return error_execution(format!("{}", err), &logs, started); + return error_execution(format!("{}", err), &logs, clock); } let event_str = match std::str::from_utf8(event_json) { Ok(s) => s, Err(_) => { - return error_execution("EventObject is not valid UTF-8".to_string(), &logs, started); + return error_execution("EventObject is not valid UTF-8".to_string(), &logs, clock); } }; // Wrap in parens so a top-level `{ ... }` object literal parses as @@ -157,7 +227,7 @@ fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsEx let event = match ctx.eval(Source::from_bytes(event_src.as_bytes())) { Ok(v) => v, Err(err) => { - return error_execution(format!("invalid EventObject JSON: {err}"), &logs, started); + return error_execution(format!("invalid EventObject JSON: {err}"), &logs, clock); } }; @@ -167,22 +237,18 @@ fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsEx return error_execution( format!("function handler is not defined: {err}"), &logs, - started, + clock, ); } }; let Some(handler_fn) = handler.as_callable() else { - return error_execution( - "function handler is not callable".to_string(), - &logs, - started, - ); + return error_execution("function handler is not callable".to_string(), &logs, clock); }; let returned = match handler_fn.call(&JsValue::undefined(), &[event], &mut ctx) { Ok(v) => v, Err(err) => { - return error_execution(format!("{}", err), &logs, started); + return error_execution(format!("{}", err), &logs, clock); } }; @@ -192,7 +258,7 @@ fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsEx return error_execution( format!("failed to JSON.stringify result: {err}"), &logs, - started, + clock, ); } }; @@ -201,7 +267,7 @@ fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsEx return error_execution( format!("function output exceeded {MAX_OUTPUT_BYTES} bytes"), &logs, - started, + clock, ); } @@ -210,16 +276,20 @@ fn run_handler_blocking(code: &str, event_json: &[u8], started: Instant) -> JsEx output: Some(stringified), error: None, logs: captured, - compute_utilization: utilization_pct(started.elapsed()), + compute_utilization: utilization_pct(clock.elapsed()), } } -fn error_execution(msg: String, logs: &Rc>>, started: Instant) -> JsExecution { +fn error_execution( + msg: String, + logs: &Rc>>, + clock: &ComputeClock, +) -> JsExecution { let mut captured = logs.borrow().clone(); captured.push(format!("ERROR: {msg}")); // Saturate past 100 on any failure so the metric alone signals the // run did not complete cleanly, regardless of how fast it failed. - let elapsed_pct = utilization_pct(started.elapsed()); + let elapsed_pct = utilization_pct(clock.elapsed()); let pct = elapsed_pct.max(101); JsExecution { output: None, @@ -380,6 +450,68 @@ mod tests { assert!(exec.error.is_some()); } + #[test] + fn a_run_over_its_compute_budget_reports_the_time_limit() { + let exec = run_handler_with_limits( + r#"function handler(e) { console.log("ran"); return e; }"#, + b"{}", + Duration::ZERO, + WALL_CLOCK_LIMIT, + ); + assert!(exec.output.is_none(), "got output {:?}", exec.output); + let err = exec.error.expect("error"); + assert!(err.contains("250ms time limit"), "got {err}"); + assert!(exec.compute_utilization > 100); + assert!( + exec.logs.first().is_some_and(|l| l == "ran"), + "logs before the limit are kept: {:?}", + exec.logs + ); + } + + #[test] + fn a_worker_that_does_not_report_back_hits_the_wall_clock_safety_net() { + let exec = run_handler_with_limits( + r#"function handler(e) { return e; }"#, + b"{}", + EXECUTION_TIMEOUT, + Duration::ZERO, + ); + let err = exec.error.expect("error"); + assert!(err.contains("time limit"), "got {err}"); + assert!(exec.compute_utilization > 100); + } + + #[cfg(unix)] + #[test] + fn a_descheduled_thread_accrues_no_compute() { + // The budget is CPU time: a thread that is not running -- here + // sleeping, on a loaded host descheduled -- does not use it up. + let clock = ComputeClock::start(); + std::thread::sleep(EXECUTION_TIMEOUT + Duration::from_millis(100)); + assert!( + clock.elapsed() < EXECUTION_TIMEOUT, + "sleeping consumed {:?} of compute", + clock.elapsed() + ); + } + + #[cfg(unix)] + #[test] + fn busy_work_accrues_compute() { + let clock = ComputeClock::start(); + let wall = Instant::now(); + let mut x: u64 = 0; + while wall.elapsed() < Duration::from_millis(50) { + x = std::hint::black_box(x.wrapping_add(1)); + } + assert!( + clock.elapsed() >= Duration::from_millis(20), + "50ms of busy work measured {:?}", + clock.elapsed() + ); + } + #[test] fn infinite_loop_is_killed_by_timeout() { let exec = run_handler(r#"function handler() { while(1){} }"#, b"{}"); diff --git a/website/content/docs/services/cloudfront.md b/website/content/docs/services/cloudfront.md index 8bddca602..59884b34e 100644 --- a/website/content/docs/services/cloudfront.md +++ b/website/content/docs/services/cloudfront.md @@ -20,7 +20,7 @@ fakecloud implements CloudFront's REST-XML control plane focused on the operatio - **Origin Request Policy** — `CreateOriginRequestPolicy`, `GetOriginRequestPolicy`, `GetOriginRequestPolicyConfig`, `UpdateOriginRequestPolicy`, `DeleteOriginRequestPolicy`, `ListOriginRequestPolicies`. Managed `Managed-CORS-S3Origin`, `Managed-CORS-CustomOrigin`, `Managed-AllViewer`, `Managed-UserAgentRefererHeaders`, `Managed-AllViewerExceptHostHeader` pre-seeded. - **Response Headers Policy** — `CreateResponseHeadersPolicy`, `GetResponseHeadersPolicy`, `GetResponseHeadersPolicyConfig`, `UpdateResponseHeadersPolicy`, `DeleteResponseHeadersPolicy`, `ListResponseHeadersPolicies`. Managed `Managed-SimpleCORS`, `Managed-CORS-With-Preflight`, `Managed-SecurityHeadersPolicy` pre-seeded. - **Continuous Deployment Policy** — `CreateContinuousDeploymentPolicy`, `GetContinuousDeploymentPolicy`, `GetContinuousDeploymentPolicyConfig`, `UpdateContinuousDeploymentPolicy`, `DeleteContinuousDeploymentPolicy`, `ListContinuousDeploymentPolicies`. -- **CloudFront Functions** — `CreateFunction`, `DescribeFunction`, `GetFunction` (returns raw source bytes), `UpdateFunction`, `DeleteFunction`, `ListFunctions`, `PublishFunction` (DEVELOPMENT -> LIVE), `TestFunction`. `TestFunction` executes the user's `function handler(event) { ... }` against the supplied `EventObject` via an embedded `boa_engine` JavaScript runtime: the handler return value is JSON-encoded into ``, thrown errors land in ``, and `console.log` / `console.error` lines are captured into ``. `` reports a 0..=100 share of the per-request CPU budget on success and saturates past 100 on failure. `Stage` (DEVELOPMENT | LIVE) picks which version's source to run: DEVELOPMENT is the latest CreateFunction / UpdateFunction body, LIVE is the snapshot frozen at the most recent PublishFunction call (so post-publish updates don't leak into the published behaviour). Execution is capped by boa's loop-iteration and recursion limits plus a 250ms wall-clock guard (looser than AWS's ~1ms production budget so CI runners with noisy neighbours don't false-alarm; tests rely on the cap to kill `while(1){}`). +- **CloudFront Functions** — `CreateFunction`, `DescribeFunction`, `GetFunction` (returns raw source bytes), `UpdateFunction`, `DeleteFunction`, `ListFunctions`, `PublishFunction` (DEVELOPMENT -> LIVE), `TestFunction`. `TestFunction` executes the user's `function handler(event) { ... }` against the supplied `EventObject` via an embedded `boa_engine` JavaScript runtime: the handler return value is JSON-encoded into ``, thrown errors land in ``, and `console.log` / `console.error` lines are captured into ``. `` reports a 0..=100 share of the per-request CPU budget on success and saturates past 100 on failure. `Stage` (DEVELOPMENT | LIVE) picks which version's source to run: DEVELOPMENT is the latest CreateFunction / UpdateFunction body, LIVE is the snapshot frozen at the most recent PublishFunction call (so post-publish updates don't leak into the published behaviour). Execution is capped by boa's loop-iteration and recursion limits plus a 250ms compute budget, measured as CPU time on the executing thread like AWS's ~1ms production budget (looser, so ordinary handlers never trip it unoptimized, while a CPU-bound handler still fails with the time-limit error); because it counts CPU rather than elapsed time, a loaded host that stalls the thread does not fail a handler that did little work, and a 10s wall-clock safety net abandons a run that never returns. - **Connection Functions** — `CreateConnectionFunction`, `DescribeConnectionFunction`, `GetConnectionFunction` (raw source), `UpdateConnectionFunction`, `DeleteConnectionFunction`, `ListConnectionFunctions`, `PublishConnectionFunction`, `TestConnectionFunction`. `TestConnectionFunction` runs the same `boa_engine` runtime against the supplied `ConnectionObject` and emits the parallel `` / `` / `` / `` shape, with the same DEVELOPMENT-vs-LIVE stage selection as `TestFunction`. - **Public Keys + Key Groups** — full CRUD with `CallerReference` immutability on update. - **Key Value Stores** — `CreateKeyValueStore` (with optional `ImportSource`), `DescribeKeyValueStore`, `UpdateKeyValueStore`, `DeleteKeyValueStore`, `ListKeyValueStores`. From d41f73456d45ee9ace793bbfbfcae303f0fb8b56 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 21:32:30 -0300 Subject: [PATCH 2/4] fix(cloudfront): stop a function at its CPU budget while it runs The budget was only checked after the handler returned, so a CPU-bound handler that boa's loop and recursion caps cannot stop (catastrophic regex backtracking inside one builtin call) held the caller until the wall-clock safety net instead of being cut off at 250ms. The worker now hands the caller a handle to its CPU clock as it starts (pthread_getcpuclockid on Linux, the thread's Mach port on macOS), and the caller polls it every 5ms, abandoning the run with the time-limit error as soon as it passes the budget. Where no thread CPU clock exists, elapsed time stands in. The worker's own final check stays for a run that crosses the budget between polls; the safety net drops to 5s. Tests: a backtracking regex that runs for seconds is cut off in ~0.26s (verified to fail, waiting the full 5s, with the watchdog check disabled), and the busy-work clock test is bounded by work done rather than elapsed time so a loaded runner cannot flake it. --- crates/fakecloud-cloudfront/src/js_runtime.rs | 251 ++++++++++++++---- website/content/docs/services/cloudfront.md | 2 +- 2 files changed, 196 insertions(+), 57 deletions(-) diff --git a/crates/fakecloud-cloudfront/src/js_runtime.rs b/crates/fakecloud-cloudfront/src/js_runtime.rs index c4ee1b194..80844527c 100644 --- a/crates/fakecloud-cloudfront/src/js_runtime.rs +++ b/crates/fakecloud-cloudfront/src/js_runtime.rs @@ -49,7 +49,10 @@ pub(crate) const EXECUTION_TIMEOUT: Duration = Duration::from_millis(250); /// A safety net for a run that never returns; the compute budget above is /// the limit a handler is actually held to, so this is generous enough /// that a descheduled thread on a loaded host still reports back. -const WALL_CLOCK_LIMIT: Duration = Duration::from_secs(10); +const WALL_CLOCK_LIMIT: Duration = Duration::from_secs(5); + +/// How often the caller checks the worker's CPU time against the budget. +const WATCHDOG_POLL: Duration = Duration::from_millis(5); /// boa loop iteration cap. Tight enough that `while(1){}` exits the VM /// well within the wall-clock budget on any reasonable host, loose @@ -93,54 +96,91 @@ fn run_handler_with_limits( ) -> JsExecution { let code = code.to_owned(); let event = event_json.to_vec(); - let (tx, rx) = mpsc::sync_channel::(1); + let (tx, rx) = mpsc::sync_channel::(2); // Each call gets its own thread because boa's `Context` holds // `Rc`s and is `!Send`. We can't pre-spawn a worker pool without // marshalling the script + event via channels anyway, so a fresh // thread per call is the simpler shape. - let _ = std::thread::Builder::new() + let spawned = std::thread::Builder::new() .name("cloudfront-js".to_string()) // Boa's bytecode VM is recursive so we want a generous stack. .stack_size(8 * 1024 * 1024) .spawn(move || { let clock = ComputeClock::start(); + let _ = tx.send(WorkerMessage::Started(clock.cpu.map(|(c, _)| c))); let mut result = run_handler_blocking(&code, &event, &clock); + // Covers a run that crossed the budget between the caller's polls. if clock.elapsed() > compute_budget { result = time_limit_exceeded(result.logs); } // If the receiver has timed out and gone away the send // simply errors; we don't care — the worker is being // abandoned. - let _ = tx.send(result); + let _ = tx.send(WorkerMessage::Done(result)); }); + if spawned.is_err() { + return worker_failed(); + } - match rx.recv_timeout(wall_limit) { - Ok(mut exec) => { - // Floor compute_utilization at 1% on success so callers - // don't mistake a successful run for an unrun one. - if exec.error.is_none() && exec.compute_utilization == 0 { - exec.compute_utilization = 1; + // Watch the worker: stop waiting as soon as its CPU time passes the + // budget, so a CPU-bound handler is cut off at the budget rather than + // holding the caller until the wall-clock safety net. Without a thread + // CPU clock on this platform, elapsed time stands in for it. + let wall_start = Instant::now(); + let mut worker_cpu: Option<(ThreadCpuClock, Duration)> = None; + loop { + let remaining = wall_limit.saturating_sub(wall_start.elapsed()); + match rx.recv_timeout(remaining.min(WATCHDOG_POLL)) { + Ok(WorkerMessage::Started(cpu)) => { + worker_cpu = cpu.and_then(|c| c.read().map(|start| (c, start))); } - exec - } - Err(mpsc::RecvTimeoutError::Timeout) => time_limit_exceeded(Vec::new()), - Err(mpsc::RecvTimeoutError::Disconnected) => { - // The worker thread either failed to spawn (`spawn()` errored - // and the sender was dropped) or panicked partway through. - // Surface that distinctly so a host-level problem doesn't get - // misdiagnosed as adversarial JS. - let msg = "function execution worker thread panicked or failed to spawn".to_string(); - JsExecution { - output: None, - error: Some(msg.clone()), - logs: vec![format!("ERROR: {msg}")], - compute_utilization: 101, + Ok(WorkerMessage::Done(mut exec)) => { + // Floor compute_utilization at 1% on success so callers + // don't mistake a successful run for an unrun one. + if exec.error.is_none() && exec.compute_utilization == 0 { + exec.compute_utilization = 1; + } + return exec; } + Err(mpsc::RecvTimeoutError::Timeout) => {} + Err(mpsc::RecvTimeoutError::Disconnected) => return worker_failed(), + } + if wall_start.elapsed() >= wall_limit { + return time_limit_exceeded(Vec::new()); + } + let used = match worker_cpu { + // A read fails once the thread has exited; its result is + // already on the way, so keep waiting for it. + Some((clock, start)) => clock.read().map(|now| now.saturating_sub(start)), + None => Some(wall_start.elapsed()), + }; + if used.is_some_and(|u| u > compute_budget) { + return time_limit_exceeded(Vec::new()); } } } +/// What the worker thread reports: first the handle to its CPU clock (so +/// the caller can watch the run), then the result. +enum WorkerMessage { + Started(Option), + Done(JsExecution), +} + +/// The worker thread either failed to spawn or panicked partway through. +/// Surfaced distinctly so a host-level problem doesn't get misdiagnosed as +/// adversarial JS. +fn worker_failed() -> JsExecution { + let msg = "function execution worker thread panicked or failed to spawn".to_string(); + JsExecution { + output: None, + error: Some(msg.clone()), + logs: vec![format!("ERROR: {msg}")], + compute_utilization: 101, + } +} + /// The error a run over its compute budget (or past the wall-clock safety /// net) reports, keeping whatever it logged before the limit. fn time_limit_exceeded(mut logs: Vec) -> JsExecution { @@ -160,43 +200,118 @@ fn time_limit_exceeded(mut logs: Vec) -> JsExecution { /// Measures the compute a run consumed: CPU time on the current thread /// where the platform reports it, elapsed time otherwise. struct ComputeClock { - cpu_start: Option, + cpu: Option<(ThreadCpuClock, Duration)>, wall_start: Instant, } impl ComputeClock { fn start() -> Self { Self { - cpu_start: thread_cpu_time(), + cpu: ThreadCpuClock::current().and_then(|c| c.read().map(|start| (c, start))), wall_start: Instant::now(), } } fn elapsed(&self) -> Duration { - match (self.cpu_start, thread_cpu_time()) { - (Some(start), Some(now)) => now.saturating_sub(start), - _ => self.wall_start.elapsed(), + match self + .cpu + .and_then(|(c, start)| c.read().map(|now| (now, start))) + { + Some((now, start)) => now.saturating_sub(start), + None => self.wall_start.elapsed(), } } } -/// CPU time consumed so far by the calling thread. -#[cfg(unix)] -fn thread_cpu_time() -> Option { - let mut ts = libc::timespec { - tv_sec: 0, - tv_nsec: 0, - }; - // SAFETY: `ts` is a valid, writable timespec; CLOCK_THREAD_CPUTIME_ID - // is supported on Linux and macOS, and a failure is reported by the - // return code rather than by leaving `ts` partially written. - let rc = unsafe { libc::clock_gettime(libc::CLOCK_THREAD_CPUTIME_ID, &mut ts) }; - (rc == 0).then(|| Duration::new(ts.tv_sec as u64, ts.tv_nsec as u32)) +/// A handle to one thread's CPU-time clock that any thread can read, so the +/// caller can watch the worker while it runs. Linux exposes a per-thread +/// clock id; macOS exposes the thread's Mach port. Elsewhere there is none +/// and callers fall back to elapsed time. +#[derive(Clone, Copy)] +struct ThreadCpuClock { + #[cfg(any(target_os = "linux", target_os = "android"))] + clock_id: libc::clockid_t, + #[cfg(any(target_os = "macos", target_os = "ios"))] + thread: libc::mach_port_t, } -#[cfg(not(unix))] -fn thread_cpu_time() -> Option { - None +impl ThreadCpuClock { + /// The calling thread's clock. + #[cfg(any(target_os = "linux", target_os = "android"))] + fn current() -> Option { + let mut clock_id: libc::clockid_t = 0; + // SAFETY: `clock_id` is a valid out-pointer, and `pthread_self()` is + // always a live thread (the caller). + let rc = unsafe { libc::pthread_getcpuclockid(libc::pthread_self(), &mut clock_id) }; + (rc == 0).then_some(Self { clock_id }) + } + + #[cfg(any(target_os = "macos", target_os = "ios"))] + fn current() -> Option { + // SAFETY: `pthread_self()` is always a live thread; the returned port + // is borrowed from the pthread (no reference to release). + let thread = unsafe { libc::pthread_mach_thread_np(libc::pthread_self()) }; + (thread != 0).then_some(Self { thread }) + } + + #[cfg(not(any( + target_os = "linux", + target_os = "android", + target_os = "macos", + target_os = "ios" + )))] + fn current() -> Option { + None + } + + /// CPU time the thread has consumed, or `None` once it has exited. + #[cfg(any(target_os = "linux", target_os = "android"))] + fn read(self) -> Option { + let mut ts = libc::timespec { + tv_sec: 0, + tv_nsec: 0, + }; + // SAFETY: `ts` is a valid, writable timespec; a stale clock id (the + // thread exited) is reported through the return code. + let rc = unsafe { libc::clock_gettime(self.clock_id, &mut ts) }; + (rc == 0).then(|| Duration::new(ts.tv_sec as u64, ts.tv_nsec as u32)) + } + + #[cfg(any(target_os = "macos", target_os = "ios"))] + fn read(self) -> Option { + // SAFETY: an all-zero `thread_basic_info` is a valid value for this + // plain-integer struct. + let mut info: libc::thread_basic_info = unsafe { std::mem::zeroed() }; + let mut count = libc::THREAD_BASIC_INFO_COUNT; + // SAFETY: `info` is large enough for THREAD_BASIC_INFO_COUNT + // integers, which `count` declares; a dead thread's port is reported + // through the return code. + let rc = unsafe { + libc::thread_info( + self.thread, + libc::THREAD_BASIC_INFO as libc::thread_flavor_t, + &mut info as *mut libc::thread_basic_info as libc::thread_info_t, + &mut count, + ) + }; + if rc != libc::KERN_SUCCESS { + return None; + } + let micros = |t: libc::time_value_t| t.seconds as u64 * 1_000_000 + t.microseconds as u64; + Some(Duration::from_micros( + micros(info.user_time) + micros(info.system_time), + )) + } + + #[cfg(not(any( + target_os = "linux", + target_os = "android", + target_os = "macos", + target_os = "ios" + )))] + fn read(self) -> Option { + None + } } fn run_handler_blocking(code: &str, event_json: &[u8], clock: &ComputeClock) -> JsExecution { @@ -452,8 +567,10 @@ mod tests { #[test] fn a_run_over_its_compute_budget_reports_the_time_limit() { + // A zero budget is exceeded by any run: either the caller's watchdog + // or the worker's own final check reports it, whichever sees it first. let exec = run_handler_with_limits( - r#"function handler(e) { console.log("ran"); return e; }"#, + r#"function handler(e) { return e; }"#, b"{}", Duration::ZERO, WALL_CLOCK_LIMIT, @@ -462,11 +579,7 @@ mod tests { let err = exec.error.expect("error"); assert!(err.contains("250ms time limit"), "got {err}"); assert!(exec.compute_utilization > 100); - assert!( - exec.logs.first().is_some_and(|l| l == "ran"), - "logs before the limit are kept: {:?}", - exec.logs - ); + assert!(exec.logs.iter().any(|l| l.starts_with("ERROR: "))); } #[test] @@ -496,22 +609,48 @@ mod tests { ); } - #[cfg(unix)] #[test] fn busy_work_accrues_compute() { + // Bounded by work done, not elapsed time, so a loaded host that + // deschedules this thread only makes it take longer. let clock = ComputeClock::start(); - let wall = Instant::now(); let mut x: u64 = 0; - while wall.elapsed() < Duration::from_millis(50) { - x = std::hint::black_box(x.wrapping_add(1)); + for i in 0..50_000_000u64 { + x = std::hint::black_box(x.wrapping_add(i)); } + std::hint::black_box(x); assert!( - clock.elapsed() >= Duration::from_millis(20), - "50ms of busy work measured {:?}", + clock.elapsed() >= Duration::from_millis(1), + "the work measured {:?}", clock.elapsed() ); } + #[test] + fn a_cpu_bound_handler_is_cut_off_at_its_budget() { + // Catastrophic regex backtracking burns CPU inside one builtin call, + // so boa's loop and recursion caps never trip; only the compute + // budget can stop it. It runs for seconds (about 6s unoptimized, + // doubling per extra `a`), well past half the safety net, so the + // caller returning quickly proves the watchdog cut it off. The + // abandoned worker then finishes in the background. + let started = Instant::now(); + let exec = run_handler( + &format!( + r#"function handler() {{ return /^(a+)+$/.test("{}b"); }}"#, + "a".repeat(24) + ), + b"{}", + ); + let took = started.elapsed(); + let err = exec.error.expect("error"); + assert!(err.contains("time limit"), "got {err}"); + assert!( + took < WALL_CLOCK_LIMIT / 2, + "caller waited {took:?}; the budget should stop it well before the safety net" + ); + } + #[test] fn infinite_loop_is_killed_by_timeout() { let exec = run_handler(r#"function handler() { while(1){} }"#, b"{}"); diff --git a/website/content/docs/services/cloudfront.md b/website/content/docs/services/cloudfront.md index 59884b34e..7c66a91c8 100644 --- a/website/content/docs/services/cloudfront.md +++ b/website/content/docs/services/cloudfront.md @@ -20,7 +20,7 @@ fakecloud implements CloudFront's REST-XML control plane focused on the operatio - **Origin Request Policy** — `CreateOriginRequestPolicy`, `GetOriginRequestPolicy`, `GetOriginRequestPolicyConfig`, `UpdateOriginRequestPolicy`, `DeleteOriginRequestPolicy`, `ListOriginRequestPolicies`. Managed `Managed-CORS-S3Origin`, `Managed-CORS-CustomOrigin`, `Managed-AllViewer`, `Managed-UserAgentRefererHeaders`, `Managed-AllViewerExceptHostHeader` pre-seeded. - **Response Headers Policy** — `CreateResponseHeadersPolicy`, `GetResponseHeadersPolicy`, `GetResponseHeadersPolicyConfig`, `UpdateResponseHeadersPolicy`, `DeleteResponseHeadersPolicy`, `ListResponseHeadersPolicies`. Managed `Managed-SimpleCORS`, `Managed-CORS-With-Preflight`, `Managed-SecurityHeadersPolicy` pre-seeded. - **Continuous Deployment Policy** — `CreateContinuousDeploymentPolicy`, `GetContinuousDeploymentPolicy`, `GetContinuousDeploymentPolicyConfig`, `UpdateContinuousDeploymentPolicy`, `DeleteContinuousDeploymentPolicy`, `ListContinuousDeploymentPolicies`. -- **CloudFront Functions** — `CreateFunction`, `DescribeFunction`, `GetFunction` (returns raw source bytes), `UpdateFunction`, `DeleteFunction`, `ListFunctions`, `PublishFunction` (DEVELOPMENT -> LIVE), `TestFunction`. `TestFunction` executes the user's `function handler(event) { ... }` against the supplied `EventObject` via an embedded `boa_engine` JavaScript runtime: the handler return value is JSON-encoded into ``, thrown errors land in ``, and `console.log` / `console.error` lines are captured into ``. `` reports a 0..=100 share of the per-request CPU budget on success and saturates past 100 on failure. `Stage` (DEVELOPMENT | LIVE) picks which version's source to run: DEVELOPMENT is the latest CreateFunction / UpdateFunction body, LIVE is the snapshot frozen at the most recent PublishFunction call (so post-publish updates don't leak into the published behaviour). Execution is capped by boa's loop-iteration and recursion limits plus a 250ms compute budget, measured as CPU time on the executing thread like AWS's ~1ms production budget (looser, so ordinary handlers never trip it unoptimized, while a CPU-bound handler still fails with the time-limit error); because it counts CPU rather than elapsed time, a loaded host that stalls the thread does not fail a handler that did little work, and a 10s wall-clock safety net abandons a run that never returns. +- **CloudFront Functions** — `CreateFunction`, `DescribeFunction`, `GetFunction` (returns raw source bytes), `UpdateFunction`, `DeleteFunction`, `ListFunctions`, `PublishFunction` (DEVELOPMENT -> LIVE), `TestFunction`. `TestFunction` executes the user's `function handler(event) { ... }` against the supplied `EventObject` via an embedded `boa_engine` JavaScript runtime: the handler return value is JSON-encoded into ``, thrown errors land in ``, and `console.log` / `console.error` lines are captured into ``. `` reports a 0..=100 share of the per-request CPU budget on success and saturates past 100 on failure. `Stage` (DEVELOPMENT | LIVE) picks which version's source to run: DEVELOPMENT is the latest CreateFunction / UpdateFunction body, LIVE is the snapshot frozen at the most recent PublishFunction call (so post-publish updates don't leak into the published behaviour). Execution is capped by boa's loop-iteration and recursion limits plus a 250ms compute budget, measured as CPU time on the executing thread like AWS's ~1ms production budget (looser, so ordinary handlers never trip it unoptimized, while a CPU-bound handler still fails with the time-limit error); because it counts CPU rather than elapsed time, a loaded host that stalls the thread does not fail a handler that did little work, and a 5s wall-clock safety net abandons a run that never returns. - **Connection Functions** — `CreateConnectionFunction`, `DescribeConnectionFunction`, `GetConnectionFunction` (raw source), `UpdateConnectionFunction`, `DeleteConnectionFunction`, `ListConnectionFunctions`, `PublishConnectionFunction`, `TestConnectionFunction`. `TestConnectionFunction` runs the same `boa_engine` runtime against the supplied `ConnectionObject` and emits the parallel `` / `` / `` / `` shape, with the same DEVELOPMENT-vs-LIVE stage selection as `TestFunction`. - **Public Keys + Key Groups** — full CRUD with `CallerReference` immutability on update. - **Key Value Stores** — `CreateKeyValueStore` (with optional `ImportSource`), `DescribeKeyValueStore`, `UpdateKeyValueStore`, `DeleteKeyValueStore`, `ListKeyValueStores`. From 8ed84b4bee4f206ac6de70fa4accd922fd9f7574 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 21:35:07 -0300 Subject: [PATCH 3/4] fix(cloudfront): measure the watchdog from the worker's own start reading - Charge nothing against the budget until the worker reports in: a thread the host is slow to schedule has used no compute, and only the wall-clock safety net applies before it starts. Previously the caller counted elapsed time from before the spawn, moving the false alarm into thread startup. - Send the worker's own starting CPU reading with its clock, so CPU it burns before the caller reads the clock is counted and both checks measure from the same point. - Test: a worker delayed longer than the whole budget before starting still runs a trivial handler successfully (fails with the time-limit error when pre-start wall time is charged). --- crates/fakecloud-cloudfront/src/js_runtime.rs | 76 +++++++++++++++---- 1 file changed, 63 insertions(+), 13 deletions(-) diff --git a/crates/fakecloud-cloudfront/src/js_runtime.rs b/crates/fakecloud-cloudfront/src/js_runtime.rs index 80844527c..1b02fedd4 100644 --- a/crates/fakecloud-cloudfront/src/js_runtime.rs +++ b/crates/fakecloud-cloudfront/src/js_runtime.rs @@ -85,14 +85,17 @@ pub(crate) struct JsExecution { /// dedicated worker thread, holding it to the `EXECUTION_TIMEOUT` compute /// budget. pub(crate) fn run_handler(code: &str, event_json: &[u8]) -> JsExecution { - run_handler_with_limits(code, event_json, EXECUTION_TIMEOUT, WALL_CLOCK_LIMIT) + run_handler_with_limits(code, event_json, EXECUTION_TIMEOUT, WALL_CLOCK_LIMIT, || {}) } +/// `before_start` runs on the worker thread before its clock starts; tests +/// use it to stand in for a thread the host is slow to schedule. fn run_handler_with_limits( code: &str, event_json: &[u8], compute_budget: Duration, wall_limit: Duration, + before_start: fn(), ) -> JsExecution { let code = code.to_owned(); let event = event_json.to_vec(); @@ -107,8 +110,9 @@ fn run_handler_with_limits( // Boa's bytecode VM is recursive so we want a generous stack. .stack_size(8 * 1024 * 1024) .spawn(move || { + before_start(); let clock = ComputeClock::start(); - let _ = tx.send(WorkerMessage::Started(clock.cpu.map(|(c, _)| c))); + let _ = tx.send(WorkerMessage::Started(clock.cpu)); let mut result = run_handler_blocking(&code, &event, &clock); // Covers a run that crossed the budget between the caller's polls. if clock.elapsed() > compute_budget { @@ -125,15 +129,21 @@ fn run_handler_with_limits( // Watch the worker: stop waiting as soon as its CPU time passes the // budget, so a CPU-bound handler is cut off at the budget rather than - // holding the caller until the wall-clock safety net. Without a thread - // CPU clock on this platform, elapsed time stands in for it. + // holding the caller until the wall-clock safety net. Until the worker + // reports in, nothing counts against the budget -- a thread the host is + // slow to schedule has used none -- and only the safety net applies. + // Without a thread CPU clock on this platform, time elapsed since the + // worker started stands in for it. let wall_start = Instant::now(); - let mut worker_cpu: Option<(ThreadCpuClock, Duration)> = None; + let mut watch = Watch::NotStarted; loop { let remaining = wall_limit.saturating_sub(wall_start.elapsed()); match rx.recv_timeout(remaining.min(WATCHDOG_POLL)) { Ok(WorkerMessage::Started(cpu)) => { - worker_cpu = cpu.and_then(|c| c.read().map(|start| (c, start))); + watch = match cpu { + Some((clock, start)) => Watch::Cpu(clock, start), + None => Watch::Elapsed(Instant::now()), + }; } Ok(WorkerMessage::Done(mut exec)) => { // Floor compute_utilization at 1% on success so callers @@ -149,11 +159,12 @@ fn run_handler_with_limits( if wall_start.elapsed() >= wall_limit { return time_limit_exceeded(Vec::new()); } - let used = match worker_cpu { + let used = match watch { + Watch::NotStarted => None, // A read fails once the thread has exited; its result is // already on the way, so keep waiting for it. - Some((clock, start)) => clock.read().map(|now| now.saturating_sub(start)), - None => Some(wall_start.elapsed()), + Watch::Cpu(clock, start) => clock.read().map(|now| now.saturating_sub(start)), + Watch::Elapsed(started) => Some(started.elapsed()), }; if used.is_some_and(|u| u > compute_budget) { return time_limit_exceeded(Vec::new()); @@ -161,13 +172,24 @@ fn run_handler_with_limits( } } -/// What the worker thread reports: first the handle to its CPU clock (so -/// the caller can watch the run), then the result. +/// What the worker thread reports: first its CPU clock and the reading it +/// started from (so the caller measures the run from the same point as the +/// worker's own check), then the result. enum WorkerMessage { - Started(Option), + Started(Option<(ThreadCpuClock, Duration)>), Done(JsExecution), } +/// How the caller measures the worker's compute while waiting on it. +enum Watch { + /// The worker has not reported in yet; nothing counts against the budget. + NotStarted, + /// The worker's CPU clock and its reading when the run started. + Cpu(ThreadCpuClock, Duration), + /// No thread CPU clock on this platform: time since the worker started. + Elapsed(Instant), +} + /// The worker thread either failed to spawn or panicked partway through. /// Surfaced distinctly so a host-level problem doesn't get misdiagnosed as /// adversarial JS. @@ -574,6 +596,7 @@ mod tests { b"{}", Duration::ZERO, WALL_CLOCK_LIMIT, + || {}, ); assert!(exec.output.is_none(), "got output {:?}", exec.output); let err = exec.error.expect("error"); @@ -589,13 +612,40 @@ mod tests { b"{}", EXECUTION_TIMEOUT, Duration::ZERO, + || {}, ); let err = exec.error.expect("error"); assert!(err.contains("time limit"), "got {err}"); assert!(exec.compute_utilization > 100); } - #[cfg(unix)] + #[cfg(any( + target_os = "linux", + target_os = "android", + target_os = "macos", + target_os = "ios" + ))] + #[test] + fn a_worker_slow_to_start_is_not_charged_for_the_wait() { + // The host took longer than the whole budget to get the worker going; + // the handler itself is trivial and must still succeed. + let exec = run_handler_with_limits( + r#"function handler(e) { return e; }"#, + b"{}", + EXECUTION_TIMEOUT, + WALL_CLOCK_LIMIT, + || std::thread::sleep(EXECUTION_TIMEOUT + Duration::from_millis(150)), + ); + assert!(exec.error.is_none(), "unexpected error: {:?}", exec.error); + assert_eq!(exec.output.as_deref(), Some("{}")); + } + + #[cfg(any( + target_os = "linux", + target_os = "android", + target_os = "macos", + target_os = "ios" + ))] #[test] fn a_descheduled_thread_accrues_no_compute() { // The budget is CPU time: a thread that is not running -- here From b6025a75a54ec062aa46ff1712f3c830235138f2 Mon Sep 17 00:00:00 2001 From: Lucas Vieira Date: Sun, 13 Sep 2026 21:39:05 -0300 Subject: [PATCH 4/4] fix(cloudfront): wait on a function handler off the async runtime TestFunction and TestConnectionFunction called the blocking runner from the async request path, holding a tokio worker thread for the whole wait. With the compute budget measured in CPU time, that wait can reach the 5s wall-clock safety net when the host stalls the worker, and on a current-thread runtime it stalls everything else. Both handlers are async now and run the wait through spawn_blocking; their state lock is scoped so no guard is held across the await. Test: a current-thread runtime keeps ticking while a CPU-bound handler runs to its budget (ticks 0 times when the runner is called directly). --- .../src/cfunctions_service.rs | 32 +++++++------ .../src/functions_service.rs | 32 +++++++------ crates/fakecloud-cloudfront/src/js_runtime.rs | 46 +++++++++++++++++++ crates/fakecloud-cloudfront/src/service.rs | 4 +- 4 files changed, 82 insertions(+), 32 deletions(-) diff --git a/crates/fakecloud-cloudfront/src/cfunctions_service.rs b/crates/fakecloud-cloudfront/src/cfunctions_service.rs index 1f7b729cb..fbdd69508 100644 --- a/crates/fakecloud-cloudfront/src/cfunctions_service.rs +++ b/crates/fakecloud-cloudfront/src/cfunctions_service.rs @@ -274,7 +274,7 @@ impl CloudFrontService { Ok(xml_with_etag(StatusCode::OK, body, &snap.etag, None)) } - pub(crate) fn test_connection_function( + pub(crate) async fn test_connection_function( &self, req: &AwsRequest, route: &Route, @@ -288,19 +288,21 @@ impl CloudFrontService { let event_bytes = base64::engine::general_purpose::STANDARD .decode(parsed.connection_object.trim().as_bytes()) .map_err(|e| invalid_argument(format!("ConnectionObject is not valid base64: {e}")))?; - let state = self.state.read(); - let f = state - .accounts - .get(DEFAULT_ACCOUNT) - .and_then(|a| a.connection_functions.get(&name).cloned()) - .ok_or_else(|| { - aws_error( - StatusCode::NOT_FOUND, - "EntityNotFound", - format!("ConnectionFunction {name} does not exist"), - ) - })?; - drop(state); + // Scoped so the lock guard is gone before the await below. + let f = { + let state = self.state.read(); + state + .accounts + .get(DEFAULT_ACCOUNT) + .and_then(|a| a.connection_functions.get(&name).cloned()) + .ok_or_else(|| { + aws_error( + StatusCode::NOT_FOUND, + "EntityNotFound", + format!("ConnectionFunction {name} does not exist"), + ) + })? + }; if f.etag != if_match { return Err(precondition_failed()); } @@ -317,7 +319,7 @@ impl CloudFrontService { }; let code = std::str::from_utf8(source) .map_err(|e| invalid_argument(format!("function code is not valid UTF-8: {e}")))?; - let exec = crate::js_runtime::run_handler(code, &event_bytes); + let exec = crate::js_runtime::run_handler_off_runtime(code.to_owned(), event_bytes).await; let mut body = String::with_capacity(1024); body.push_str(XML_DECL); diff --git a/crates/fakecloud-cloudfront/src/functions_service.rs b/crates/fakecloud-cloudfront/src/functions_service.rs index d72c64419..6cc32f56b 100644 --- a/crates/fakecloud-cloudfront/src/functions_service.rs +++ b/crates/fakecloud-cloudfront/src/functions_service.rs @@ -256,7 +256,7 @@ impl CloudFrontService { Ok(xml_with_etag(StatusCode::OK, body, &snap.etag, None)) } - pub(crate) fn test_function( + pub(crate) async fn test_function( &self, req: &AwsRequest, route: &Route, @@ -269,19 +269,21 @@ impl CloudFrontService { .decode(parsed.event_object.trim().as_bytes()) .map_err(|e| invalid_argument(format!("EventObject is not valid base64: {e}")))?; - let state = self.state.read(); - let f = state - .accounts - .get(DEFAULT_ACCOUNT) - .and_then(|a| a.functions.get(&name).cloned()) - .ok_or_else(|| { - aws_error( - StatusCode::NOT_FOUND, - "NoSuchFunctionExists", - format!("The specified function does not exist: {name}"), - ) - })?; - drop(state); + // Scoped so the lock guard is gone before the await below. + let f = { + let state = self.state.read(); + state + .accounts + .get(DEFAULT_ACCOUNT) + .and_then(|a| a.functions.get(&name).cloned()) + .ok_or_else(|| { + aws_error( + StatusCode::NOT_FOUND, + "NoSuchFunctionExists", + format!("The specified function does not exist: {name}"), + ) + })? + }; if f.etag != if_match { return Err(precondition_failed()); } @@ -303,7 +305,7 @@ impl CloudFrontService { .unwrap_or_else(|_| source_b64.as_bytes().to_vec()); let code = String::from_utf8(code_bytes) .map_err(|e| invalid_argument(format!("function code is not valid UTF-8: {e}")))?; - let exec = crate::js_runtime::run_handler(&code, &event_bytes); + let exec = crate::js_runtime::run_handler_off_runtime(code, event_bytes).await; let mut body = String::with_capacity(1024); body.push_str(XML_DECL); diff --git a/crates/fakecloud-cloudfront/src/js_runtime.rs b/crates/fakecloud-cloudfront/src/js_runtime.rs index 1b02fedd4..6ccff067c 100644 --- a/crates/fakecloud-cloudfront/src/js_runtime.rs +++ b/crates/fakecloud-cloudfront/src/js_runtime.rs @@ -81,6 +81,17 @@ pub(crate) struct JsExecution { pub compute_utilization: u32, } +/// [`run_handler`] for async request handlers. The wait for the worker +/// blocks -- up to the compute budget, or the wall-clock safety net when the +/// host is stalling the worker -- so it runs on tokio's blocking pool rather +/// than holding a runtime worker thread (or, on a current-thread runtime, +/// the whole runtime) for that long. +pub(crate) async fn run_handler_off_runtime(code: String, event_json: Vec) -> JsExecution { + tokio::task::spawn_blocking(move || run_handler(&code, &event_json)) + .await + .unwrap_or_else(|_| worker_failed()) +} + /// Run `handler(event)` defined in `code` against `event_json` on a /// dedicated worker thread, holding it to the `EXECUTION_TIMEOUT` compute /// budget. @@ -701,6 +712,41 @@ mod tests { ); } + #[tokio::test] + async fn waiting_on_a_handler_leaves_the_runtime_free() { + // A current-thread runtime: if the wait blocked it, the ticker below + // could not run until the handler returned. + use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; + use std::sync::Arc; + let done = Arc::new(AtomicBool::new(false)); + let ticks = Arc::new(AtomicU32::new(0)); + let ticker = tokio::spawn({ + let (done, ticks) = (done.clone(), ticks.clone()); + async move { + while !done.load(Ordering::SeqCst) { + tokio::time::sleep(Duration::from_millis(10)).await; + ticks.fetch_add(1, Ordering::SeqCst); + } + } + }); + let exec = run_handler_off_runtime( + format!( + r#"function handler() {{ return /^(a+)+$/.test("{}b"); }}"#, + "a".repeat(24) + ), + b"{}".to_vec(), + ) + .await; + let ticked = ticks.load(Ordering::SeqCst); + done.store(true, Ordering::SeqCst); + ticker.await.unwrap(); + assert!(exec.error.is_some_and(|e| e.contains("time limit"))); + assert!( + ticked >= 5, + "the runtime ticked {ticked} times while the handler ran for the budget" + ); + } + #[test] fn infinite_loop_is_killed_by_timeout() { let exec = run_handler(r#"function handler() { while(1){} }"#, b"{}"); diff --git a/crates/fakecloud-cloudfront/src/service.rs b/crates/fakecloud-cloudfront/src/service.rs index c7f50b1be..2cbc933cd 100644 --- a/crates/fakecloud-cloudfront/src/service.rs +++ b/crates/fakecloud-cloudfront/src/service.rs @@ -524,7 +524,7 @@ impl AwsService for CloudFrontService { "DeleteFunction" => self.delete_function(&req, &resolved), "ListFunctions" => self.list_functions(&req), "PublishFunction" => self.publish_function(&req, &resolved), - "TestFunction" => self.test_function(&req, &resolved), + "TestFunction" => self.test_function(&req, &resolved).await, "CreatePublicKey" => self.create_public_key(&req), "GetPublicKey" => self.get_public_key(&resolved), "GetPublicKeyConfig" => self.get_public_key_config(&resolved), @@ -649,7 +649,7 @@ impl AwsService for CloudFrontService { "DeleteConnectionFunction" => self.delete_connection_function(&req, &resolved), "ListConnectionFunctions" => self.list_connection_functions(&req), "PublishConnectionFunction" => self.publish_connection_function(&req, &resolved), - "TestConnectionFunction" => self.test_connection_function(&req, &resolved), + "TestConnectionFunction" => self.test_connection_function(&req, &resolved).await, other => Err(aws_error( StatusCode::NOT_IMPLEMENTED, "InvalidAction",