Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 9 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -660,8 +660,13 @@ Options accept exactly one of `entry: string` or inline `source: string`, plus `
inline `source` and defaults to JavaScript. With the generated crate's `typescript-runtime`
feature enabled, TypeScript entry files are selected by their `.ts`, `.mts`, or `.cts` extension,
and inline jobs can select `language: "typescript"`.
Output and completion wake the parent without periodic polling. CPU deadlines start
after child-runtime initialization and are enforced by the QuickJS interrupt handler. Entry modules run for their side effects; an exported
Output and completion wake the parent without periodic polling. Timeout budgets start
after child-runtime initialization. An async wait observes deadlines while jobs yield; for
non-yielding JavaScript, the QuickJS interrupt handler samples the clock at an adaptive stride
targeting one check per 10% of `timeoutMs` during steady CPU work. This is a best-effort budget,
not an exact cutoff or a maximum-overrun guarantee. A job that finishes between checks can
succeed after its literal deadline; no extra clock check occurs at completion. Entry modules run
for their side effects; an exported
`default` function (or `run` function when there is no default function) is invoked and awaited.
Relative entry paths and imports resolve from `cwd`. Inline `source` is an async function body:
top-level `await` and `return` are supported, while static `import`/`export` declarations are not.
Expand All @@ -673,7 +678,8 @@ With `overflow: "terminate"`, exceeding either stream's bound rejects the job. W
`overflow: "truncate"`, execution continues, captured output is bounded, and `overflowed` is
`true` when either stream exceeded the bound.
Cancellation is cooperative for queued or yielding code and cannot preempt a tight loop already
running on the same thread; use a positive `timeoutMs` when that guarantee is required. A job that
running on the same thread; a positive `timeoutMs` is the best-effort escape mechanism for such a
loop. Cancellation remains a guest-local flag check on every QuickJS interrupt. A job that
never settles must be cancelled or given a timeout before the enclosing component invocation can
finish. A runtime accepts at
most eight active jobs, and child executions cannot recursively create more execution jobs.
Expand Down
22 changes: 18 additions & 4 deletions crates/wasm-rquickjs/skeleton/src/builtin/execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,10 @@ use std::sync::atomic::{AtomicBool, Ordering};
use std::task::Poll;
use std::time::{Duration, Instant};

#[path = "execution_timeout.rs"]
mod execution_timeout;
use execution_timeout::ExecutionTimeoutSampler;

const MAX_ACTIVE_JOBS: usize = 8;
const MAX_TIMEOUT_MS: u64 = u64::MAX / 1_000_000;
const MAX_OUTPUT_BYTES: usize = 64 * 1024 * 1024;
Expand Down Expand Up @@ -483,9 +487,17 @@ async fn run_job(options: ExecutionOptions, job: Rc<ExecutionJob>) {
// The public timeout budget starts when user code begins, after runtime and
// builtin initialization. The interrupt handler is also needed for tight
// loops that cannot cooperatively yield to the timer future.
let deadline = options
.timeout_ms
.and_then(|ms| Instant::now().checked_add(Duration::from_millis(ms)));
let timeout = options.timeout_ms.and_then(|ms| {
let started_at = Instant::now();
let duration = Duration::from_millis(ms);
started_at
.checked_add(duration)
.map(|deadline| (started_at, deadline, duration))
});
let deadline = timeout.map(|(_, deadline, _)| deadline);
let mut timeout_sampler = timeout.map(|(started_at, deadline, duration)| {
ExecutionTimeoutSampler::new(started_at, deadline, duration)
});
let cancelled = job.cancel.clone();
let timed_out = job.timed_out.clone();
runtime
Expand All @@ -494,7 +506,9 @@ async fn run_job(options: ExecutionOptions, job: Rc<ExecutionJob>) {
if cancelled.load(Ordering::Relaxed) {
return true;
}
if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
if timeout_sampler.as_mut().is_some_and(|sampler| {
sampler.expired(Instant::now)
}) {
timed_out.store(true, Ordering::Relaxed);
return true;
}
Expand Down
65 changes: 65 additions & 0 deletions crates/wasm-rquickjs/skeleton/src/builtin/execution_timeout.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
use std::time::{Duration, Instant};

const INITIAL_INTERRUPT_STRIDE: u64 = 64;
const MAX_STRIDE_GROWTH: u64 = 16;

/// Samples a wall-clock deadline without importing the clock on every QuickJS
/// interrupt. The stride is deliberately a soft target: callback frequency can
/// change, and finishing between samples does not trigger another clock read.
pub(super) struct ExecutionTimeoutSampler {
deadline: Instant,
target_interval: Duration,
last_sample_at: Instant,
stride: u64,
interrupts_until_sample: u64,
expired: bool,
}

impl ExecutionTimeoutSampler {
pub(super) fn new(started_at: Instant, deadline: Instant, timeout: Duration) -> Self {
Self {
deadline,
target_interval: timeout / 10,
last_sample_at: started_at,
stride: INITIAL_INTERRUPT_STRIDE,
interrupts_until_sample: INITIAL_INTERRUPT_STRIDE,
expired: false,
}
}

/// Returns true only when a sampled clock value reaches the deadline.
/// `clock` is never called on unsampled interrupts.
pub(super) fn expired(&mut self, clock: impl FnOnce() -> Instant) -> bool {
if self.expired {
return true;
}
self.interrupts_until_sample -= 1;
if self.interrupts_until_sample != 0 {
return false;
}

let now = clock();
if now >= self.deadline {
self.expired = true;
return true;
}

let elapsed_nanos = now
.saturating_duration_since(self.last_sample_at)
.as_nanos()
.max(1);
let estimated_stride = self
.target_interval
.as_nanos()
.saturating_mul(u128::from(self.stride))
/ elapsed_nanos;
let growth_limit = self.stride.saturating_mul(MAX_STRIDE_GROWTH);
self.stride = estimated_stride
.clamp(1, u128::from(growth_limit))
.try_into()
.unwrap_or(u64::MAX);
self.interrupts_until_sample = self.stride;
self.last_sample_at = now;
false
}
}
64 changes: 64 additions & 0 deletions tests/runtime/execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,70 @@ use crate::common::{CompiledTest, FeatureCombination, invoke_and_capture_output}
use camino::Utf8Path;
use test_r::{test, test_dep};

#[path = "../../crates/wasm-rquickjs/skeleton/src/builtin/execution_timeout.rs"]
mod execution_timeout;

#[test]
fn timeout_sampler_uses_sparse_clock_reads() {
let started = std::time::Instant::now();
let timeout = std::time::Duration::from_secs(1);
let mut sampler =
execution_timeout::ExecutionTimeoutSampler::new(started, started + timeout, timeout);
let mut now = started;
let mut reads = 0;
for _ in 0..200_000 {
now += std::time::Duration::from_micros(10);
if sampler.expired(|| {
reads += 1;
now
}) {
assert!(now >= started + timeout);
assert!(reads <= 30, "too many clock reads: {reads}");
assert!(sampler.expired(|| panic!("expired sampler read the clock again")));
return;
}
}
panic!("continuing CPU work never observed the deadline");
}

#[test]
fn timeout_sampler_adapts_to_changed_interrupt_rate() {
let started = std::time::Instant::now();
let timeout = std::time::Duration::from_secs(1);
let mut sampler =
execution_timeout::ExecutionTimeoutSampler::new(started, started + timeout, timeout);
let mut now = started;
let mut reads = 0;
for interrupt in 0..200_000 {
now += std::time::Duration::from_micros(if interrupt < 50_000 { 10 } else { 100 });
if sampler.expired(|| {
reads += 1;
now
}) {
assert!(now >= started + timeout);
assert!(reads < 50, "too many clock reads: {reads}");
return;
}
}
panic!("continuing CPU work never observed the deadline after rate change");
}

#[test]
fn timeout_sampler_does_not_check_at_completion() {
let started = std::time::Instant::now();
let timeout = std::time::Duration::from_millis(1);
let mut sampler =
execution_timeout::ExecutionTimeoutSampler::new(started, started + timeout, timeout);
let mut reads = 0;
for _ in 0..63 {
assert!(!sampler.expired(|| {
reads += 1;
started + timeout * 2
}));
}
assert_eq!(reads, 0);
}

#[test_dep(tagged_as = "execution", scope = Cloneable)]
async fn compiled_execution() -> CompiledTest {
CompiledTest::new_with_features(
Expand Down
Loading