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/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 b7c7a08bc..6ccff067c 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,21 @@ 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(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 /// enough that small loops in real handlers run to completion. @@ -68,68 +81,273 @@ 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, 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, || {}) +} + +/// `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(); - 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 started = Instant::now(); - 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 result = run_handler_blocking(&code, &event, started); + before_start(); + let clock = ComputeClock::start(); + 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 { + 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(EXECUTION_TIMEOUT) { - 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. 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 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)) => { + watch = match cpu { + Some((clock, start)) => Watch::Cpu(clock, start), + None => Watch::Elapsed(Instant::now()), + }; } - 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, + 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(), } - 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, - } + if wall_start.elapsed() >= wall_limit { + return time_limit_exceeded(Vec::new()); + } + 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. + 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()); + } + } +} + +/// 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<(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. +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 { + 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: Option<(ThreadCpuClock, Duration)>, + wall_start: Instant, +} + +impl ComputeClock { + fn start() -> Self { + Self { + cpu: ThreadCpuClock::current().and_then(|c| c.read().map(|start| (c, start))), + wall_start: Instant::now(), + } + } + + fn elapsed(&self) -> Duration { + 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(), + } + } +} + +/// 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, +} + +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], started: Instant) -> JsExecution { +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 +356,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 +375,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 +385,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 +406,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 +415,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 +424,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 +598,155 @@ mod tests { assert!(exec.error.is_some()); } + #[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) { 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.iter().any(|l| l.starts_with("ERROR: "))); + } + + #[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(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 + // 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() + ); + } + + #[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 mut x: u64 = 0; + 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(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" + ); + } + + #[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", diff --git a/website/content/docs/services/cloudfront.md b/website/content/docs/services/cloudfront.md index 8bddca602..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 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 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`.