diff --git a/CHANGELOG.md b/CHANGELOG.md index 4ddf711..f5cf381 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,36 @@ All notable changes to WASM-OJ are recorded here. Releases follow ## Unreleased +- WebKit now stops a terminated or killed Worker that is running a metered program within well under + a millisecond, so the Web Lock liveness check reports it at once. JavaScriptCore never acts on + `Worker.terminate()` inside Wasm, only at JavaScript checkpoints such as `Atomics.wait`, so a + killed runner or interactive side kept computing until its instruction budget ran out (up to tens + of seconds) and the session ended at its wall limit. Every 2^20 units, at the next function entry + or loop iteration, the meter now calls an uncharged `wasm_oj_metering.safepoint` import, which the + browser runtime turns into such a checkpoint and the native runtime ignores. Costs, cost profiles' + cost values and the exhaustion point are unchanged; the refreshed runtime identity changes cost + profile identifiers. +- Browser runner, compiler, compiler stage (rustc, Go, Java) and interactive side Workers that die + without an `error` event, for example when the browser terminates them, now reject their + operation promptly as a runner or compiler failure. A silently killed Worker used to leave the operation waiting for its wall-time + limit, which reported the student's program as `wall-time-limit`. Each module Worker holds a Web + Lock for its lifetime and its owner treats the lock's release as a crash; without Web Locks + nothing changes. Interactive pipes now wake every 100 ms while waiting, so a terminated side + Worker stops promptly in WebKit, which otherwise keeps it blocked in `Atomics.wait`. +- Fixed WebKit taking about 25 s to stop an empty C++ `for(;;);` at the default instruction budget, + so the run usually hit its wall limit instead. JavaScriptCore never enters optimized code inside a + loop of a function that has no parameters or locals, and keeps calling its tier-up slow path, so + such a metered loop ran about 20 times slower than in Chromium. Instrumentation now gives these + functions one unused local, which changes neither behaviour nor cost; WebKit stops the loop at the + budget in about 0.2 s. The refreshed runtime identity changes cost profiles. +- Interactive writes no longer fail with `EPIPE` after the other side exits or closes its stdin. The + bytes are recorded in the transcript once and dropped, and the writer keeps running, as with a + judge that keeps draining both pipes. An interactor that replies to a contestant that already + exited now reads EOF and exits with its own verdict; a CPython interactor used to exit 120 and + repeat its reply up to five times in the transcript. Reads still return the remaining buffered + bytes and then EOF. This applies to the server and the browser. The refreshed runtime identity + changes cost profiles. + ## 0.2.4 - 2026-10-09 - Fixed a host process crash (`Uncaught Error: write EPIPE`) when `ServerRunner` cancelled or diff --git a/crates/runtime-core/src/interactive.rs b/crates/runtime-core/src/interactive.rs index 77a9f29..429bbb1 100644 --- a/crates/runtime-core/src/interactive.rs +++ b/crates/runtime-core/src/interactive.rs @@ -3,7 +3,9 @@ use crate::deterministic::{VirtualClock, attach_interactive_deterministic_import use crate::filesystem::{ RuntimeProjectFilesystem, is_normalized_guest_path, runtime_project_files, }; -use crate::meter::{CostPoints, MeterState, instrument_wasm, meter_state, remaining_points}; +use crate::meter::{ + CostPoints, MeterState, attach_safepoint, instrument_wasm, meter_state, remaining_points, +}; use crate::module_imports::attach_declared_memory_imports; use crate::module_policy::{ DEFERRED_START_EXPORT, defer_start_section, enforce_memory_limit, @@ -121,6 +123,9 @@ struct InteractiveOutput { } impl AsyncWrite for InteractiveOutput { + /// Once the peer has closed its stdin, for example by exiting, the pipe reports a broken + /// pipe. The bytes are already in the transcript, so the write succeeds and they are dropped, + /// as when a judge keeps draining a pipe whose reader is gone. fn poll_write( mut self: Pin<&mut Self>, context: &mut Context<'_>, @@ -128,7 +133,12 @@ impl AsyncWrite for InteractiveOutput { ) -> Poll> { match Pin::new(&mut self.capture).poll_write(context, buffer) { Poll::Ready(Ok(written)) => { - Pin::new(&mut self.pipe).poll_write(context, &buffer[..written]) + match Pin::new(&mut self.pipe).poll_write(context, &buffer[..written]) { + Poll::Ready(Err(error)) if error.kind() == io::ErrorKind::BrokenPipe => { + Poll::Ready(Ok(written)) + } + result => result, + } } result => result, } @@ -563,6 +573,7 @@ fn interactive_runtime( clock.clone(), startup_entropy_bytes, ); + attach_safepoint(store, &mut imports); attach_capability_denials(store, module, &mut imports).map_err(io::Error::other)?; attach_declared_memory_imports(store, module, &mut imports).map_err(io::Error::other)?; Ok(imports) @@ -1120,6 +1131,128 @@ mod tests { assert_eq!(interactive.metrics.cost, standalone.metrics.cost); } + const DRAIN_STDIN: &str = r#" + (func $drain_stdin (result i32) + (local $total i32) + (loop $again + (i32.store (i32.const 0) (i32.add (i32.const 256) (local.get $total))) + (i32.store (i32.const 4) (i32.const 64)) + (if (call $fd_read (i32.const 0) (i32.const 0) (i32.const 1) (i32.const 8)) + (then (call $exit (i32.const 3)))) + (local.set $total (i32.add (local.get $total) (i32.load (i32.const 8)))) + (br_if $again (i32.load (i32.const 8)))) + local.get $total)"#; + + fn peer_program(body: &str, data: &str) -> Vec { + wat::parse_str(format!( + r#"(module + (import "wasi_snapshot_preview1" "fd_read" + (func $fd_read (param i32 i32 i32 i32) (result i32))) + (import "wasi_snapshot_preview1" "fd_write" + (func $fd_write (param i32 i32 i32 i32) (result i32))) + (import "wasi_snapshot_preview1" "proc_exit" (func $exit (param i32))) + (memory (export "memory") 1) + (data (i32.const 128) "{data}") + {DRAIN_STDIN} + (func $write (param $length i32) (result i32) + (i32.store (i32.const 16) (i32.const 128)) + (i32.store (i32.const 20) (local.get $length)) + (call $fd_write (i32.const 1) (i32.const 16) (i32.const 1) (i32.const 24))) + (func (export "_start") {body}))"# + )) + .unwrap() + } + + fn interact_pair(contestant: Vec, interactor: Vec) -> crate::InteractiveResult { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(interact(InteractiveRequest { + contestant: program(contestant), + interactor: program(interactor), + determinism: DeterminismConfig { + random_seed: 7, + realtime_epoch_ms: 946_684_800_000, + clock_step_ns: 1_000_000, + }, + })) + .unwrap() + } + + #[test] + fn interactor_writes_to_an_exited_contestant_without_a_broken_pipe() { + let contestant = peer_program("(drop (call $write (i32.const 2)))", "1\\n"); + let interactor = peer_program( + r#"(local $errno i32) + (drop (call $drain_stdin)) + (local.set $errno (call $write (i32.const 8))) + (if (local.get $errno) (then (call $exit (local.get $errno)))) + (if (i32.ne (i32.load (i32.const 24)) (i32.const 8)) (then (call $exit (i32.const 4)))) + (call $exit (i32.const 42))"#, + "correct\\n", + ); + + let result = interact_pair(contestant, interactor); + + assert_eq!(result.contestant.termination, ExecutionTermination::Exited); + assert_eq!(result.contestant.code, 0); + assert_eq!(result.interactor.termination, ExecutionTermination::Exited); + assert_eq!(result.interactor.code, 42); + assert_eq!(result.contestant_to_interactor, b"1\n"); + assert_eq!(result.interactor_to_contestant, b"correct\n"); + assert_eq!(result.interactor.metrics.protocol_bytes, 8); + } + + #[test] + fn contestant_reads_buffered_bytes_then_eof_after_the_interactor_exits() { + let contestant = peer_program( + r#"(local $total i32) + (local.set $total (call $drain_stdin)) + (call $exit (select (i32.const 0) (i32.const 1) + (i32.and + (i32.eq (local.get $total) (i32.const 4)) + (i32.eq (i32.load (i32.const 256)) (i32.const 0x0a657962)))))"#, + "", + ); + let interactor = peer_program("(drop (call $write (i32.const 4)))", "bye\\n"); + + let result = interact_pair(contestant, interactor); + + assert_eq!(result.interactor.termination, ExecutionTermination::Exited); + assert_eq!(result.interactor.code, 0); + assert_eq!(result.contestant.termination, ExecutionTermination::Exited); + assert_eq!(result.contestant.code, 0); + assert_eq!(result.interactor_to_contestant, b"bye\n"); + } + + #[test] + fn contestant_writes_to_an_exited_interactor_without_a_broken_pipe() { + let contestant = peer_program( + r#"(local $round i32) + (local $errno i32) + (drop (call $drain_stdin)) + (loop $again + (local.set $errno (call $write (i32.const 2))) + (if (local.get $errno) (then (call $exit (local.get $errno)))) + (if (i32.ne (i32.load (i32.const 24)) (i32.const 2)) (then (call $exit (i32.const 4)))) + (local.set $round (i32.add (local.get $round) (i32.const 1))) + (br_if $again (i32.lt_u (local.get $round) (i32.const 3)))) + (call $exit (i32.const 42))"#, + "x\\n", + ); + let interactor = peer_program("(drop (call $write (i32.const 4)))", "bye\\n"); + + let result = interact_pair(contestant, interactor); + + assert_eq!(result.interactor.termination, ExecutionTermination::Exited); + assert_eq!(result.interactor.code, 0); + assert_eq!(result.contestant.termination, ExecutionTermination::Exited); + assert_eq!(result.contestant.code, 42); + assert_eq!(result.contestant_to_interactor, b"x\nx\nx\n"); + assert_eq!(result.interactor_to_contestant, b"bye\n"); + } + fn program(wasm: Vec) -> InteractiveProgram { InteractiveProgram { wasm, diff --git a/crates/runtime-core/src/meter.rs b/crates/runtime-core/src/meter.rs index 6ca57a7..77db4f7 100644 --- a/crates/runtime-core/src/meter.rs +++ b/crates/runtime-core/src/meter.rs @@ -9,7 +9,7 @@ use wasm_encoder::reencode::{Error as ReencodeError, Reencode}; use wasm_encoder::{Encode, Section}; #[cfg(target_arch = "wasm32")] use wasmer::js::AsJs; -use wasmer::{AsStoreMut, Global, Instance}; +use wasmer::{AsStoreMut, Function, Global, Imports, Instance}; pub const METER_MODEL: &str = "weighted"; const METERING_MODULE: &str = "wasm_oj_metering"; @@ -88,30 +88,424 @@ impl std::fmt::Display for MeterInitializationError { impl std::error::Error for MeterInitializationError {} +/// Units a program may spend between safepoints. +pub const SAFEPOINT_INTERVAL: i64 = 1 << 20; +const SAFEPOINT_NAME: &str = "safepoint"; + +/// Final instrumentation pass: sets the meter's budget, adds the uncharged safepoint, and gives +/// every function that has a loop but no parameters or locals one unused `i32` local. +/// +/// JavaScriptCore (Safari 26) never enters optimized code inside a loop of a function without +/// parameters or locals: its baseline tier asks to tier up on almost every iteration and keeps +/// running the slow path, so a metered empty loop runs about 20 times slower than in Chromium and +/// hits the wall deadline before its instruction budget. The local changes neither behaviour nor +/// cost. +/// +/// JavaScriptCore also never stops a terminated Worker while it runs Wasm, only when it next +/// reaches a JavaScript checkpoint such as `Atomics.wait`. So the charge that opens every function +/// and loop body is inlined, and once the counter drops below a threshold it reaches a cold +/// function that traps on exhaustion, or calls the imported `wasm_oj_metering.safepoint` (which +/// the browser host turns into such a checkpoint) and lowers the threshold by +/// `SAFEPOINT_INTERVAL`. A call inside a hot loop makes the whole loop slower in JavaScriptCore and +/// V8 even when it never runs, so a loop body leaves the loop to reach it: each loop becomes +/// `block (loop (block (loop ...) br 2) safepoint, refund, br 0)`, and every branch out of it +/// skips the three new labels. Functions that use other kinds of labels keep the call inside the +/// loop. Every other block keeps radix's gas function. Exhaustion is exactly where it was, and +/// none of these instructions is charged. struct MeterInitializer { budget: i64, + parameterized_types: Vec, + function_types: Vec, + defined_functions: usize, + original_types: u32, + imported_functions: u32, + imported_globals: u32, + defined_globals: u32, + import_section_written: bool, +} + +impl MeterInitializer { + fn new(budget: i64) -> Self { + Self { + budget, + parameterized_types: Vec::new(), + function_types: Vec::new(), + defined_functions: 0, + original_types: 0, + imported_functions: 0, + imported_globals: 0, + defined_globals: 0, + import_section_written: false, + } + } + + fn safepoint_type(&self) -> u32 { + self.original_types + } + + fn safepoint_function(&self) -> u32 { + self.imported_functions + } + + fn gas_function(&self) -> u32 { + self.imported_functions + self.function_types.len() as u32 + } + + fn slow_function(&self) -> u32 { + self.gas_function() + 1 + } + + fn gas_global(&self) -> u32 { + self.imported_globals + self.defined_globals - 1 + } + + fn threshold_global(&self) -> u32 { + self.imported_globals + self.defined_globals + } + + fn safepoint_import(&self, imports: &mut wasm_encoder::ImportSection) { + imports.import( + METERING_MODULE, + SAFEPOINT_NAME, + wasm_encoder::EntityType::Function(self.safepoint_type()), + ); + } + + fn slow_body(&self) -> wasm_encoder::Function { + use wasm_encoder::Instruction::*; + let gas = self.gas_global(); + let mut function = wasm_encoder::Function::new([(1, wasm_encoder::ValType::I64)]); + for instruction in [ + GlobalGet(gas), + I64Const(0), + I64LtS, + If(wasm_encoder::BlockType::Empty), + I64Const(-1), + GlobalSet(gas), + Unreachable, + End, + GlobalGet(gas), + I64Const(SAFEPOINT_INTERVAL), + I64Sub, + LocalTee(0), + I64Const(0), + LocalGet(0), + I64Const(0), + I64GtS, + Select, + GlobalSet(self.threshold_global()), + Call(self.safepoint_function()), + End, + ] { + function.instruction(&instruction); + } + function + } + + /// Re-encodes a body, replacing the gas call that opens the function and each loop body with + /// an inline charge and safepoint check. + fn encode_body( + &mut self, + body: &wasmparser::FunctionBody<'_>, + extra_local: bool, + ) -> Result> { + use wasm_encoder::Instruction::*; + use wasmparser::Operator as Op; + let mut function = self.new_function_with_parsed_locals(body)?; + if extra_local { + function = wasm_encoder::Function::new([(1, wasm_encoder::ValType::I32)]); + } + let original_gas = self.imported_functions + self.function_types.len() as u32 - 1; + let gas = self.gas_global(); + let rotate = !has_other_labels(body)?; + let mut frames: Vec> = Vec::new(); + let mut head = Some(false); + let mut pending = None; + let adjust = |frames: &[Option], depth: u32| { + let crossed = frames[frames.len() - depth as usize..] + .iter() + .filter(|frame| frame.is_some()) + .count(); + depth + 3 * crossed as u32 + }; + let mut operators = body.get_operators_reader()?; + while !operators.eof() { + let operator = operators.read()?; + if let Some((cost, wrapped)) = pending.take() { + if matches!(operator, Op::Call { function_index } if function_index == original_gas) + { + for instruction in [ + GlobalGet(gas), + I64Const(cost), + I64Sub, + GlobalSet(gas), + GlobalGet(gas), + GlobalGet(self.threshold_global()), + I64LtS, + ] { + function.instruction(&instruction); + } + if wrapped { + function.instruction(&BrIf(1)); + if let Some(frame) = frames.last_mut() { + *frame = Some(cost); + } + } else { + function.instruction(&If(wasm_encoder::BlockType::Empty)); + function.instruction(&Call(self.slow_function())); + function.instruction(&End); + } + continue; + } + function.instruction(&I64Const(cost)); + } + let at_head = head.take(); + match operator { + Op::I64Const { value } if at_head.is_some() => { + pending = Some((value, at_head == Some(true))); + continue; + } + Op::Loop { blockty } => { + let parameterized = match blockty { + wasmparser::BlockType::FuncType(ty) => self + .parameterized_types + .get(ty as usize) + .copied() + .unwrap_or(true), + _ => false, + }; + let wrapped = rotate && !parameterized; + let blockty = self.block_type(blockty)?; + if wrapped { + function.instruction(&Block(blockty)); + function.instruction(&Loop(wasm_encoder::BlockType::Empty)); + function.instruction(&Block(wasm_encoder::BlockType::Empty)); + frames.push(Some(0)); + } else { + frames.push(None); + } + function.instruction(&Loop(blockty)); + head = Some(wrapped); + continue; + } + Op::Block { .. } | Op::If { .. } | Op::Try { .. } | Op::TryTable { .. } => { + frames.push(None) + } + Op::End => { + if let Some(Some(cost)) = frames.pop() { + for instruction in [ + End, + Br(2), + End, + Call(self.slow_function()), + GlobalGet(gas), + I64Const(cost), + I64Add, + GlobalSet(gas), + Br(0), + End, + Unreachable, + End, + ] { + function.instruction(&instruction); + } + continue; + } + } + Op::Br { relative_depth } if rotate => { + function.instruction(&Br(adjust(&frames, relative_depth))); + continue; + } + Op::BrIf { relative_depth } if rotate => { + function.instruction(&BrIf(adjust(&frames, relative_depth))); + continue; + } + Op::BrTable { targets } if rotate => { + let mut depths = Vec::new(); + for target in targets.targets() { + depths.push(adjust(&frames, target?)); + } + function + .instruction(&BrTable(depths.into(), adjust(&frames, targets.default()))); + continue; + } + _ => {} + } + function.instruction(&self.instruction(operator)?); + } + Ok(function) + } +} + +/// Whether a body uses labels other than through `br`, `br_if` and `br_table`. +fn has_other_labels( + body: &wasmparser::FunctionBody<'_>, +) -> Result { + use wasmparser::Operator as Op; + let mut operators = body.get_operators_reader()?; + while !operators.eof() { + if matches!( + operators.read()?, + Op::Try { .. } + | Op::TryTable { .. } + | Op::Delegate { .. } + | Op::Rethrow { .. } + | Op::BrOnNull { .. } + | Op::BrOnNonNull { .. } + | Op::BrOnCast { .. } + | Op::BrOnCastFail { .. } + | Op::BrOnCastDescEq { .. } + | Op::BrOnCastDescEqFail { .. } + ) { + return Ok(true); + } + } + Ok(false) +} + +fn meter_error(message: &str) -> ReencodeError { + ReencodeError::UserError(MeterInitializationError(message.to_string())) } impl Reencode for MeterInitializer { type Error = MeterInitializationError; + fn function_index(&mut self, func: u32) -> Result> { + Ok(if func >= self.imported_functions { + func + 1 + } else { + func + }) + } + + fn intersperse_section_hook( + &mut self, + module: &mut wasm_encoder::Module, + _after: Option, + before: Option, + ) -> Result<(), ReencodeError> { + use wasm_encoder::SectionId; + if self.import_section_written + || matches!( + before, + Some(SectionId::Custom | SectionId::Type | SectionId::Import) + ) + { + return Ok(()); + } + let mut imports = wasm_encoder::ImportSection::new(); + self.safepoint_import(&mut imports); + module.section(&imports); + self.import_section_written = true; + Ok(()) + } + + fn parse_type_section( + &mut self, + types: &mut wasm_encoder::TypeSection, + section: wasmparser::TypeSectionReader<'_>, + ) -> Result<(), ReencodeError> { + for group in section.clone() { + for ty in group?.types() { + self.parameterized_types + .push(match &ty.composite_type.inner { + wasmparser::CompositeInnerType::Func(function) => { + !function.params().is_empty() + } + _ => true, + }); + } + } + self.original_types = self.parameterized_types.len() as u32; + wasm_encoder::reencode::utils::parse_type_section(self, types, section)?; + types.ty().function([], []); + Ok(()) + } + + fn parse_import_section( + &mut self, + imports: &mut wasm_encoder::ImportSection, + section: wasmparser::ImportSectionReader<'_>, + ) -> Result<(), ReencodeError> { + for import in section.clone().into_imports() { + match import?.ty { + wasmparser::TypeRef::Func(_) => self.imported_functions += 1, + wasmparser::TypeRef::Global(_) => self.imported_globals += 1, + _ => {} + } + } + wasm_encoder::reencode::utils::parse_import_section(self, imports, section)?; + self.safepoint_import(imports); + self.import_section_written = true; + Ok(()) + } + + fn parse_function_section( + &mut self, + functions: &mut wasm_encoder::FunctionSection, + section: wasmparser::FunctionSectionReader<'_>, + ) -> Result<(), ReencodeError> { + for function in section.clone() { + self.function_types.push(function?); + } + if self.function_types.is_empty() { + return Err(meter_error("instrumented module has no gas function")); + } + wasm_encoder::reencode::utils::parse_function_section(self, functions, section)?; + functions.function(self.safepoint_type()); + Ok(()) + } + + fn parse_code_section( + &mut self, + code: &mut wasm_encoder::CodeSection, + section: wasmparser::CodeSectionReader<'_>, + ) -> Result<(), ReencodeError> { + wasm_encoder::reencode::utils::parse_code_section(self, code, section)?; + code.function(&self.slow_body()); + Ok(()) + } + + fn parse_function_body( + &mut self, + code: &mut wasm_encoder::CodeSection, + body: wasmparser::FunctionBody<'_>, + ) -> Result<(), ReencodeError> { + let ordinal = self.defined_functions; + self.defined_functions += 1; + if ordinal + 1 == self.function_types.len() { + return wasm_encoder::reencode::utils::parse_function_body(self, code, body); + } + let parameterized = self + .function_types + .get(ordinal) + .and_then(|ty| self.parameterized_types.get(*ty as usize)) + .copied() + .unwrap_or(true); + let loops = has_loop(&body)?; + let extra_local = !parameterized && body.get_locals_reader()?.get_count() == 0 && loops; + let function = self.encode_body(&body, extra_local)?; + code.function(&function); + Ok(()) + } + fn parse_global_section( &mut self, globals: &mut wasm_encoder::GlobalSection, section: wasmparser::GlobalSectionReader<'_>, ) -> Result<(), ReencodeError> { - let meter_ordinal = section.count().checked_sub(1).ok_or_else(|| { - ReencodeError::UserError(MeterInitializationError( - "instrumented module has no meter global".to_string(), - )) - })?; + self.defined_globals = section.count(); + let meter_ordinal = section + .count() + .checked_sub(1) + .ok_or_else(|| meter_error("instrumented module has no meter global"))?; for (ordinal, global) in section.into_iter().enumerate() { let global = global?; if u32::try_from(ordinal).ok() == Some(meter_ordinal) { if global.ty.content_type != wasmparser::ValType::I64 || !global.ty.mutable { - return Err(ReencodeError::UserError(MeterInitializationError( - "instrumented meter global has an unexpected type".to_string(), - ))); + return Err(meter_error( + "instrumented meter global has an unexpected type", + )); } globals.global( self.global_type(global.ty)?, @@ -121,6 +515,14 @@ impl Reencode for MeterInitializer { wasm_encoder::reencode::utils::parse_global(self, globals, global)?; } } + globals.global( + wasm_encoder::GlobalType { + val_type: wasm_encoder::ValType::I64, + mutable: true, + shared: false, + }, + &wasm_encoder::ConstExpr::i64_const((self.budget - SAFEPOINT_INTERVAL).max(0)), + ); Ok(()) } } @@ -128,12 +530,22 @@ impl Reencode for MeterInitializer { fn set_initial_meter_budget(wasm: &[u8], budget: i64) -> Result, String> { validate_meter_global_position(wasm)?; let mut module = wasm_encoder::Module::new(); - MeterInitializer { budget } + MeterInitializer::new(budget) .parse_core_module(&mut module, wasmparser::Parser::new(0), wasm) .map_err(|error| format!("failed to initialize weighted meter: {error}"))?; Ok(module.finish()) } +fn has_loop(body: &wasmparser::FunctionBody<'_>) -> Result { + let mut operators = body.get_operators_reader()?; + while !operators.eof() { + if matches!(operators.read()?, wasmparser::Operator::Loop { .. }) { + return Ok(true); + } + } + Ok(false) +} + fn validate_meter_global_position(wasm: &[u8]) -> Result<(), String> { let mut imported_globals = 0_u32; let mut defined_globals = 0_u32; @@ -271,6 +683,36 @@ pub fn remaining_points( } } +/// Provides `wasm_oj_metering.safepoint`, which the meter calls every `SAFEPOINT_INTERVAL` units. +/// Natively it does nothing. In a browser it calls `Atomics.wait` with a value that never matches, +/// which returns at once; JavaScriptCore checks for Worker termination there, which it never does +/// inside Wasm, so `terminate()` stops a running program within one interval in WebKit too. +pub fn attach_safepoint(store: &mut impl AsStoreMut, imports: &mut Imports) { + imports.define( + METERING_MODULE, + SAFEPOINT_NAME, + Function::new_typed(store, safepoint), + ); +} + +#[cfg(not(target_arch = "wasm32"))] +fn safepoint() {} + +#[cfg(target_arch = "wasm32")] +fn safepoint() { + thread_local! { + static CHECKPOINT: Option = js_sys::Reflect::get(&js_sys::global(), &"SharedArrayBuffer".into()) + .ok() + .filter(|constructor| constructor.is_function()) + .map(|_| js_sys::Int32Array::new(&js_sys::SharedArrayBuffer::new(4))); + } + CHECKPOINT.with(|checkpoint| { + if let Some(checkpoint) = checkpoint { + let _ = js_sys::Atomics::wait_with_timeout(checkpoint, 0, 1, 0.0); + } + }); +} + fn weighted_instruction_cost(operator: &Operator) -> Option { let debug = format!("{operator:?}"); weighted_opcode_cost(debug.split_whitespace().next().unwrap_or("UNKNOWN")) @@ -406,10 +848,12 @@ fn wark_v03_opcode_cost(opcode: &str) -> u32 { #[cfg(test)] mod tests { - use super::{METER_MODEL, instrument_wasm, weighted_opcode_cost}; + use super::{METER_MODEL, SAFEPOINT_INTERVAL, instrument_wasm, weighted_opcode_cost}; + use crate::{DeterminismConfig, ExecutionTermination, ResourcePolicy, RunRequest}; use std::borrow::Cow; + use std::collections::BTreeMap; use wasm_encoder::{Encode, Section}; - use wasmparser::{ExternalKind, Parser, Payload}; + use wasmparser::{ExternalKind, Parser, Payload, ValType}; #[test] fn instrumentation_adds_the_metering_global() { @@ -488,6 +932,48 @@ mod tests { assert_eq!(weighted_opcode_cost("AtomicFence"), Some(1000)); } + #[test] + fn functions_with_a_loop_and_no_params_or_locals_get_one_unused_local() { + let wasm = wat::parse_str( + r#"(module + (memory (export "memory") 1) + (func (export "_start") (loop (br 0))) + (func (local i64) (loop (br 0))) + (func (param i32) (loop (br 0))) + (func nop) + (func (block (loop (br 1)))))"#, + ) + .unwrap(); + let metered = instrument_wasm(&wasm, 1_000_000).unwrap(); + let locals = Parser::new(0) + .parse_all(&metered.wasm) + .filter_map(Result::ok) + .filter_map(|payload| match payload { + Payload::CodeSectionEntry(body) => Some( + body.get_locals_reader() + .unwrap() + .into_iter() + .map(|local| local.unwrap()) + .collect::>(), + ), + _ => None, + }) + .collect::>(); + let i32_local = vec![(1, ValType::I32)]; + assert_eq!( + locals, + [ + i32_local.clone(), + vec![(1, ValType::I64)], + vec![], + vec![], + i32_local, + vec![], + vec![(1, ValType::I64)], + ] + ); + } + #[test] fn reports_wark_compatible_static_operation_counts() { let wasm = wat::parse_str( @@ -499,4 +985,235 @@ mod tests { assert_eq!(metered.operations.get("I32Add"), Some(&1)); assert_eq!(metered.operations.get("Drop"), Some(&1)); } + + fn looping(iterations: u32) -> Vec { + wat::parse_str(format!( + r#"(module + (import "wasi_snapshot_preview1" "proc_exit" (func $exit (param i32))) + (memory (export "memory") 1) + (func (export "_start") + (local $remaining i32) + i32.const {iterations} local.set $remaining + (loop $again + local.get $remaining i32.const 1 i32.sub local.tee $remaining + br_if $again)))"# + )) + .unwrap() + } + + fn run(wasm: &[u8], instruction_budget: u64) -> crate::RunResult { + crate::run(RunRequest { + wasm: wasm.to_vec(), + args: Vec::new(), + env: BTreeMap::new(), + stdin: Vec::new(), + files: BTreeMap::new(), + output_paths: Vec::new(), + cwd: Some("/".to_string()), + startup_entropy_bytes: 0, + determinism: DeterminismConfig { + random_seed: 7, + realtime_epoch_ms: 946_684_800_000, + clock_step_ns: 1_000_000, + }, + resources: ResourcePolicy { + instruction_budget, + logical_time_limit_ms: 60_000, + memory_limit_bytes: 64 * 1024 * 1024, + output_limit_bytes: 1024, + filesystem_write_limit_bytes: 64 * 1024 * 1024, + filesystem_entry_limit: 4_096, + }, + }) + .unwrap() + } + + #[test] + fn costs_and_exhaustion_do_not_depend_on_safepoint_intervals() { + let wasm = looping(3_000_000); + let generous = run(&wasm, 10_000_000_000); + assert_eq!(generous.termination, ExecutionTermination::Exited); + let cost = generous.metrics.cost; + assert!(cost > 5 * SAFEPOINT_INTERVAL as u64); + assert_eq!(run(&wasm, cost * 7).metrics.cost, cost); + + let exact = run(&wasm, cost); + assert_eq!(exact.termination, ExecutionTermination::Exited); + assert_eq!(exact.metrics.cost, cost); + + let short = run(&wasm, cost - 1); + assert_eq!(short.termination, ExecutionTermination::InstructionLimit); + assert_eq!(short.metrics.cost, cost - 1); + + let within_one_interval = looping(1_000); + let small = run(&within_one_interval, 1_000_000).metrics.cost; + assert_eq!(run(&within_one_interval, small).metrics.cost, small); + assert_eq!( + run(&within_one_interval, small - 1).termination, + ExecutionTermination::InstructionLimit + ); + } + + #[test] + fn the_meter_calls_the_safepoint_once_per_interval() { + let recursing = wat::parse_str( + r#"(module + (import "wasi_snapshot_preview1" "proc_exit" (func $exit (param i32))) + (memory (export "memory") 1) + (func $fib (param $n i32) (result i32) + local.get $n i32.const 2 i32.lt_u + if (result i32) local.get $n + else + local.get $n i32.const 1 i32.sub call $fib + local.get $n i32.const 2 i32.sub call $fib + i32.add + end) + (func (export "_start") i32.const 25 call $fib drop))"#, + ) + .unwrap(); + for wasm in [looping(3_000_000), recursing] { + let (calls, cost) = count_safepoints(&wasm); + assert!(cost > 5 * SAFEPOINT_INTERVAL as u64); + assert!(calls.abs_diff(cost / SAFEPOINT_INTERVAL as u64) <= 1); + } + } + + #[test] + fn safepoint_loops_keep_their_branches_and_results() { + use wasmer::{Function, Imports, Instance, Module, Store, Value}; + + let wasm = wat::parse_str( + r#"(module + (memory (export "memory") 1) + (func (export "compute") (result i32) + (local $i i32) (local $j i32) (local $acc i32) + (block $done + (loop $outer + (local.set $j (i32.const 0)) + (block $inner_done + (loop $inner + (local.set $acc + (i32.add (local.get $acc) (i32.mul (local.get $i) (local.get $j)))) + (local.set $j (i32.add (local.get $j) (i32.const 1))) + (br_if $done (i32.eq (local.get $acc) (i32.const -1))) + (br_table $inner $inner_done $done + (i32.gt_u (local.get $j) (i32.const 500))))) + (local.set $i (i32.add (local.get $i) (i32.const 1))) + (br_if $outer (i32.lt_u (local.get $i) (i32.const 3000))))) + (local.get $acc) + (block $value (result i32) + (loop $early (result i32) (br $value (i32.const 3)))) + (i32.add) + (loop $fallthrough (result i32) (i32.const 7)) + (i32.add) + (i32.const 5) + (loop $param (param i32) (result i32) + i32.const 1 + i32.sub + local.tee $j + local.get $j + br_if $param) + (i32.add)))"#, + ) + .unwrap(); + let mut results = Vec::new(); + for module in [ + wasm.clone(), + instrument_wasm(&wasm, 10_000_000_000).unwrap().wasm, + ] { + let mut store = Store::default(); + let module = Module::new(&store, &module).unwrap(); + let mut imports = Imports::new(); + imports.define( + "wasm_oj_metering", + "safepoint", + Function::new_typed(&mut store, || {}), + ); + let instance = Instance::new(&mut store, &module, &imports).unwrap(); + let compute = instance.exports.get_function("compute").unwrap(); + results.push(compute.call(&mut store, &[]).unwrap().to_vec()); + if let Ok(state) = super::meter_state(&instance) { + let super::CostPoints::Remaining(remaining) = + super::remaining_points(&mut store, &state).unwrap() + else { + panic!("meter exhausted"); + }; + assert!(10_000_000_000 - remaining > 5 * SAFEPOINT_INTERVAL as u64); + } + } + assert_eq!(results[0], results[1]); + assert!(matches!(results[0][..], [Value::I32(_)])); + } + + fn count_safepoints(wasm: &[u8]) -> (u64, u64) { + use std::sync::Arc; + use std::sync::atomic::{AtomicU64, Ordering}; + use wasmer::{Function, FunctionEnv, FunctionEnvMut, Imports, Instance, Module, Store}; + + let budget = 10_000_000_000; + let metered = instrument_wasm(wasm, budget).unwrap(); + let mut store = Store::default(); + let module = Module::new(&store, &metered.wasm).unwrap(); + let calls = Arc::new(AtomicU64::new(0)); + let env = FunctionEnv::new(&mut store, calls.clone()); + let mut imports = Imports::new(); + imports.define( + "wasm_oj_metering", + "safepoint", + Function::new_typed_with_env( + &mut store, + &env, + |env: FunctionEnvMut>| { + env.data().fetch_add(1, Ordering::Relaxed); + }, + ), + ); + imports.define( + "wasi_snapshot_preview1", + "proc_exit", + Function::new_typed(&mut store, |_: i32| {}), + ); + let instance = Instance::new(&mut store, &module, &imports).unwrap(); + instance + .exports + .get_function("_start") + .unwrap() + .call(&mut store, &[]) + .unwrap(); + let state = super::meter_state(&instance).unwrap(); + let super::CostPoints::Remaining(remaining) = + super::remaining_points(&mut store, &state).unwrap() + else { + panic!("meter exhausted"); + }; + (calls.load(Ordering::Relaxed), budget - remaining) + } + + #[test] + fn modules_without_imports_receive_the_safepoint_import() { + let wasm = + wat::parse_str("(module (memory (export \"memory\") 1) (func (export \"_start\")))") + .unwrap(); + let metered = instrument_wasm(&wasm, 1_000_000).unwrap(); + let imports = Parser::new(0) + .parse_all(&metered.wasm) + .filter_map(Result::ok) + .filter_map(|payload| match payload { + Payload::ImportSection(section) => Some( + section + .into_imports() + .map(|import| { + let import = import.unwrap(); + format!("{}.{}", import.module, import.name) + }) + .collect::>(), + ), + _ => None, + }) + .collect::>(); + assert_eq!(imports, [vec!["wasm_oj_metering.safepoint".to_string()]]); + wasmparser::Validator::new() + .validate_all(&metered.wasm) + .unwrap(); + } } diff --git a/crates/runtime-core/src/run/native.rs b/crates/runtime-core/src/run/native.rs index 2267ae5..e1a431c 100644 --- a/crates/runtime-core/src/run/native.rs +++ b/crates/runtime-core/src/run/native.rs @@ -2,7 +2,9 @@ use crate::capabilities::attach_capability_denials; use crate::deterministic::{VirtualClock, attach_deterministic_imports}; use crate::filesystem::{read_files_bounded, runtime_project_files}; use crate::memory::LimitingTunables; -use crate::meter::{CostPoints, METER_MODEL, instrument_wasm, meter_state, remaining_points}; +use crate::meter::{ + CostPoints, METER_MODEL, attach_safepoint, instrument_wasm, meter_state, remaining_points, +}; use crate::module_imports::attach_imported_memory; use crate::module_policy::{DEFERRED_START_EXPORT, defer_start_section, enforce_memory_limit}; use crate::output::{OutputBudget, OutputCapture}; @@ -105,6 +107,7 @@ fn run_in_runtime(request: RunRequest) -> Result { clock.clone(), request.startup_entropy_bytes, ); + attach_safepoint(&mut store, &mut imports); attach_capability_denials(&mut store, &module, &mut imports).map_err(RunError::Compile)?; let imported_memory = attach_imported_memory(&mut store, &module, &mut imports).map_err(RunError::Compile)?; diff --git a/crates/runtime-core/src/run/web.rs b/crates/runtime-core/src/run/web.rs index 7cca6ea..54396c0 100644 --- a/crates/runtime-core/src/run/web.rs +++ b/crates/runtime-core/src/run/web.rs @@ -2,7 +2,9 @@ use super::web_runtime::runtime_with_engine; use crate::capabilities::attach_capability_denials; use crate::deterministic::{VirtualClock, attach_deterministic_imports}; use crate::filesystem::{read_files_bounded, runtime_project_files}; -use crate::meter::{CostPoints, METER_MODEL, instrument_wasm, meter_state, remaining_points}; +use crate::meter::{ + CostPoints, METER_MODEL, attach_safepoint, instrument_wasm, meter_state, remaining_points, +}; use crate::module_imports::attach_imported_memory; use crate::module_policy::{DEFERRED_START_EXPORT, defer_start_section, enforce_memory_limit}; use crate::output::{CappedOutput, OutputBudget, OutputCapture}; @@ -102,6 +104,7 @@ pub(super) fn execute( clock.clone(), request.startup_entropy_bytes, ); + attach_safepoint(&mut store, &mut imports); attach_capability_denials(&mut store, &module, &mut imports).map_err(RunError::Compile)?; let imported_memory = attach_imported_memory(&mut store, &module, &mut imports).map_err(RunError::Compile)?; diff --git a/crates/runtime-core/src/run/web_interactive.rs b/crates/runtime-core/src/run/web_interactive.rs index c16f73b..6820922 100644 --- a/crates/runtime-core/src/run/web_interactive.rs +++ b/crates/runtime-core/src/run/web_interactive.rs @@ -291,13 +291,18 @@ impl std::fmt::Debug for StreamOutput { } impl AsyncWrite for StreamOutput { + /// Drops bytes written after the peer closed its stdin, as native does; they are already in + /// the transcript. fn poll_write( mut self: Pin<&mut Self>, context: &mut Context<'_>, buffer: &[u8], ) -> Poll> { match Pin::new(&mut self.capture).poll_write(context, buffer) { - Poll::Ready(Ok(written)) => Poll::Ready(self.streams.write(&buffer[..written])), + Poll::Ready(Ok(written)) => Poll::Ready(match self.streams.write(&buffer[..written]) { + Err(error) if error.kind() == io::ErrorKind::BrokenPipe => Ok(written), + result => result, + }), result => result, } } diff --git a/docs/architecture.md b/docs/architecture.md index 9406b48..b383e54 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -57,7 +57,7 @@ complete browser Worker-generation boundary. Browser interaction runs each side Worker as a standalone metered run. The two sides exchange bytes through shared-memory ring buffers whose reads block with `Atomics.wait`, so neither side ever yields to the other on one thread. Each ring holds its writer's whole output budget, so a write never waits, as on the server's unbounded -pipes. A poll checks the input without blocking, so another ready subscription is reported first; if +pipes. On both hosts a write after the reader exited is dropped instead of failing. A poll checks the input without blocking, so another ready subscription is reported first; if nothing is ready, a poll with a clock advances the virtual clock to its deadline, as the server does. On both hosts a read from stdin opened with `O_NONBLOCK` waits for input instead of failing with `EAGAIN`. @@ -154,8 +154,16 @@ or toolchain source. Before instantiation the runtime validates the module, removes non-semantic debug/name sections, preserves required runtime metadata, and injects a mutable 64-bit weighted instruction meter. The -budget is present before a start section can execute. Static original-opcode counts and normalized -cost are reported separately from injected meter instructions. +budget is present before a start section can execute. A function that has a loop but no parameters +or locals also gets one unused local: JavaScriptCore never optimizes such a loop and runs it about 20 +times slower than Chromium, so it would reach the wall deadline before its budget. Each function entry and +loop iteration also compares the counter with a threshold. Once it drops 2^20 units below the last +safepoint, a cold function calls the imported `wasm_oj_metering.safepoint`, which does nothing +natively and in browsers performs an `Atomics.wait` that returns at once, a point where +JavaScriptCore acts on `Worker.terminate()`. Loops reach it through a branch out of the loop, because +a call inside a hot loop slows the whole loop in JavaScriptCore and V8 even when it never runs. The +check and the safepoint are not charged, so costs and the exhaustion point are unchanged. Static +original-opcode counts and normalized cost are reported separately from injected meter instructions. Contract 2 enforces: diff --git a/docs/library-contract.md b/docs/library-contract.md index 214a5fa..8fe647a 100644 --- a/docs/library-contract.md +++ b/docs/library-contract.md @@ -125,6 +125,17 @@ nested stages with bounded lifetime. Python and JavaScript package source files crossing a stage budget, cancellation, timeout, restart, cache clearing, disposal, or infrastructure failure establishes a complete Worker-generation boundary. +A module Worker that dies without an `error` event, for example because the browser terminated it, +is reported like a crash: its operation rejects with a `runner-failure` or `compiler-failure` +instead of waiting for its wall-time or build deadline, so a killed Worker is never reported as a +time limit. Each Worker holds a Web Lock for its lifetime, and its owner learns of the death when +that lock frees. Without Web Locks the previous behaviour remains. WebKit stops a terminated Worker +only at JavaScript checkpoints, never inside Wasm, so a metered program reaches the runtime's +safepoint every 2^20 instruction units (see the runtime policy in the architecture guide) and a +killed Worker running one stops within that interval. Chromium reports a busy Worker about 2 s after +the kill. Unmetered toolchain code, such as clang or rustc running in a compiler stage, has no +safepoint, so in WebKit such a stage is reported only when it next returns to JavaScript. + Wasmer secondary Workers are host implementation details. They use the SDK's supported `workerUrl` protocol and do not grant guest thread-spawn capability. The host page must be cross-origin isolated. @@ -244,6 +255,15 @@ resource policies, process-local deterministic clocks, and secret inputs mounted interactor side. Either side may be a standalone Wasm module or a runtime bundle such as CPython; runtime bundles that cannot provide streaming fd 0 are rejected for interaction. +A side that has exited, or closed its stdin, no longer reads, but its peer's writes to it still +succeed: the bytes are dropped and the peer keeps running without `EPIPE`, as with a judge that keeps +draining both programs' output until they exit. Reads from a side that has exited, or closed its +stdout, return the bytes still buffered and then EOF. An interactor can therefore reply to a +contestant that already exited, read EOF, and exit with its own verdict. `contestantToInteractor` and +`interactorToContestant` record every byte each side wrote to stdout exactly once, up to its output +limit, including bytes written after the peer exited. A transcript depends only on what its writer +wrote, not on when the reader exited. + Each case executes under the broad hard policy once. Correct output and the same normalized metrics are evaluated against ordered cumulative `baseline`, `efficient`, and `optimal` policies. The portable compute metric is `RunResult.metrics.cost`; wall time is only a safety boundary. diff --git a/scripts/verify-browser-csp.mjs b/scripts/verify-browser-csp.mjs index 5f9a545..1a1be3e 100644 --- a/scripts/verify-browser-csp.mjs +++ b/scripts/verify-browser-csp.mjs @@ -27,6 +27,9 @@ const bootstrap = `import { createBrowserEngine, WASM_OJ_LIBCXX_PCH_HEADER } fro window.cspViolations = []; addEventListener('securitypolicyviolation', e => window.cspViolations.push({directive:e.effectiveDirective, blockedURI:e.blockedURI, source:e.sourceFile})); try { new Function('return 1')(); window.evalBlocked = false; } catch { window.evalBlocked = true; } window.header = WASM_OJ_LIBCXX_PCH_HEADER; +const NativeWorker = Worker; window.createdWorkers = []; +window.Worker = class extends NativeWorker { constructor(url, options) { super(url, options); window.createdWorkers.push({ name:options?.name, worker:this }); } }; +window.killWorker = name => NativeWorker.prototype.terminate.call(window.createdWorkers.findLast(entry => entry.name === name).worker); window.engine = await createBrowserEngine({ artifactCache:false, toolchains: ${JSON.stringify(sources)} }); window.ready = true;`; const server = createServer(async (req, res) => { @@ -90,12 +93,18 @@ fixtures.push( { language:"c", label:"logical-clock-limit", source:'#include \n#include \nint main(){for(int i=0;i<10000;i++) clock();puts("42");}', input:"", termination:"logical-time-limit", resources:{logicalTimeLimitMs:1} }, { language:"c", label:"cpu-work", source:'#include \nint main(){volatile unsigned long long s=0;for(unsigned i=0;i<10000000;i++)s+=i;printf("%llu\\n",s);}', input:"", expected:"49999995000000\n" }, { language:"python", label:"memory-16mb", source:'print(42)', input:"", expected:"42\n", resources:{memoryLimitBytes:16*1024*1024} }, + { language:"cpp", label:"empty-loop-budget", source:'int main(){for(;;);}', input:"", termination:"instruction-limit", resources:{wallTimeLimitMs:15000} }, ); const guessInteractor = { language:"cpp", source:'#include \nint main(int argc,char**argv){std::FILE*f=argc>1?std::fopen(argv[1],"r"):nullptr;long long secret,guess;int limit;if(!f||std::fscanf(f,"%lld %d",&secret,&limit)!=2)return 3;for(int round=0;round");std::fflush(stdout);}return 1;}' }; const guessInput = secret => ({ args:["/judge/input.txt"], files:{"/judge/input.txt":`${secret} 25\n`} }); const guessCpp = { language:"cpp", source:'#include \n#include \nint main(){long long lo=1,hi=1<<20;std::string r;while(lo<=hi){long long mid=(lo+hi)/2;std::cout<>r))return 2;if(r=="=")return 0;if(r=="<")lo=mid+1;else hi=mid-1;}return 1;}' }; const guessC = { language:"c", source:'#include \nint main(void){long long lo=1,hi=1<<20;char r[4];while(lo<=hi){long long mid=(lo+hi)/2;printf("%lld\\n",mid);fflush(stdout);if(scanf("%3s",r)!=1)return 2;if(r[0]==\'=\')return 0;if(r[0]==\'<\')lo=mid+1;else hi=mid-1;}return 1;}' }; const readOne = { language:"c", source:'#include \nint main(void){int x;return scanf("%d",&x)==1?0:1;}' }; +const finalGuess = { language:"c", source:'#include \nint main(void){puts("7");return 0;}' }; +const secretInput = secret => ({ args:["/judge/input.txt"], files:{"/judge/input.txt":`${secret}\n`} }); +const replyInteractor = (language, afterEof) => language === "c" + ? { language, source:`#include \nint main(int argc,char**argv){FILE*f=argc>1?fopen(argv[1],"r"):NULL;long long secret,guess;if(!f||fscanf(f,"%lld",&secret)!=1)return 2;if(scanf("%lld",&guess)!=1)return 3;${afterEof ? "while(getchar()!=EOF){}" : ""}puts(guess==secret?"correct":"wrong");if(fflush(stdout)!=0)return 4;${afterEof ? "" : "if(scanf(\"%lld\",&guess)!=EOF)return 5;"}return guess==secret?42:43;}` } + : { language, source:`import sys\nsecret = int(open(sys.argv[1]).read())\nguess = int(input())\n${afterEof ? "sys.stdin.read()\n" : ""}print("correct" if guess == secret else "wrong", flush=True)\n${afterEof ? "" : "if sys.stdin.read().strip():\n sys.exit(5)\n"}sys.exit(42 if guess == secret else 43)\n` }; const exited = (process, code) => process.termination === "exited" && process.code === code; const guessed = result => exited(result.contestant, 0) && exited(result.interactor, 0) && result.interactorToContestant.endsWith("=\n"); const interactiveFixtures = [ @@ -110,8 +119,11 @@ const interactiveFixtures = [ { label:"interactive-batch", contestant:{ language:"c", source:'#include \nint main(void){long long x;for(int i=0;i<30000;i++){if(scanf("%lld",&x)!=1)return 2;printf("%lld\\n",2*x);}fflush(stdout);char b[8];if(scanf("%7s",b)!=1)return 3;return b[0]==\'o\'?0:4;}' }, interactor:{ language:"c", source:'#include \nint main(void){for(int i=0;i<30000;i++)printf("%d\\n",1000000+i);fflush(stdout);for(int i=0;i<30000;i++){long long x;if(scanf("%lld",&x)!=1)return 2;if(x!=2LL*(1000000+i))return 1;}puts("ok");fflush(stdout);return 0;}' }, options:{ contestant:{ resources:{ wallTimeLimitMs:20000 } }, interactor:{ resources:{ wallTimeLimitMs:20000 } } }, check:result => exited(result.contestant, 0) && exited(result.interactor, 0) && result.contestantToInteractor.length === 240000 && result.interactorToContestant.length === 240003 && result.interactorToContestant.endsWith("ok\n") }, { label:"interactive-poll-clockless", contestant:{ language:"c", source:'#include \n#include \nint main(void){struct pollfd p[2]={{0,POLLIN,0},{1,POLLOUT,0}};if(poll(p,2,-1)<1||!(p[1].revents&POLLOUT))return 4;puts("ping");fflush(stdout);if(poll(p,1,-1)!=1||!(p[0].revents&POLLIN))return 5;char b[8];if(scanf("%7s",b)!=1)return 2;return b[0]==\'p\'&&b[1]==\'o\'?0:3;}' }, interactor:{ language:"c", source:'#include \nint main(void){char b[8];if(scanf("%7s",b)!=1)return 2;puts("pong");fflush(stdout);return 0;}' }, options:{ contestant:{ resources:{ wallTimeLimitMs:10000 } }, interactor:{ resources:{ wallTimeLimitMs:10000 } } }, check:result => exited(result.contestant, 0) && exited(result.interactor, 0) && result.contestantToInteractor === "ping\n" && result.interactorToContestant === "pong\n" }, { label:"interactive-instruction-limit", contestant:{ language:"cpp", source:'int main(){volatile unsigned long long spin=0;for(;;)spin=spin+1;}' }, interactor:guessInteractor, options:{ contestant:{ resources:{ wallTimeLimitMs:30000 } }, interactor:{ ...guessInput(1), resources:{ wallTimeLimitMs:30000 } } }, check:result => result.contestant.termination === "instruction-limit" && result.contestant.code === 137 && exited(result.interactor, 4) }, + { label:"interactive-empty-loop-budget", contestant:{ language:"cpp", source:'int main(){for(;;);}' }, interactor:readOne, options:{ contestant:{ resources:{ wallTimeLimitMs:15000 } }, interactor:{ resources:{ wallTimeLimitMs:15000 } } }, check:(result, elapsedMs) => result.contestant.termination === "instruction-limit" && result.contestant.code === 137 && exited(result.interactor, 1) && elapsedMs < 10000 }, { label:"interactive-contestant-exits", contestant:{ language:"c", source:'int main(void){return 0;}' }, interactor:guessInteractor, options:{ interactor:guessInput(1) }, check:result => exited(result.contestant, 0) && exited(result.interactor, 4) && result.contestantToInteractor === "" }, - { label:"interactive-interactor-exits", contestant:{ language:"c", source:'#include \n#include \n#include \nint main(void){char b[8];if(scanf("%7s",b)!=1)return 2;for(int i=0;i<1000000;i++)if(write(1,"x\\n",2)<0)return errno==EPIPE?32:33;return 34;}' }, interactor:{ language:"cpp", source:'#include \nint main(){std::puts("bye");std::fflush(stdout);return 0;}' }, options:{}, check:result => exited(result.contestant, 32) && exited(result.interactor, 0) && result.interactorToContestant === "bye\n" }, + { label:"interactive-interactor-exits-eof", contestant:{ language:"c", source:'#include \n#include \nint main(void){char b[8];if(scanf("%7s",b)!=1||strcmp(b,"bye")!=0)return 2;return scanf("%7s",b)==EOF&&feof(stdin)?0:3;}' }, interactor:{ language:"c", source:'#include \nint main(void){puts("bye");return 0;}' }, options:{}, check:result => exited(result.contestant, 0) && exited(result.interactor, 0) && result.interactorToContestant === "bye\n" }, + { label:"interactive-interactor-exits", contestant:{ language:"c", source:'#include \n#include \n#include \n#include \nint main(void){char b[8];if(scanf("%7s",b)!=1||strcmp(b,"bye")!=0)return 2;if(scanf("%7s",b)!=EOF)return 3;for(int i=0;i<3;i++)if(write(1,"x\\n",2)!=2)return errno==EPIPE?32:33;return 0;}' }, interactor:{ language:"c", source:'#include \nint main(void){puts("bye");return 0;}' }, options:{}, check:result => exited(result.contestant, 0) && exited(result.interactor, 0) && result.contestantToInteractor === "x\nx\nx\n" && result.interactorToContestant === "bye\n" }, + ...[["c", true, 7], ["python", true, 8], ["c", false, 8], ["python", false, 7]].map(([language, afterEof, secret]) => ({ label:`interactive-reply-${afterEof ? "after-exit" : "race"}-${language}`, contestant:finalGuess, interactor:replyInteractor(language, afterEof), options:{ interactor:secretInput(secret) }, check:result => exited(result.contestant, 0) && exited(result.interactor, secret === 7 ? 42 : 43) && result.contestantToInteractor === "7\n" && result.interactorToContestant === (secret === 7 ? "correct\n" : "wrong\n") })), { label:"interactive-output-flood", contestant:{ language:"c", source:'#include \nint main(void){while(fputs("flood\\n",stdout)>=0&&fflush(stdout)==0){}return 0;}' }, interactor:{ language:"c", source:'#include \nint main(void){while(getchar()!=EOF){}return 0;}' }, options:{ contestant:{ resources:{ outputLimitBytes:65536 } } }, check:result => result.contestant.termination === "output-limit" && result.contestant.code === 137 && result.contestantToInteractor.length === 65536 && exited(result.interactor, 0) }, { label:"interactive-wall-time", contestant:readOne, interactor:readOne, options:{ contestant:{ resources:{ wallTimeLimitMs:2000 } }, interactor:{ resources:{ wallTimeLimitMs:2000 } } }, check:(result, elapsedMs) => result.contestant.termination === "wall-time-limit" && result.interactor.termination === "wall-time-limit" && elapsedMs >= 2000 && elapsedMs < 15000 }, { label:"interactive-cancel", contestant:readOne, interactor:readOne, options:{}, cancelAfterMs:1000, recovery:{ contestant:guessC, interactor:guessInteractor, options:{ interactor:guessInput(5) } }, check:(result, elapsedMs) => result.cancelled && elapsedMs < 10000 && guessed(result.recovery) }, @@ -219,6 +231,106 @@ try { await writeFile(path.join(output,"results.json"),JSON.stringify(record,null,2)+"\n"); console.log(JSON.stringify({ label:fixture.label, pass, elapsedMs:outcome.elapsedMs, error:outcome.error, summary })); } + if (selected.length === 0 || selected.includes("liveness")) { + record.liveness = []; + const sources = { + readOne:{ language:"c", source:'#include \nint main(void){int x;return scanf("%d",&x)==1?0:1;}' }, + yieldLoop:{ language:"c", source:'#include \nint main(void){for(;;)sched_yield();}' }, + computeLoop:{ language:"cpp", source:'int main(){volatile unsigned long long spin=0;for(;;)spin=spin+1;}' }, + guessContestant:guessC, + guessInteractor, + }; + const prepared = await page.evaluate(async sources => { + window.livenessBuilds = {}; + for (const [name, { language, source }] of Object.entries(sources)) { + const entry = language === "cpp" ? "main.cpp" : "main.c"; + const files = { [entry]:source }; + if (language === "cpp") files["src/bits/stdc++.h"] = window.header; + const built = await window.engine.compile({ language, target:"wasip1", optimization:"release", entry, files, projectId:`csp-liveness-${name}` }, { cache:false }); + if (!built.success || !built.artifact) throw new Error(`liveness build ${name} failed: ${built.stderr}`); + window.livenessBuilds[name] = built.artifact; + } + const summary = value => value.termination ?? (value.contestant ? `${value.contestant.termination}/${value.interactor.code}` : `compiled:${value.success}`); + window.settle = promise => promise.then(value => ({ ok:true, summary:summary(value), stderr:value.stderr, at:performance.now() }), error => ({ ok:false, error:String(error), at:performance.now() })); + const blocked = { contestant:{ resources:{ wallTimeLimitMs:20000 } }, interactor:{ resources:{ wallTimeLimitMs:20000 } } }; + window.livenessOperation = operation => { + const builds = window.livenessBuilds; + if (operation === "interact") return window.engine.interact(builds.readOne, builds.readOne, blocked); + if (operation === "interact-compute") return window.engine.interact(builds.computeLoop, builds.readOne, { contestant:{ resources:{ instructionBudget:1e15, wallTimeLimitMs:20000 } }, interactor:{ resources:{ wallTimeLimitMs:20000 } } }); + if (operation === "run-yielding") return window.engine.run(builds.yieldLoop, { resources:{ instructionBudget:1e15, wallTimeLimitMs:20000 } }); + if (operation === "run-compute") return window.engine.run(builds.computeLoop, { resources:{ instructionBudget:1e15, wallTimeLimitMs:15000 } }); + if (operation === "compile-rust") return window.engine.compile({ language:"rust", target:"wasip1", optimization:"release", entry:"main.rs", files:{ "main.rs":'fn main(){println!("{}", 42);}' }, projectId:"csp-liveness-rust" }, { cache:false }); + return window.engine.compile({ language:"cpp", target:"wasip1", optimization:"release", entry:"main.cpp", files:{ "main.cpp":"#include \n#include \nint main(){std::regex r(\"a+\");std::cout< undefined, error => String(error)); + if (prepared) record.liveness.push({ label:"liveness-preparation", pass:false, error:prepared }); + const workerNamed = async name => { + const nameOf = worker => Promise.race([worker.evaluate(() => self.name).catch(() => ""), new Promise(resolve => setTimeout(resolve, 3000, ""))]); + for (let attempt = 0; attempt < 5; attempt++) { + for (const worker of page.workers().reverse()) if (await nameOf(worker) === name) return worker; + await page.waitForTimeout(500); + } + throw new Error(`The ${name} Worker is not visible to Playwright.`); + }; + const nestedKiller = async (parentName, childName, settleMs = 0) => { + const parent = await workerNamed(parentName); + await parent.evaluate(() => { + if (self.livenessWorkers) return; + const Native = self.Worker; self.livenessWorkers = []; self.nativeTerminate = Native.prototype.terminate; + self.Worker = class extends Native { constructor(url, options) { super(url, options); self.livenessWorkers.push({ name:options?.name, worker:this }); } }; + }); + return () => parent.evaluate(async ([name, settleMs]) => { + for (let attempt = 0; attempt < 600 && !self.livenessWorkers.some(entry => entry.name === name); attempt++) await new Promise(resolve => setTimeout(resolve, 100)); + await new Promise(resolve => setTimeout(resolve, settleMs)); + self.nativeTerminate.call(self.livenessWorkers.findLast(entry => entry.name === name).worker); + }, [childName, settleMs]); + }; + const pageKiller = name => () => page.evaluate(name => window.killWorker(name), name); + const browserName = process.env.WASM_OJ_BROWSER ?? "chromium"; + const crashed = (outcome, killMs, limitMs) => /stopped without reporting an error/.test(outcome.ok ? outcome.stderr ?? "" : outcome.error) && killMs < limitMs; + const livenessCases = [ + { label:"liveness-interactive-contestant", operation:"interact", killer:() => nestedKiller("wasm-oj-runner", "wasm-oj-interactive-contestant"), check:(outcome, killMs) => crashed(outcome, killMs, 3000) && outcome.error.includes("interactive contestant Worker") }, + { label:"liveness-interactive-interactor", operation:"interact", killer:() => nestedKiller("wasm-oj-runner", "wasm-oj-interactive-interactor"), check:(outcome, killMs) => crashed(outcome, killMs, 3000) && outcome.error.includes("interactive interactor Worker") }, + { label:"liveness-runner-interact", operation:"interact", killer:async () => pageKiller("wasm-oj-runner"), check:(outcome, killMs) => crashed(outcome, killMs, 3000) }, + { label:"liveness-runner-run-yielding", operation:"run-yielding", killer:async () => pageKiller("wasm-oj-runner"), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + { label:"liveness-runner-run-compute", operation:"run-compute", killer:async () => pageKiller("wasm-oj-runner"), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + { label:"liveness-interactive-compute", operation:"interact-compute", killer:() => nestedKiller("wasm-oj-runner", "wasm-oj-interactive-contestant"), check:(outcome, killMs) => crashed(outcome, killMs, 4000) && outcome.error.includes("interactive contestant Worker") }, + { label:"liveness-compiler", operation:"compile", killAfterMs:300, killer:async () => pageKiller("wasm-oj-compiler"), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + { label:"liveness-compiler-stage", operation:"compile-rust", killAfterMs:0, killer:() => nestedKiller("wasm-oj-compiler", "wasm-oj-rustc-stage", 1000), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + ]; + for (const fixture of prepared ? [] : livenessCases) { + console.log(`START ${fixture.label}`); + let outcome; let killMs; let error; + try { + const kill = await fixture.killer(); + await page.evaluate(operation => { window.livenessPending = window.settle(window.livenessOperation(operation)); }, fixture.operation); + await page.waitForTimeout(fixture.killAfterMs ?? 1500); + await kill(); + const killedAt = await page.evaluate(() => performance.now()); + outcome = await page.evaluate(() => window.livenessPending); + killMs = Math.round(outcome.at - killedAt); + } catch (caught) { error = String(caught); } + const recovery = await page.evaluate(() => window.settle(window.engine.run(window.livenessBuilds.readOne, { stdin:"7\n" }))); + const pass = !error && fixture.check(outcome, killMs) && recovery.summary === "exited"; + record.liveness.push({ label:fixture.label, pass, killMs, outcome, error, recovery }); + await writeFile(path.join(output,"results.json"),JSON.stringify(record,null,2)+"\n"); + console.log(JSON.stringify({ label:fixture.label, pass, killMs, error, outcome, recovery:recovery.summary ?? recovery.error })); + } + if (!prepared) { + console.log("START liveness-no-false-positive"); + const steady = await page.evaluate(async () => { + const outcomes = []; + const builds = window.livenessBuilds; + for (let index = 0; index < 20; index++) outcomes.push(await window.settle(window.engine.run(builds.readOne, { stdin:`${index}\n` }))); + for (let index = 0; index < 5; index++) outcomes.push(await window.settle(window.engine.interact(builds.guessContestant, builds.guessInteractor, { interactor:{ args:["/judge/input.txt"], files:{ "/judge/input.txt":`${index + 1} 25\n` } } }))); + outcomes.push(await window.settle(window.engine.compile({ language:"c", target:"wasip1", optimization:"release", entry:"main.c", files:{ "main.c":"int main(void){return 0;}" }, projectId:"csp-liveness-steady" }, { cache:false }))); + return outcomes.map(outcome => outcome.summary ?? outcome.error); + }); + const steadyPass = steady.length === 26 && steady.slice(0, 20).every(value => value === "exited") && steady.slice(20, 25).every(value => value === "exited/0") && steady[25] === "compiled:true"; + record.liveness.push({ label:"liveness-no-false-positive", pass:steadyPass, outcomes:steady }); + console.log(JSON.stringify({ label:"liveness-no-false-positive", pass:steadyPass, outcomes:steady })); + } else console.log(JSON.stringify({ label:"liveness-preparation", pass:false, error:prepared })); + } record.capabilities = []; for (const invoke of [false, true]) { const wasmPath = path.join(output, `capability-${invoke}.wasm`); @@ -251,5 +363,5 @@ finally { await browser?.close(); await new Promise(resolve => server.close(resolve)); } -if(record.results.some(result=>!result.pass)||record.capabilities?.some(result=>!result.pass)||record.executionTiming?.pass===false||record.interactive?.some(result=>!result.pass))process.exitCode=1; +if(record.results.some(result=>!result.pass)||record.capabilities?.some(result=>!result.pass)||record.executionTiming?.pass===false||record.interactive?.some(result=>!result.pass)||record.liveness?.some(result=>!result.pass))process.exitCode=1; console.log(`EVIDENCE ${path.join(output,"results.json")}`); diff --git a/src/core/runtime-identity.ts b/src/core/runtime-identity.ts index 6323dee..78d9d25 100644 --- a/src/core/runtime-identity.ts +++ b/src/core/runtime-identity.ts @@ -3,8 +3,8 @@ import { sha256Hex } from "./sha256.ts"; /** Executable runtime components covered by deterministic cost calibration. */ export const WASM_OJ_RUNTIME_COMPONENTS = Object.freeze({ - runtimeCoreWasmSha256: "0ae00fc646f529aafb5d1a68758220e6428f6c565e8e621eeb6c8f14ac5d4a10", - runtimeSourceRootSha256: "a7e9071511d127e3b14920558c7c38c27a1fa19d98906c8476c1df8041a9c3bb", + runtimeCoreWasmSha256: "077aa936a3382a242eee7222fc17369769cb5ad265303d730ffc45ee192367e9", + runtimeSourceRootSha256: "c65b3bbcad41359f86b8f1d1315af8f079b44f0e9084ace0d6d50c0ec418415b", wasmerNativeVersion: "7.2.1", wasmerSdkVersion: "0.10.0", wasmerSdkWasmSha256: "49a6646209f5ab5e7c737eac33407d87d9a9959ac83e5ecaaab9261b2323589e", @@ -17,7 +17,7 @@ export const WASM_OJ_RUNTIME_COMPONENTS = Object.freeze({ * identity is admitted into a calibrated release. */ export const WASM_OJ_RUNTIME_IDENTITY_SHA256 = - "e5680beae9bd832ed2b679267250ce1efc7d104dd005bb155135664b40f961fb"; + "47684bee57aafd163b18a787fa16eed33cab3ec3c0005c3c46bb4f1803ffc835"; /** Exact canonical serialization hashed by `WASM_OJ_RUNTIME_IDENTITY_SHA256`. */ export function runtimeIdentityBytes(): Uint8Array { diff --git a/src/runner/generated/runtime-core.js b/src/runner/generated/runtime-core.js index 4aaddfc..891edf0 100644 --- a/src/runner/generated/runtime-core.js +++ b/src/runner/generated/runtime-core.js @@ -560,6 +560,10 @@ function __wbg_get_imports() { const ret = new Error(); return ret; }, + __wbg_new_3867db84790e6cb4: function(arg0) { + const ret = new SharedArrayBuffer(arg0 >>> 0); + return ret; + }, __wbg_new_52e39b693b131692: function() { return handleError(function (arg0) { const ret = new WebAssembly.Tag(arg0); return ret; @@ -592,6 +596,10 @@ function __wbg_get_imports() { const ret = new Object(); return ret; }, + __wbg_new_f545e2e1ea609f29: function(arg0) { + const ret = new Int32Array(arg0); + return ret; + }, __wbg_new_f9d6489212f3b2b3: function(arg0) { const ret = new Date(arg0); return ret; @@ -769,23 +777,27 @@ function __wbg_get_imports() { const ret = arg0.value; return ret; }, + __wbg_wait_5f03a1b29e76c8cb: function() { return handleError(function (arg0, arg1, arg2, arg3) { + const ret = Atomics.wait(arg0, arg1 >>> 0, arg2, arg3); + return ret; + }, arguments); }, __wbindgen_cast_0000000000000001: function(arg0, arg1) { - // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Externref], shim_idx: 6916, ret: Result(Unit), inner_ret: Some(Result(Unit)) }, mutable: true }) -> Externref`. + // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Externref], shim_idx: 6915, ret: Result(Unit), inner_ret: Some(Result(Unit)) }, mutable: true }) -> Externref`. const ret = makeMutClosure(arg0, arg1, wasm_bindgen_68458880a41dd4bb___convert__closures_____invoke___wasm_bindgen_68458880a41dd4bb___JsValue__core_9b3796e30d99ddb7___result__Result_____wasm_bindgen_68458880a41dd4bb___JsError___true_); return ret; }, __wbindgen_cast_0000000000000002: function(arg0, arg1) { - // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Ref(NamedExternref("Array"))], shim_idx: 2393, ret: Result(Externref), inner_ret: Some(Result(Externref)) }, mutable: true }) -> Externref`. + // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Ref(NamedExternref("Array"))], shim_idx: 2326, ret: Result(Externref), inner_ret: Some(Result(Externref)) }, mutable: true }) -> Externref`. const ret = makeMutClosure(arg0, arg1, wasm_bindgen_68458880a41dd4bb___convert__closures________invoke___js_sys_c11fba41208799d1___Array__core_9b3796e30d99ddb7___result__Result_js_sys_c11fba41208799d1___Array__wasm_bindgen_68458880a41dd4bb___JsValue___true_); return ret; }, __wbindgen_cast_0000000000000003: function(arg0, arg1) { - // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Ref(NamedExternref("Array"))], shim_idx: 2393, ret: Result(NamedExternref("Array")), inner_ret: Some(Result(NamedExternref("Array"))) }, mutable: true }) -> Externref`. + // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Ref(NamedExternref("Array"))], shim_idx: 2326, ret: Result(NamedExternref("Array")), inner_ret: Some(Result(NamedExternref("Array"))) }, mutable: true }) -> Externref`. const ret = makeMutClosure(arg0, arg1, wasm_bindgen_68458880a41dd4bb___convert__closures________invoke___js_sys_c11fba41208799d1___Array__core_9b3796e30d99ddb7___result__Result_js_sys_c11fba41208799d1___Array__wasm_bindgen_68458880a41dd4bb___JsValue___true__2); return ret; }, __wbindgen_cast_0000000000000004: function(arg0, arg1) { - // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Ref(NamedExternref("Array"))], shim_idx: 2394, ret: Result(Unit), inner_ret: Some(Result(Unit)) }, mutable: true }) -> Externref`. + // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Ref(NamedExternref("Array"))], shim_idx: 2327, ret: Result(Unit), inner_ret: Some(Result(Unit)) }, mutable: true }) -> Externref`. const ret = makeMutClosure(arg0, arg1, wasm_bindgen_68458880a41dd4bb___convert__closures________invoke___js_sys_c11fba41208799d1___Array__core_9b3796e30d99ddb7___result__Result_____wasm_bindgen_68458880a41dd4bb___JsValue___true_); return ret; }, diff --git a/src/runner/generated/runtime-core_bg.wasm b/src/runner/generated/runtime-core_bg.wasm index bebed1b..1c3dcbb 100644 --- a/src/runner/generated/runtime-core_bg.wasm +++ b/src/runner/generated/runtime-core_bg.wasm @@ -1,3 +1,3 @@ version https://git-lfs.github.com/spec/v1 -oid sha256:0ae00fc646f529aafb5d1a68758220e6428f6c565e8e621eeb6c8f14ac5d4a10 -size 14599854 +oid sha256:077aa936a3382a242eee7222fc17369769cb5ad265303d730ffc45ee192367e9 +size 14612510 diff --git a/src/runtime/client-lifecycle.test.ts b/src/runtime/client-lifecycle.test.ts index bcb901e..f1ee543 100644 --- a/src/runtime/client-lifecycle.test.ts +++ b/src/runtime/client-lifecycle.test.ts @@ -44,6 +44,7 @@ const TEST_TOOLCHAINS = Object.freeze([{ interface FakeWorker { readonly messages: unknown[]; readonly listeners: Map void>>; + readonly lostListeners: Set<(error: Error) => void>; terminated: boolean; addEventListener(type: string, listener: (event: unknown) => void): void; postMessage(message: unknown): void; @@ -65,6 +66,7 @@ vi.mock("./module-worker", () => ({ const worker: FakeWorker = { messages: [], listeners: new Map(), + lostListeners: new Set(), terminated: false, addEventListener: () => undefined, postMessage: () => undefined, @@ -87,6 +89,9 @@ vi.mock("./module-worker", () => ({ collection.push(worker); return worker; }, + onModuleWorkerLost(worker: FakeWorker, listener: (error: Error) => void): void { + worker.lostListeners.add(listener); + }, })); beforeEach(() => { @@ -153,6 +158,46 @@ describe("browser client lifecycle", () => { compiler.dispose(); }); + it("rejects a running execution at once when the runner Worker is lost", async () => { + vi.useFakeTimers(); + const runner = new BrowserRunner({ toolchains: TEST_TOOLCHAINS, additionalCostBaselines: { [TEST_COST_PROFILE]: 0 } }); + try { + const worker = workerState.runners[0]!; + respondToInitialization(worker); + await runner.ready(); + const pending = runner.run(wasmArtifact(), runConfig()); + await Promise.resolve(); + const { requestId } = requestOfType(worker, "run"); + respond(worker, { type: "progress", requestId, progress: { phase: "running", label: "guest" } }); + await vi.advanceTimersByTimeAsync(10); + + lose(worker, "The wasm-oj-runner Worker stopped without reporting an error."); + + await expect(pending).rejects.toThrow("The wasm-oj-runner Worker stopped without reporting an error."); + expect(worker.terminated).toBe(true); + expect(workerState.runners).toHaveLength(2); + } finally { + runner.dispose(); + vi.useRealTimers(); + } + }); + + it("rejects a build at once when the compiler Worker is lost", async () => { + const compiler = new BrowserCompiler({ toolchains: TEST_TOOLCHAINS }); + const worker = workerState.compilers[0]!; + respondToInitialization(worker); + await compiler.ready(); + const pending = compiler.build(javascriptProject(), "cache-key"); + await vi.waitFor(() => expect(requestsOfType(worker, "build")).toHaveLength(1)); + + lose(worker, "The wasm-oj-compiler Worker stopped without reporting an error."); + + await expect(pending).rejects.toThrow("The wasm-oj-compiler Worker stopped without reporting an error."); + expect(worker.terminated).toBe(true); + expect(workerState.compilers).toHaveLength(2); + compiler.dispose(); + }); + it("rejects malformed direct compiler inputs before crossing the Worker boundary", async () => { const compiler = new BrowserCompiler({ toolchains: TEST_TOOLCHAINS }); const worker = workerState.compilers[0]!; @@ -711,6 +756,10 @@ function dispatch(worker: FakeWorker, type: string, event: unknown): void { for (const listener of worker.listeners.get(type) ?? []) listener(event); } +function lose(worker: FakeWorker, message: string): void { + for (const listener of worker.lostListeners) listener(new Error(message)); +} + function respond(worker: FakeWorker, data: unknown): void { for (const listener of worker.listeners.get("message") ?? []) listener({ data }); } diff --git a/src/runtime/compiler-client.ts b/src/runtime/compiler-client.ts index e7c84e7..5a6d1a7 100644 --- a/src/runtime/compiler-client.ts +++ b/src/runtime/compiler-client.ts @@ -29,7 +29,7 @@ import { maximumOutputReadyRustStages, } from "../compiler/browser-rust-policy"; import CompilerWorkerUrl from "./compiler.worker?worker&url"; -import { createModuleWorker } from "./module-worker"; +import { createModuleWorker, onModuleWorkerLost } from "./module-worker"; import { prefetchBrowserToolchain } from "./toolchain-prefetch"; import { clearClangBuildGraphCache } from "../compiler/indexeddb-build-graph-cache"; @@ -273,13 +273,14 @@ export class BrowserCompiler implements Compiler { worker.addEventListener("message", (event: MessageEvent) => { if (!this.disposed && !this.workerDormant && this.worker === worker) this.handleMessage(event.data); }); - worker.addEventListener("error", (event) => { - const error = new Error(event.message || "The compiler worker crashed."); + const crashed = (error: Error) => { if (this.disposed || this.workerDormant || this.worker !== worker) return; const canRecover = this.workerInitialized; this.stopWorker(error); if (canRecover) this.installWorker(); - }); + }; + worker.addEventListener("error", (event) => crashed(new Error(event.message || "The compiler worker crashed."))); + onModuleWorkerLost(worker, crashed); return worker; } diff --git a/src/runtime/interactive-pipe.ts b/src/runtime/interactive-pipe.ts index 1512990..fe16c34 100644 --- a/src/runtime/interactive-pipe.ts +++ b/src/runtime/interactive-pipe.ts @@ -7,6 +7,12 @@ const WRITER_CLOSED = 2; const READER_CLOSED = 3; const SEQUENCE = 4; const HEADER_BYTES = 32; +/** + * WebKit does not stop a Worker that `terminate()` catches inside `Atomics.wait` until the wait + * returns, and Chromium waits up to 2 s. Waking periodically lets a terminated side Worker stop, + * and release its liveness lock, promptly. + */ +const WAIT_SLICE_MS = 100; /** * The smallest ring that holds a writer's whole output budget. Only budgeted stdout bytes enter @@ -49,7 +55,7 @@ class InteractivePipeEnd { } protected sleep(sequence: number): void { - Atomics.wait(this.header, SEQUENCE, sequence); + Atomics.wait(this.header, SEQUENCE, sequence, WAIT_SLICE_MS); } } diff --git a/src/runtime/isolated-stage.ts b/src/runtime/isolated-stage.ts index b207bc5..6d4d729 100644 --- a/src/runtime/isolated-stage.ts +++ b/src/runtime/isolated-stage.ts @@ -1,3 +1,5 @@ +import { onModuleWorkerLost } from "./module-worker"; + export type IsolatedStageResponse = | { type: "result"; result: Result } | { type: "shutdown-complete" } @@ -51,6 +53,7 @@ export class PersistentIsolatedStage { this.worker.addEventListener("message", this.onMessage); this.worker.addEventListener("error", this.onError); this.worker.addEventListener("messageerror", this.onMessageError); + onModuleWorkerLost(this.worker, (error) => this.fail(error)); } run(request: Request): Promise { @@ -247,6 +250,7 @@ export function runIsolatedStage( worker.addEventListener("message", onMessage); worker.addEventListener("error", onError); worker.addEventListener("messageerror", onMessageError); + onModuleWorkerLost(worker, (error) => finish(() => reject(error))); try { worker.postMessage(request); } catch (error) { diff --git a/src/runtime/module-worker.test.ts b/src/runtime/module-worker.test.ts index 36be42a..4481b1a 100644 --- a/src/runtime/module-worker.test.ts +++ b/src/runtime/module-worker.test.ts @@ -3,6 +3,7 @@ import { createModuleWorker, createModuleWorkerBootstrap, moduleWorkerBaseUrl, + onModuleWorkerLost, } from "./module-worker"; interface WorkerConstruction { @@ -11,15 +12,42 @@ interface WorkerConstruction { } const constructions: WorkerConstruction[] = []; +const workers: FakeWorker[] = []; + +class FakeWorker extends EventTarget { + terminated = false; -class FakeWorker { constructor(url: string | URL, options?: WorkerOptions) { + super(); constructions.push({ url, options }); + workers.push(this); + } + + terminate(): void { + this.terminated = true; } } +interface LockRequest { + name: string; + signal: AbortSignal; + grant(): void; +} + +const lockRequests: LockRequest[] = []; +const fakeLocks = { + request(name: string, options: { signal: AbortSignal }, callback: () => void): Promise { + return new Promise((resolve) => { + lockRequests.push({ name, signal: options.signal, grant: () => resolve(callback()) }); + }); + }, +}; + beforeEach(() => { constructions.length = 0; + workers.length = 0; + lockRequests.length = 0; + vi.stubGlobal("navigator", { locks: fakeLocks }); vi.stubGlobal("location", { href: "https://wasm-oj.example/judge", origin: "https://wasm-oj.example", @@ -50,6 +78,7 @@ describe("module Worker bootstrap", () => { "const queueMessage = (event) => { event.stopImmediatePropagation(); pendingMessages.push(event.data); };", 'globalThis.addEventListener("message", queueMessage);', 'Object.defineProperty(globalThis, "__wasmOjModuleWorkerBaseUrl", { value: "https://wasm-oj.example/judge" });', + 'try { const lock = "wasm-oj-worker-" + crypto.randomUUID(); navigator.locks.request(lock, () => { postMessage({ __wasmOjWorkerLiveness: lock }); return new Promise(() => {}); }).catch(() => {}); } catch {}', 'try { await import("https://wasm-oj.example/assets/compiler.worker.js"); } finally { globalThis.removeEventListener("message", queueMessage); }', 'for (const data of pendingMessages) globalThis.dispatchEvent(new MessageEvent("message", { data }));', "", @@ -90,6 +119,7 @@ describe("module Worker bootstrap", () => { expect(await (source as Blob).text()).toContain( 'await import("https://wasm-oj.example/assets/wasmer-thread.worker.js")', ); + expect(await (source as Blob).text()).not.toContain("__wasmOjWorkerLiveness"); bootstrap.revoke(); bootstrap.revoke(); @@ -97,6 +127,52 @@ describe("module Worker bootstrap", () => { expect(URL.revokeObjectURL).toHaveBeenCalledWith("blob:https://wasm-oj.example/bootstrap"); }); + it("reports a Worker whose liveness lock frees before its owner terminates it", async () => { + const worker = createModuleWorker("/assets/runner.worker.js", { name: "wasm-oj-runner" }); + const messages: unknown[] = []; + const lost: string[] = []; + worker.addEventListener("message", (event) => messages.push((event as MessageEvent).data)); + onModuleWorkerLost(worker, (error) => lost.push(error.message)); + + worker.dispatchEvent(new MessageEvent("message", { data: { __wasmOjWorkerLiveness: "wasm-oj-worker-1" } })); + worker.dispatchEvent(new MessageEvent("message", { data: { type: "ready" } })); + expect(messages).toEqual([{ type: "ready" }]); + expect(lockRequests.map(({ name }) => name)).toEqual(["wasm-oj-worker-1"]); + expect(lost).toEqual([]); + + lockRequests[0]!.grant(); + expect(lost).toEqual(["The wasm-oj-runner Worker stopped without reporting an error."]); + + const late = await new Promise((resolve) => onModuleWorkerLost(worker, (error) => resolve(error.message))); + expect(late).toBe("The wasm-oj-runner Worker stopped without reporting an error."); + }); + + it("stops watching a Worker once its owner terminates it", () => { + const worker = createModuleWorker("/assets/runner.worker.js", { name: "wasm-oj-runner" }); + const lost: Error[] = []; + onModuleWorkerLost(worker, (error) => lost.push(error)); + worker.dispatchEvent(new MessageEvent("message", { data: { __wasmOjWorkerLiveness: "wasm-oj-worker-2" } })); + + worker.terminate(); + lockRequests[0]!.grant(); + + expect(workers[0]?.terminated).toBe(true); + expect(lockRequests[0]?.signal.aborted).toBe(true); + expect(lost).toEqual([]); + }); + + it("keeps the previous behaviour when Web Locks are unavailable", () => { + vi.stubGlobal("navigator", {}); + const worker = createModuleWorker("/assets/runner.worker.js"); + const messages: unknown[] = []; + worker.addEventListener("message", (event) => messages.push((event as MessageEvent).data)); + + worker.dispatchEvent(new MessageEvent("message", { data: { __wasmOjWorkerLiveness: "wasm-oj-worker-3" } })); + + expect(messages).toEqual([]); + expect(lockRequests).toEqual([]); + }); + it("uses the injected browser base inside a blob Worker", () => { vi.stubGlobal("location", { href: "blob:https://wasm-oj.example/bootstrap", diff --git a/src/runtime/module-worker.ts b/src/runtime/module-worker.ts index 235fd4d..08bece9 100644 --- a/src/runtime/module-worker.ts +++ b/src/runtime/module-worker.ts @@ -24,19 +24,27 @@ export function createModuleWorker( scriptUrl: string | URL, options: ModuleWorkerOptions = {}, ): Worker { - const bootstrap = createModuleWorkerBootstrap(scriptUrl); + const bootstrap = createModuleWorkerBootstrap(scriptUrl, { liveness: true }); + let worker: Worker; try { - return new Worker(bootstrap.url, { ...options, type: "module" }); + worker = new Worker(bootstrap.url, { ...options, type: "module" }); } finally { bootstrap.revoke(); } + superviseLiveness(worker, options.name); + return worker; } /** * Creates a reusable blob bootstrap for APIs such as the Wasmer SDK that own * Worker construction and may defer it until a command is instantiated. + * `liveness` makes the Worker hold and report its liveness lock; only + * `createModuleWorker`, which consumes that report, enables it. */ -export function createModuleWorkerBootstrap(scriptUrl: string | URL): ModuleWorkerBootstrap { +export function createModuleWorkerBootstrap( + scriptUrl: string | URL, + { liveness = false }: { liveness?: boolean } = {}, +): ModuleWorkerBootstrap { const baseUrl = moduleWorkerBaseUrl(); const absoluteScriptUrl = resolveModuleWorkerUrl(scriptUrl); const bootstrap = new Blob( @@ -45,6 +53,7 @@ export function createModuleWorkerBootstrap(scriptUrl: string | URL): ModuleWork "const queueMessage = (event) => { event.stopImmediatePropagation(); pendingMessages.push(event.data); };\n", 'globalThis.addEventListener("message", queueMessage);\n', `Object.defineProperty(globalThis, "__wasmOjModuleWorkerBaseUrl", { value: ${JSON.stringify(baseUrl.href)} });\n`, + ...(liveness ? [LIVENESS_BOOTSTRAP] : []), `try { await import(${JSON.stringify(absoluteScriptUrl)}); } finally { globalThis.removeEventListener("message", queueMessage); }\n`, 'for (const data of pendingMessages) globalThis.dispatchEvent(new MessageEvent("message", { data }));\n', ], @@ -62,6 +71,76 @@ export function createModuleWorkerBootstrap(scriptUrl: string | URL): ModuleWork }; } +const LIVENESS_MESSAGE_KEY = "__wasmOjWorkerLiveness"; + +const LIVENESS_BOOTSTRAP = "try { const lock = \"wasm-oj-worker-\" + crypto.randomUUID(); " + + `navigator.locks.request(lock, () => { postMessage({ ${LIVENESS_MESSAGE_KEY}: lock }); return new Promise(() => {}); }).catch(() => {}); } catch {}\n`; + +interface WorkerLiveness { + lost?: Error; + readonly listeners: Set<(error: Error) => void>; +} + +const liveness = new WeakMap(); + +/** + * Calls `listener` once if `worker`, created by `createModuleWorker`, stops + * without an `error` event before its owner calls `terminate()`, for example + * because the browser terminated it. Owners treat this like a crash. Without + * Web Locks the listener is never called and the previous behaviour remains. + */ +export function onModuleWorkerLost(worker: Worker, listener: (error: Error) => void): void { + const state = liveness.get(worker); + if (!state) return; + if (state.lost) { + const lost = state.lost; + queueMicrotask(() => listener(lost)); + return; + } + state.listeners.add(listener); +} + +/** + * The bootstrap holds a uniquely named Web Lock until its context is destroyed + * and reports the name. Requesting the same lock here is granted only once the + * Worker is gone; unless its owner terminated it first, that is reported to + * `onModuleWorkerLost` listeners. They are called directly rather than through + * an `error` event because WebKit drops events dispatched on a Worker after + * `terminate()`. + */ +function superviseLiveness(worker: Worker, name: string | undefined): void { + const state: WorkerLiveness = { listeners: new Set() }; + liveness.set(worker, state); + const owner = new AbortController(); + const terminate = worker.terminate.bind(worker); + worker.terminate = () => { + owner.abort(); + state.listeners.clear(); + terminate(); + }; + worker.addEventListener("message", (event: MessageEvent) => { + const lock = livenessLock(event.data); + if (lock === undefined) return; + event.stopImmediatePropagation(); + const locks = globalThis.navigator?.locks; + if (!locks) return; + void locks.request(lock, { signal: owner.signal }, () => { + if (owner.signal.aborted) return; + const lost = new Error(`The ${name ?? "module"} Worker stopped without reporting an error.`); + state.lost = lost; + const listeners = [...state.listeners]; + state.listeners.clear(); + for (const listener of listeners) listener(lost); + }).catch(() => undefined); + }); +} + +function livenessLock(data: unknown): string | undefined { + if (typeof data !== "object" || data === null) return undefined; + const lock = (data as Record)[LIVENESS_MESSAGE_KEY]; + return typeof lock === "string" ? lock : undefined; +} + export function moduleWorkerBaseUrl(): URL { const locationHref = globalThis.location?.href; if (locationHref) { diff --git a/src/runtime/runner-client.ts b/src/runtime/runner-client.ts index 5130218..a2bf7a8 100644 --- a/src/runtime/runner-client.ts +++ b/src/runtime/runner-client.ts @@ -20,7 +20,7 @@ import type { CostBaselineRegistry, } from "@wasm-oj/core"; import RunnerWorkerUrl from "./runner.worker?worker&url"; -import { createModuleWorker } from "./module-worker"; +import { createModuleWorker, onModuleWorkerLost } from "./module-worker"; import { validateBrowserRuntimeDriverPlugins } from "./browser-runtime-plugin"; import { runtimePreparationTimeoutMs } from "../runner/preparation-timeout-policy"; @@ -278,13 +278,14 @@ export class BrowserRunner implements Runner { worker.addEventListener("message", (event: MessageEvent) => { if (!this.disposed && this.worker === worker) this.handleMessage(event.data); }); - worker.addEventListener("error", (event) => { - const error = new Error(event.message || "The runner worker crashed."); + const crashed = (error: Error) => { if (this.disposed || this.worker !== worker) return; const canRecover = this.workerInitialized; this.stopWorker(error); if (canRecover) this.installWorker(); - }); + }; + worker.addEventListener("error", (event) => crashed(new Error(event.message || "The runner worker crashed."))); + onModuleWorkerLost(worker, crashed); return worker; } diff --git a/src/runtime/runner.worker.ts b/src/runtime/runner.worker.ts index 2c32fde..ddb6a53 100644 --- a/src/runtime/runner.worker.ts +++ b/src/runtime/runner.worker.ts @@ -60,6 +60,7 @@ import { createModuleWorkerBootstrap, type ModuleWorkerBootstrap, moduleWorkerBaseUrl, + onModuleWorkerLost, } from "./module-worker"; import { createInteractivePipe, interactivePipeCapacity } from "./interactive-pipe"; import type { InteractiveSideMessage, InteractiveSideStart } from "./interactive-side.worker"; @@ -515,10 +516,14 @@ function startInteractiveSide(role: InteractiveRole, start: InteractiveSideStart } } }); + const crashed = (detail: string) => { + reject(Object.assign(new Error(`The interactive ${role} Worker crashed${detail ? `: ${detail}` : "."}`), { code: "RUNTIME_ERROR" })); + }; worker.addEventListener("error", (event) => { event.preventDefault(); - reject(Object.assign(new Error(`The interactive ${role} Worker crashed${event.message ? `: ${event.message}` : "."}`), { code: "RUNTIME_ERROR" })); + crashed(event.message); }); + onModuleWorkerLost(worker, (error) => crashed(error.message)); }); try { worker.postMessage(start); diff --git a/src/server/judge.integration.test.ts b/src/server/judge.integration.test.ts index 8f9f00e..6145bf6 100644 --- a/src/server/judge.integration.test.ts +++ b/src/server/judge.integration.test.ts @@ -217,6 +217,102 @@ describe.skipIf(!enabled)("real server judge contracts", () => { }); }); + it("lets C and CPython interactors reply to a contestant that already exited", { timeout: 300_000 }, async () => { + const contestant = await compileC("final-guess", [ + "#include ", + "int main(void) {", + " puts(\"7\");", + " return 0;", + "}", + ].join("\n")); + const interactors = { + c: await compileC("c-reply-after-exit", [ + "#include ", + "int main(int argc, char **argv) {", + " long long secret = 0, guess = 0;", + " FILE *input = argc > 1 ? fopen(argv[1], \"r\") : NULL;", + " if (!input || fscanf(input, \"%lld\", &secret) != 1) return 2;", + " if (scanf(\"%lld\", &guess) != 1) return 3;", + " while (getchar() != EOF) {}", + " puts(guess == secret ? \"correct\" : \"wrong\");", + " if (fflush(stdout) != 0) return 4;", + " return guess == secret ? 42 : 43;", + "}", + ].join("\n")), + python: await compilePython("python-reply-after-exit", [ + "import sys", + "secret = int(open(sys.argv[1]).read())", + "guess = int(input())", + "sys.stdin.read()", + "print(\"correct\" if guess == secret else \"wrong\", flush=True)", + "sys.exit(42 if guess == secret else 43)", + ].join("\n")), + }; + + for (const [language, interactor] of Object.entries(interactors)) { + for (const [secret, code, reply] of [[7, 42, "correct\n"], [8, 43, "wrong\n"]] as const) { + const result = await engine.interact(contestant, interactor, { + interactor: { args: ["/judge/input.txt"], files: { "/judge/input.txt": `${secret}\n` } }, + }); + expect({ language, secret, result }).toMatchObject({ + language, + secret, + result: { + contestantToInteractor: "7\n", + interactorToContestant: reply, + contestant: { code: 0, termination: "exited" }, + interactor: { code, termination: "exited" }, + }, + }); + } + } + }); + + it("gives the contestant EOF and accepts its writes after the interactor exits", { timeout: 300_000 }, async () => { + const contestant = await compileC("write-after-interactor", [ + "#include ", + "#include ", + "#include ", + "#include ", + "int main(void) {", + " char word[8];", + " if (scanf(\"%7s\", word) != 1 || strcmp(word, \"bye\") != 0) return 2;", + " if (scanf(\"%7s\", word) != EOF || !feof(stdin)) return 3;", + " for (int i = 0; i < 3; ++i)", + " if (write(1, \"x\\n\", 2) != 2) return errno == EPIPE ? 32 : 33;", + " return 0;", + "}", + ].join("\n")); + const interactor = await compileC("early-exit", [ + "#include ", + "int main(void) {", + " puts(\"bye\");", + " return 0;", + "}", + ].join("\n")); + + const result = await engine.interact(contestant, interactor, {}); + + expect(result).toMatchObject({ + contestantToInteractor: "x\nx\nx\n", + interactorToContestant: "bye\n", + contestant: { code: 0, termination: "exited" }, + interactor: { code: 0, termination: "exited" }, + }); + }); + + async function compilePython(name: string, source: string): Promise { + const built = await engine.compile({ + projectId: `judge-integration:${name}`, + name, + language: "python", + entry: "main.py", + files: { "main.py": `${source}\n` }, + }, { cache: false }); + if (!built.success || !built.artifact) throw new Error(`Failed to compile ${name}: ${built.stderr}`); + return built.artifact; + } + async function compileC(name: string, source: string): Promise { const entry = `src/${name}.c`; const input: CompileInput = {