diff --git a/CHANGELOG.md b/CHANGELOG.md index 80dac74..88081b0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -25,6 +25,13 @@ All notable changes to WASM-OJ are recorded here. Releases follow self-delimiting archive or JSON compiler response is complete instead of waiting for `Instance.wait()`. A guest that exits before completing its output still fails with its stderr, about 2 s after it exits. +- Fix browser `Engine.interact`, which failed on every dialogue (#98). Each side now runs in its + own nested Worker as a standalone metered run, connected to the other through shared-memory pipes + that block with `Atomics.wait`. Browser interactive costs equal `run` costs for the same program. + Each pipe holds the writing side's whole output budget, and closing stdin or stdout signals the + peer at once, so browser and server give the same verdicts. The runner Worker also sends + `startupEntropyBytes` for interactive programs. The refreshed runtime identity changes cost + profiles. ## 0.2.3 - 2026-10-05 diff --git a/crates/runtime-core/src/deterministic.rs b/crates/runtime-core/src/deterministic.rs index 0110a58..0ff964d 100644 --- a/crates/runtime-core/src/deterministic.rs +++ b/crates/runtime-core/src/deterministic.rs @@ -6,6 +6,7 @@ use crate::{ }, types::DeterminismConfig, }; +use std::cell::Cell; use std::sync::{Arc, Mutex}; use wasmer::{ AsStoreMut, Extern, Function, FunctionEnv, FunctionEnvMut, Imports, Memory, Memory32, Memory64, @@ -239,6 +240,23 @@ struct PollEnv { clock: VirtualClock, } +thread_local! { + static STDIN_READINESS_DEFERRALS: Cell = const { Cell::new(0) }; +} + +/// Whether an empty stdin readiness check may return `Pending` instead of +/// blocking. A clock probe never blocks. A poll without a clock grants one +/// deferral per read subscription, so WASIX reports any other ready +/// subscription first and blocks on stdin only when it would wait anyway. +#[cfg(target_arch = "wasm32")] +pub(crate) fn defer_stdin_readiness() -> bool { + STDIN_READINESS_DEFERRALS.with(|deferrals| { + let remaining = deferrals.get(); + deferrals.set(remaining.saturating_sub(1)); + remaining > 0 + }) +} + #[derive(Clone, Copy)] struct ClockSubscription { subscription: Subscription, @@ -251,6 +269,7 @@ fn call_original_poll( events: WasmPtr, subscription_count: M::Offset, event_count: WasmPtr, + stdin_deferrals: u64, ) -> Result { let Some(original) = env.data().original.clone() else { return Ok(Errno::Notsup as i32); @@ -262,6 +281,7 @@ fn call_original_poll( Value::I32(offset as i32) } }; + STDIN_READINESS_DEFERRALS.with(|deferrals| deferrals.set(stdin_deferrals)); let results = original.call( env, &[ @@ -270,7 +290,9 @@ fn call_original_poll( pointer(subscription_count.into()), pointer(event_count.offset().into()), ], - )?; + ); + STDIN_READINESS_DEFERRALS.with(|deferrals| deferrals.set(0)); + let results = results?; match results.as_ref() { [Value::I32(errno)] => Ok(*errno), _ => Err(RuntimeError::new( @@ -302,6 +324,7 @@ fn deterministic_poll_oneoff( let mut originals = Vec::with_capacity(input.len() as usize); let mut clock_subscriptions = Vec::new(); let mut has_non_clock = false; + let mut read_subscriptions = 0; for index in 0..input.len() { let subscription = input .index(index) @@ -320,6 +343,7 @@ fn deterministic_poll_oneoff( }); } else { has_non_clock = true; + read_subscriptions += u64::from(subscription.type_ == Eventtype::FdRead); } } @@ -330,6 +354,7 @@ fn deterministic_poll_oneoff( events, subscription_count, event_count, + read_subscriptions, ); } @@ -355,6 +380,7 @@ fn deterministic_poll_oneoff( events, subscription_count, event_count, + u64::MAX, ); let restore_view = memory.view(&env); let restore_input = subscriptions diff --git a/crates/runtime-core/src/interactive.rs b/crates/runtime-core/src/interactive.rs index 252b8ac..77a9f29 100644 --- a/crates/runtime-core/src/interactive.rs +++ b/crates/runtime-core/src/interactive.rs @@ -971,6 +971,66 @@ mod tests { assert_eq!(result.interactor.metrics.logical_time_ns, 0); } + #[test] + fn interactive_empty_stdin_poll_times_out_on_the_process_clock() { + let poller = wat::parse_str( + r#"(module + (import "wasi_snapshot_preview1" "poll_oneoff" + (func $poll (param i32 i32 i32 i32) (result i32))) + (import "wasi_snapshot_preview1" "fd_write" + (func $fd_write (param i32 i32 i32 i32) (result i32))) + (memory (export "memory") 1) + (func (export "_start") + i32.const 0 i64.const 1 i64.store + i32.const 8 i32.const 1 i32.store8 + i32.const 16 i32.const 0 i32.store + i32.const 48 i64.const 2 i64.store + i32.const 56 i32.const 0 i32.store8 + i32.const 64 i32.const 1 i32.store + i32.const 72 i64.const 5000000000 i64.store + i32.const 80 i64.const 1 i64.store + i32.const 88 i32.const 0 i32.store16 + i32.const 0 i32.const 128 i32.const 2 i32.const 240 call $poll drop + i32.const 200 i32.const 138 i32.store + i32.const 204 i32.const 1 i32.store + i32.const 208 i32.const 240 i32.store + i32.const 212 i32.const 1 i32.store + i32.const 1 i32.const 200 i32.const 2 i32.const 216 call $fd_write drop))"#, + ) + .unwrap(); + let listener = wat::parse_str( + r#"(module + (import "wasi_snapshot_preview1" "fd_read" + (func $fd_read (param i32 i32 i32 i32) (result i32))) + (memory (export "memory") 1) + (func (export "_start") + (i32.store (i32.const 0) (i32.const 64)) + (i32.store (i32.const 4) (i32.const 8)) + (drop (call $fd_read (i32.const 0) (i32.const 0) (i32.const 1) (i32.const 8)))))"#, + ) + .unwrap(); + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap(); + let result = runtime + .block_on(interact(InteractiveRequest { + contestant: program(poller), + interactor: program(listener), + determinism: DeterminismConfig { + random_seed: 7, + realtime_epoch_ms: 946_684_800_000, + clock_step_ns: 1_000_000, + }, + })) + .unwrap(); + + assert_eq!(result.contestant_to_interactor, [0, 1]); + assert_eq!(result.contestant.termination, ExecutionTermination::Exited); + assert_eq!(result.contestant.metrics.logical_time_ns, 5_000_000_000); + assert_eq!(result.interactor.termination, ExecutionTermination::Exited); + } + #[test] fn interactive_contestant_is_metered_like_a_standalone_run() { let looping = |ending: &str| { diff --git a/crates/runtime-core/src/run/mod.rs b/crates/runtime-core/src/run/mod.rs index 4907948..8cfa0fa 100644 --- a/crates/runtime-core/src/run/mod.rs +++ b/crates/runtime-core/src/run/mod.rs @@ -15,6 +15,8 @@ mod native; #[cfg(target_arch = "wasm32")] mod web; #[cfg(target_arch = "wasm32")] +pub(crate) mod web_interactive; +#[cfg(target_arch = "wasm32")] pub(crate) mod web_runtime; pub fn run(request: RunRequest) -> Result { diff --git a/crates/runtime-core/src/run/web.rs b/crates/runtime-core/src/run/web.rs index a7d6d95..7cca6ea 100644 --- a/crates/runtime-core/src/run/web.rs +++ b/crates/runtime-core/src/run/web.rs @@ -5,10 +5,11 @@ use crate::filesystem::{read_files_bounded, runtime_project_files}; use crate::meter::{CostPoints, METER_MODEL, 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}; +use crate::output::{CappedOutput, OutputBudget, OutputCapture}; use crate::{ExecutionMetrics, ExecutionTermination, RunError, RunRequest, RunResult}; use std::io::Write; use std::sync::{Arc, Mutex}; +use virtual_fs::VirtualFile; use wasmer::{Instance, Memory, Module, Store}; use wasmer_wasix::{ Pipe, WasiEnv, WasiError, WasiModuleInstanceHandles, WasiModuleTreeHandles, wasmer_wasix_types, @@ -16,8 +17,29 @@ use wasmer_wasix::{ pub fn run( request: RunRequest, - mut on_execution: impl FnMut(bool) -> Result<(), RunError>, + on_execution: impl FnMut(bool) -> Result<(), RunError>, ) -> Result { + let (mut stdin_writer, stdin_reader) = Pipe::channel(); + stdin_writer + .write_all(&request.stdin) + .map_err(|error| RunError::Io(error.to_string()))?; + drop(stdin_writer); + execute(request, Box::new(stdin_reader), None, on_execution).map(|execution| execution.result) +} + +pub(super) type StdoutPeer = Box Box>; + +pub(super) struct Execution { + pub result: RunResult, + pub exit_code: Option, +} + +pub(super) fn execute( + request: RunRequest, + stdin: Box, + stdout_peer: Option, + mut on_execution: impl FnMut(bool) -> Result<(), RunError>, +) -> Result { let limited = enforce_memory_limit(&request.wasm, request.resources.memory_limit_bytes) .map_err(RunError::Compile)?; let metered = instrument_wasm(&limited, request.resources.instruction_budget) @@ -29,11 +51,6 @@ pub fn run( })?; let runtime = runtime_with_engine(store.engine().clone()); - let (mut stdin_writer, stdin_reader) = Pipe::channel(); - stdin_writer - .write_all(&request.stdin) - .map_err(|error| RunError::Io(error.to_string()))?; - drop(stdin_writer); let output_limit = usize::try_from(request.resources.output_limit_bytes) .map_err(|_| RunError::InvalidRequest("output limit exceeds host range".to_string()))?; let output_budget = OutputBudget::new(output_limit); @@ -51,8 +68,11 @@ pub fn run( .runtime(runtime) .args(request.args.clone()) .envs(request.env.clone()) - .stdin(Box::new(stdin_reader)) - .stdout(Box::new(stdout_file)) + .stdin(stdin) + .stdout(match stdout_peer { + Some(peer) => peer(stdout_file), + None => Box::new(stdout_file), + }) .stderr(Box::new(stderr_file)) .fs(filesystem.clone()); builder @@ -138,12 +158,14 @@ pub fn run( let logical_time_exceeded = clock.limit_exceeded()?; let mut code = 0; + let mut exit_code = None; let mut termination = ExecutionTermination::Exited; let mut trap_message = None; if let Err(error) = execution { if let Some(wasi_error) = crate::wasi_error(&error) { match wasi_error { WasiError::Exit(exit) => { + exit_code = Some(exit.raw()); let errno: wasmer_wasix_types::wasi::Errno = (*exit).into(); if errno != wasmer_wasix_types::wasi::Errno::Success { code = errno as i32; @@ -211,7 +233,7 @@ pub fn run( CostPoints::Exhausted => request.resources.instruction_budget, }; let filesystem_metrics = project_filesystem.metrics(); - Ok(RunResult { + let result = RunResult { code, metrics: ExecutionMetrics { cost, @@ -231,5 +253,6 @@ pub fn run( trap_message, determinism: request.determinism, resources: request.resources, - }) + }; + Ok(Execution { result, exit_code }) } diff --git a/crates/runtime-core/src/run/web_interactive.rs b/crates/runtime-core/src/run/web_interactive.rs new file mode 100644 index 0000000..c16f73b --- /dev/null +++ b/crates/runtime-core/src/run/web_interactive.rs @@ -0,0 +1,370 @@ +use super::web::{Execution, execute}; +use crate::output::CappedOutput; +use crate::{ + DeterminismConfig, ExecutionTermination, InteractiveMetrics, InteractiveProcessResult, + InteractiveProgram, RunError, RunRequest, +}; +use serde::{Deserialize, Serialize}; +use std::io; +use std::pin::Pin; +use std::sync::Arc; +use std::task::{Context, Poll}; +use tokio::io::{AsyncRead, AsyncSeek, AsyncWrite, ReadBuf}; +use virtual_fs::{FsError, VirtualFile}; +use wasm_bindgen::JsValue; + +#[derive(Clone, Debug, Deserialize, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct InteractiveSideRequest { + pub program: InteractiveProgram, + pub determinism: DeterminismConfig, +} + +#[derive(Clone, Debug, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct InteractiveSideResult { + pub process: InteractiveProcessResult, + #[serde(with = "serde_bytes")] + pub protocol: Vec, +} + +/// Stream callbacks supplied by the side Worker. Reads and writes wait with +/// `Atomics.wait`, so they never return `Pending` to WASIX. +pub struct HostStreams { + read: js_sys::Function, + poll: js_sys::Function, + wait: js_sys::Function, + write: js_sys::Function, + close: js_sys::Function, +} + +// SAFETY: the web build of runtime-core has no threads; these JS handles never +// leave the Worker thread that created them. +unsafe impl Send for HostStreams {} +// SAFETY: see `Send`. +unsafe impl Sync for HostStreams {} + +impl HostStreams { + pub fn new( + read: js_sys::Function, + poll: js_sys::Function, + wait: js_sys::Function, + write: js_sys::Function, + close: js_sys::Function, + ) -> Self { + Self { + read, + poll, + wait, + write, + close, + } + } + + fn read(&self, maximum: usize) -> io::Result> { + let chunk = self + .read + .call1(&JsValue::UNDEFINED, &JsValue::from(maximum as u32)) + .map_err(host_error)?; + Ok(js_sys::Uint8Array::new(&chunk).to_vec()) + } + + fn poll(&self) -> io::Result> { + let available = self + .poll + .call0(&JsValue::UNDEFINED) + .map_err(host_error)? + .as_f64() + .ok_or_else(|| io::Error::other("interactive input poll returned no size"))?; + Ok((available >= 0.0).then_some(available as usize)) + } + + fn wait(&self) -> io::Result { + let available = self + .wait + .call0(&JsValue::UNDEFINED) + .map_err(host_error)? + .as_f64() + .ok_or_else(|| io::Error::other("interactive input wait returned no size"))?; + Ok(available as usize) + } + + fn write(&self, bytes: &[u8]) -> io::Result { + let written = self + .write + .call1(&JsValue::UNDEFINED, &js_sys::Uint8Array::from(bytes)) + .map_err(host_error)? + .as_f64() + .ok_or_else(|| io::Error::other("interactive output write returned no size"))?; + if written < 0.0 { + return Err(io::ErrorKind::BrokenPipe.into()); + } + Ok(written as usize) + } + + /// Closes the input (fd 0) or output (fd 1) pipe end, so the peer sees + /// EOF or a broken pipe at once, as when native drops its pipe end. + fn close(&self, fd: u32) { + let _ = self.close.call1(&JsValue::UNDEFINED, &JsValue::from(fd)); + } +} + +fn host_error(error: JsValue) -> io::Error { + io::Error::other(format!("interactive stream failed: {error:?}")) +} + +pub fn run_side( + request: InteractiveSideRequest, + streams: HostStreams, + on_execution: impl FnMut(bool) -> Result<(), RunError>, +) -> Result { + let program = request.program; + if program.wasm.is_empty() { + return Err(RunError::InvalidRequest( + "interactive Wasm must not be empty".to_string(), + )); + } + let run = RunRequest { + wasm: program.wasm, + args: program.args, + env: program.env, + stdin: Vec::new(), + files: program.files, + output_paths: Vec::new(), + cwd: program.cwd, + startup_entropy_bytes: program.startup_entropy_bytes, + determinism: request.determinism, + resources: program.resources, + }; + super::validate(&run)?; + let streams = Arc::new(streams); + let input = StreamInput { + streams: streams.clone(), + }; + let Execution { result, exit_code } = execute( + run, + Box::new(input), + Some(Box::new(move |capture| { + Box::new(StreamOutput { streams, capture }) + })), + on_execution, + )?; + let code = match exit_code { + Some(code) if result.termination == ExecutionTermination::Exited => code, + _ => result.code, + }; + Ok(InteractiveSideResult { + process: InteractiveProcessResult { + code, + stderr: result.stderr, + termination: result.termination, + metrics: InteractiveMetrics { + cost: result.metrics.cost, + operations: result.metrics.operations, + logical_time_ns: result.metrics.logical_time_ns, + filesystem_bytes: result.metrics.filesystem_bytes, + filesystem_entries: result.metrics.filesystem_entries, + protocol_bytes: result.metrics.stdout_bytes, + stderr_bytes: result.metrics.stderr_bytes, + }, + }, + protocol: result.stdout, + }) +} + +struct StreamInput { + streams: Arc, +} + +// WASIX drops the stdio handle when the last descriptor that refers to it is +// closed, which is when native drops its `PipeRx` or `PipeTx`. +impl Drop for StreamInput { + fn drop(&mut self) { + self.streams.close(0); + } +} + +impl std::fmt::Debug for StreamInput { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("StreamInput") + } +} + +impl AsyncRead for StreamInput { + fn poll_read( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + buffer: &mut ReadBuf<'_>, + ) -> Poll> { + Poll::Ready(self.streams.read(buffer.remaining()).map(|chunk| { + buffer.put_slice(&chunk); + })) + } +} + +impl AsyncWrite for StreamInput { + fn poll_write( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + _buffer: &[u8], + ) -> Poll> { + Poll::Ready(Err(io::ErrorKind::Unsupported.into())) + } + + fn poll_flush(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } +} + +impl AsyncSeek for StreamInput { + fn start_seek(self: Pin<&mut Self>, _position: io::SeekFrom) -> io::Result<()> { + Ok(()) + } + + fn poll_complete(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(0)) + } +} + +impl VirtualFile for StreamInput { + fn last_accessed(&self) -> u64 { + 0 + } + fn last_modified(&self) -> u64 { + 0 + } + fn created_time(&self) -> u64 { + 0 + } + fn size(&self) -> u64 { + 0 + } + fn set_len(&mut self, _new_size: u64) -> Result<(), FsError> { + Ok(()) + } + fn unlink(&mut self) -> Result<(), FsError> { + Ok(()) + } + fn get_special_fd(&self) -> Option { + Some(0) + } + + fn poll_read_ready(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll> { + match self.streams.poll() { + Ok(Some(available)) => Poll::Ready(Ok(available)), + Ok(None) if crate::deterministic::defer_stdin_readiness() => { + context.waker().wake_by_ref(); + Poll::Pending + } + Ok(None) => Poll::Ready(self.streams.wait()), + Err(error) => Poll::Ready(Err(error)), + } + } + + fn poll_write_ready( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + ) -> Poll> { + Poll::Ready(Err(io::ErrorKind::Unsupported.into())) + } +} + +struct StreamOutput { + streams: Arc, + capture: CappedOutput, +} + +impl Drop for StreamOutput { + fn drop(&mut self) { + self.streams.close(1); + } +} + +impl std::fmt::Debug for StreamOutput { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("StreamOutput") + } +} + +impl AsyncWrite for StreamOutput { + 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])), + result => result, + } + } + + fn poll_flush(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } +} + +impl AsyncRead for StreamOutput { + fn poll_read( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + _buffer: &mut ReadBuf<'_>, + ) -> Poll> { + Poll::Ready(Ok(())) + } +} + +impl AsyncSeek for StreamOutput { + fn start_seek(self: Pin<&mut Self>, _position: io::SeekFrom) -> io::Result<()> { + Ok(()) + } + + fn poll_complete(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(0)) + } +} + +impl VirtualFile for StreamOutput { + fn last_accessed(&self) -> u64 { + 0 + } + fn last_modified(&self) -> u64 { + 0 + } + fn created_time(&self) -> u64 { + 0 + } + fn size(&self) -> u64 { + 0 + } + fn set_len(&mut self, _new_size: u64) -> Result<(), FsError> { + Ok(()) + } + fn unlink(&mut self) -> Result<(), FsError> { + Ok(()) + } + fn get_special_fd(&self) -> Option { + Some(1) + } + + fn poll_read_ready( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + ) -> Poll> { + Poll::Ready(Ok(0)) + } + + fn poll_write_ready( + mut self: Pin<&mut Self>, + context: &mut Context<'_>, + ) -> Poll> { + Pin::new(&mut self.capture).poll_write_ready(context) + } +} diff --git a/crates/runtime-core/src/web.rs b/crates/runtime-core/src/web.rs index 53c5ebe..e890bb9 100644 --- a/crates/runtime-core/src/web.rs +++ b/crates/runtime-core/src/web.rs @@ -1,6 +1,7 @@ +use crate::run::web_interactive::{HostStreams, InteractiveSideRequest, run_side}; use crate::{ GoCompilerSession as CoreGoCompilerSession, GoCompilerSessionConfig, GoCompilerSessionRequest, - InteractiveRequest, RunError, RunRequest, interactive_response, run_response_from_result, + RunError, RunFailure, RunRequest, run_response_from_result, }; use serde::Serialize; use wasm_bindgen::prelude::*; @@ -72,17 +73,62 @@ impl WebGoCompilerSession { } } +/// Runs one side of an interactive session in the calling Worker. `read`, +/// `wait` and `write` block on the session's shared ring buffers; `poll` +/// checks the input without blocking; `close(fd)` closes the input (0) or +/// output (1) end when the guest drops it. #[wasm_bindgen] -pub async fn interact_wasm_oj(request: JsValue) -> Result { +pub fn run_interactive_side( + request: JsValue, + read: js_sys::Function, + poll: js_sys::Function, + wait: js_sys::Function, + write: js_sys::Function, + close: js_sys::Function, + on_execution: js_sys::Function, +) -> Result { console_error_panic_hook::set_once(); - let request: InteractiveRequest = serde_wasm_bindgen::from_value(request) - .map_err(|error| JsValue::from_str(&format!("invalid interactive request: {error}")))?; - let response = interactive_response(request).await; + let request: InteractiveSideRequest = + serde_wasm_bindgen::from_value(request).map_err(|error| { + JsValue::from_str(&format!("invalid interactive side request: {error}")) + })?; + let response = match run_side( + request, + HostStreams::new(read, poll, wait, write, close), + |running| { + on_execution + .call1(&JsValue::UNDEFINED, &JsValue::from_bool(running)) + .map(|_| ()) + .map_err(|error| RunError::Runtime(format!("execution observer failed: {error:?}"))) + }, + ) { + Ok(result) => InteractiveSideResponse { + ok: true, + result: Some(result), + error: None, + }, + Err(error) => InteractiveSideResponse { + ok: false, + result: None, + error: Some(RunFailure { + code: error.code(), + message: error.to_string(), + }), + }, + }; response .serialize(&serde_wasm_bindgen::Serializer::new().serialize_maps_as_objects(true)) .map_err(|error| { JsValue::from_str(&format!( - "failed to serialize interactive response: {error}" + "failed to serialize interactive side response: {error}" )) }) } + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct InteractiveSideResponse { + ok: bool, + result: Option, + error: Option, +} diff --git a/docs/architecture.md b/docs/architecture.md index 2c60251..9406b48 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -53,7 +53,14 @@ Browser requests cross dedicated module-Worker boundaries. C/C++ retains bounded and content-addressed graph state; Rust and Go use serialized nested stages with bounded generation lifetime. Python and JavaScript package source files directly. Server builds use a fresh isolated child. Cancellation, restart, timeout, cache clearing, disposal, family switch, and stage-budget exhaustion establish a -complete browser Worker-generation boundary. +complete browser Worker-generation boundary. Browser interaction runs each side in its own nested +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 +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`. The runner is one Rust codebase compiled to browser Wasm and native executables. Both forms admit the same artifacts, deterministic inputs, resource limits, denied capabilities, filesystem model, diff --git a/scripts/verify-browser-csp.mjs b/scripts/verify-browser-csp.mjs index 767ca5b..5f9a545 100644 --- a/scripts/verify-browser-csp.mjs +++ b/scripts/verify-browser-csp.mjs @@ -91,6 +91,31 @@ fixtures.push( { 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} }, ); +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 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 = [ + { label:"interactive-ac-cpp", contestant:guessCpp, interactor:guessInteractor, options:{ interactor:guessInput(1) }, check:result => guessed(result) && result.contestantToInteractor.split("\n").length - 1 === 20 }, + { label:"interactive-ac-c", contestant:guessC, interactor:guessInteractor, options:{ interactor:guessInput(777777) }, check:guessed }, + { label:"interactive-wa-c", contestant:{ language:"c", source:'#include \nint main(void){char r[4];for(;;){puts("7");fflush(stdout);if(scanf("%3s",r)!=1)return 0;}}' }, interactor:guessInteractor, options:{ interactor:guessInput(1) }, check:result => exited(result.interactor, 1) && exited(result.contestant, 0) && result.contestantToInteractor === "7\n".repeat(26) }, + { label:"interactive-python-contestant", contestant:{ language:"python", source:'lo, hi = 1, 1 << 20\nwhile lo <= hi:\n mid = (lo + hi) // 2\n print(mid, flush=True)\n reply = input().strip()\n if reply == "=":\n break\n if reply == "<":\n lo = mid + 1\n else:\n hi = mid - 1\n' }, interactor:guessInteractor, options:{ interactor:guessInput(31337) }, check:guessed }, + { label:"interactive-python-interactor", contestant:guessC, interactor:{ language:"python", source:'import sys\nsecret, limit = map(int, open(sys.argv[1]).read().split())\nfor _ in range(limit):\n try:\n guess = int(input())\n except EOFError:\n sys.exit(4)\n if guess == secret:\n print("=", flush=True)\n sys.exit(0)\n print("<" if guess < secret else ">", flush=True)\nsys.exit(1)\n' }, options:{ interactor:guessInput(424242) }, check:guessed }, + { label:"interactive-poll-timeout", contestant:{ language:"c", source:'#include \n#include \nint main(void){struct pollfd p={0,POLLIN,0};int r=poll(&p,1,1000);printf("%d\\n",r);fflush(stdout);char b[8];if(scanf("%7s",b)!=1)return 2;return b[0]==\'o\'?0:3;}' }, interactor:{ language:"c", source:'#include \nint main(void){char b[8];if(scanf("%7s",b)!=1)return 2;puts("ok");fflush(stdout);return b[0]==\'0\'?0:1;}' }, options:{ contestant:{ resources:{ wallTimeLimitMs:10000 } }, interactor:{ resources:{ wallTimeLimitMs:10000 } } }, check:(result, elapsedMs) => exited(result.contestant, 0) && exited(result.interactor, 0) && result.contestantToInteractor === "0\n" && result.interactorToContestant === "ok\n" && result.contestant.metrics.logicalTimeNs >= 1e9 && elapsedMs < 5000 }, + { label:"interactive-poll-ready", contestant:{ language:"c", source:'#include \n#include \nint main(void){puts("ping");fflush(stdout);struct pollfd p={0,POLLIN,0};while(poll(&p,1,0)==0){}if(!(p.revents&POLLIN))return 4;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-close-stdout", contestant:{ language:"c", source:'#include \nint main(void){puts("42");if(fclose(stdout)!=0)return 4;char b[8];if(scanf("%7s",b)!=1)return 2;return b[0]==\'o\'&&b[1]==\'k\'?0:3;}' }, interactor:{ language:"c", source:'#include \nint main(void){int x,n=0,s=0;while(scanf("%d",&x)==1){n++;s+=x;}int ok=n==1&&s==42;puts(ok?"ok":"no");fflush(stdout);return ok?0:1;}' }, options:{ contestant:{ resources:{ wallTimeLimitMs:10000 } }, interactor:{ resources:{ wallTimeLimitMs:10000 } } }, check:result => exited(result.contestant, 0) && exited(result.interactor, 0) && result.contestantToInteractor === "42\n" && result.interactorToContestant === "ok\n" }, + { 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-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-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) }, +]; const record = { policy, results:[], pageErrors:[], consoleErrors:[] }; let browser; try { @@ -145,6 +170,55 @@ try { }, [...await readFile(wasmPath)]); console.log(JSON.stringify({executionTiming:record.executionTiming})); } + const interactiveSelected = interactiveFixtures.filter(fixture => selected.length === 0 || selected.includes("interactive") || selected.includes(fixture.label)); + if (interactiveSelected.length > 0) record.interactive = []; + for (const fixture of interactiveSelected) { + console.log(`START ${fixture.label}`); + const { check, ...input } = fixture; + const outcome = await page.evaluate(async fixture => { + const entries = { c:"main.c", cpp:"main.cpp", python:"main.py" }; + window.interactiveBuilds ??= new Map(); + const build = async program => { + const key = `${program.language}\0${program.source}`; + if (!window.interactiveBuilds.has(key)) { + const entry = entries[program.language]; + const files = { [entry]:program.source }; + if (program.language === "cpp") files["src/bits/stdc++.h"] = window.header; + const built = await window.engine.compile({ language:program.language, target:"wasip1", optimization:"release", entry, files, projectId:`csp-interactive-${window.interactiveBuilds.size}` }, { cache:false }); + if (!built.success || !built.artifact) throw new Error(`${program.language} build failed: ${built.stderr}`); + window.interactiveBuilds.set(key, built.artifact); + } + return window.interactiveBuilds.get(key); + }; + try { + const contestant = await build(fixture.contestant); + const interactor = await build(fixture.interactor); + const start = performance.now(); + if (fixture.cancelAfterMs === undefined) { + const result = await window.engine.interact(contestant, interactor, fixture.options); + return { result, elapsedMs:Math.round(performance.now() - start) }; + } + const pending = window.engine.interact(contestant, interactor, fixture.options); + const timer = setTimeout(() => window.engine.cancel(), fixture.cancelAfterMs); + let cancelled = false; + let error; + try { await pending; } catch (caught) { cancelled = true; error = String(caught); } finally { clearTimeout(timer); } + const elapsedMs = Math.round(performance.now() - start); + const recovery = await window.engine.interact(await build(fixture.recovery.contestant), await build(fixture.recovery.interactor), fixture.recovery.options); + return { result:{ cancelled, error, recovery }, elapsedMs }; + } catch (error) { return { error:String(error) }; } + }, input); + const pass = !outcome.error && check(outcome.result, outcome.elapsedMs); + const summary = outcome.result && !outcome.result.recovery ? { + contestant:{ code:outcome.result.contestant.code, termination:outcome.result.contestant.termination, cost:outcome.result.contestant.metrics?.cost }, + interactor:{ code:outcome.result.interactor.code, termination:outcome.result.interactor.termination, cost:outcome.result.interactor.metrics?.cost }, + contestantToInteractorBytes:outcome.result.contestantToInteractor.length, + interactorToContestantBytes:outcome.result.interactorToContestant.length, + } : outcome.result; + record.interactive.push({ label:fixture.label, pass, elapsedMs:outcome.elapsedMs, error:outcome.error, result:outcome.result }); + 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 })); + } record.capabilities = []; for (const invoke of [false, true]) { const wasmPath = path.join(output, `capability-${invoke}.wasm`); @@ -177,5 +251,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)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))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 85ef21a..6323dee 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: "522924286709c7661f9c23155f3269b7679d50573c62a5ea19d4500c7b37bbef", - runtimeSourceRootSha256: "0493153cecc1ca16349a8bb5ee1a84a1eb3a5aebf34b1e290630cf646a39a3bc", + runtimeCoreWasmSha256: "0ae00fc646f529aafb5d1a68758220e6428f6c565e8e621eeb6c8f14ac5d4a10", + runtimeSourceRootSha256: "a7e9071511d127e3b14920558c7c38c27a1fa19d98906c8476c1df8041a9c3bb", 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 = - "98e10a30ed5f406536c70464a7300bab96c0f13b43fe15f0f2d11894ec051299"; + "e5680beae9bd832ed2b679267250ce1efc7d104dd005bb155135664b40f961fb"; /** Exact canonical serialization hashed by `WASM_OJ_RUNTIME_IDENTITY_SHA256`. */ export function runtimeIdentityBytes(): Uint8Array { diff --git a/src/runner/generated/runtime-core.d.ts b/src/runner/generated/runtime-core.d.ts index 5cd2d50..40c38e8 100644 --- a/src/runner/generated/runtime-core.d.ts +++ b/src/runner/generated/runtime-core.d.ts @@ -15,7 +15,13 @@ export class GoCompilerSession { readonly generation: number; } -export function interact_wasm_oj(request: any): Promise; +/** + * Runs one side of an interactive session in the calling Worker. `read`, + * `wait` and `write` block on the session's shared ring buffers; `poll` + * checks the input without blocking; `close(fd)` closes the input (0) or + * output (1) end when the guest drops it. + */ +export function run_interactive_side(request: any, read: Function, poll: Function, wait: Function, write: Function, close: Function, on_execution: Function): any; export function run_wasm_oj(request: any, on_execution: Function): any; @@ -28,7 +34,7 @@ export interface InitOutput { readonly gocompilersession_digest: (a: number) => [number, number]; readonly gocompilersession_generation: (a: number) => [number, number, number]; readonly gocompilersession_new: (a: any) => [number, number, number]; - readonly interact_wasm_oj: (a: any) => any; + readonly run_interactive_side: (a: any, b: any, c: any, d: any, e: any, f: any, g: any) => [number, number, number]; readonly run_wasm_oj: (a: any, b: any) => [number, number, number]; readonly canonical_abi_free: (a: number, b: number, c: number) => void; readonly canonical_abi_realloc: (a: number, b: number, c: number, d: number) => number; diff --git a/src/runner/generated/runtime-core.js b/src/runner/generated/runtime-core.js index 773a163..4aaddfc 100644 --- a/src/runner/generated/runtime-core.js +++ b/src/runner/generated/runtime-core.js @@ -95,12 +95,25 @@ export class Trap { if (Symbol.dispose) Trap.prototype[Symbol.dispose] = Trap.prototype.free; /** + * Runs one side of an interactive session in the calling Worker. `read`, + * `wait` and `write` block on the session's shared ring buffers; `poll` + * checks the input without blocking; `close(fd)` closes the input (0) or + * output (1) end when the guest drops it. * @param {any} request - * @returns {Promise} + * @param {Function} read + * @param {Function} poll + * @param {Function} wait + * @param {Function} write + * @param {Function} close + * @param {Function} on_execution + * @returns {any} */ -export function interact_wasm_oj(request) { - const ret = wasm.interact_wasm_oj(request); - return ret; +export function run_interactive_side(request, read, poll, wait, write, close, on_execution) { + const ret = wasm.run_interactive_side(request, read, poll, wait, write, close, on_execution); + if (ret[2]) { + throw takeFromExternrefTable0(ret[1]); + } + return takeFromExternrefTable0(ret[0]); } /** @@ -587,6 +600,10 @@ function __wbg_get_imports() { const ret = new WebAssembly.Memory(arg0); return ret; }, arguments); }, + __wbg_new_from_slice_3eea173078478cfe: function(arg0, arg1) { + const ret = new Uint8Array(getArrayU8FromWasm0(arg0, arg1)); + return ret; + }, __wbg_new_typed_cceaf62d8d95e9f2: function(arg0, arg1) { try { var state0 = {a: arg0, b: arg1}; @@ -753,22 +770,22 @@ function __wbg_get_imports() { return ret; }, __wbindgen_cast_0000000000000001: function(arg0, arg1) { - // Cast intrinsic for `Closure(Closure { owned: true, function: Function { arguments: [Externref], shim_idx: 6970, ret: Result(Unit), inner_ret: Some(Result(Unit)) }, mutable: true }) -> Externref`. + // 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`. 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: 2445, 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: 2393, 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: 2445, 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: 2393, 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: 2446, 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: 2394, 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 dd82d58..bebed1b 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:522924286709c7661f9c23155f3269b7679d50573c62a5ea19d4500c7b37bbef -size 14907700 +oid sha256:0ae00fc646f529aafb5d1a68758220e6428f6c565e8e621eeb6c8f14ac5d4a10 +size 14599854 diff --git a/src/runner/generated/runtime-core_bg.wasm.d.ts b/src/runner/generated/runtime-core_bg.wasm.d.ts index b3ea8cd..597b8ec 100644 --- a/src/runner/generated/runtime-core_bg.wasm.d.ts +++ b/src/runner/generated/runtime-core_bg.wasm.d.ts @@ -6,7 +6,7 @@ export const gocompilersession_compilePipeline: (a: number, b: any) => any; export const gocompilersession_digest: (a: number) => [number, number]; export const gocompilersession_generation: (a: number) => [number, number, number]; export const gocompilersession_new: (a: any) => [number, number, number]; -export const interact_wasm_oj: (a: any) => any; +export const run_interactive_side: (a: any, b: any, c: any, d: any, e: any, f: any, g: any) => [number, number, number]; export const run_wasm_oj: (a: any, b: any) => [number, number, number]; export const canonical_abi_free: (a: number, b: number, c: number) => void; export const canonical_abi_realloc: (a: number, b: number, c: number, d: number) => number; diff --git a/src/runtime/interactive-pipe.test.ts b/src/runtime/interactive-pipe.test.ts new file mode 100644 index 0000000..9226342 --- /dev/null +++ b/src/runtime/interactive-pipe.test.ts @@ -0,0 +1,111 @@ +import { describe, expect, it } from "vitest"; +import { + createInteractivePipe, + interactivePipeCapacity, + InteractivePipeReader, + InteractivePipeWriter, +} from "./interactive-pipe"; + +const bytes = (text: string) => new TextEncoder().encode(text); +const text = (data: Uint8Array) => new TextDecoder().decode(data); + +describe("interactive pipe", () => { + it("delivers bytes in order across the ring boundary", () => { + const buffer = createInteractivePipe(8); + const writer = new InteractivePipeWriter(buffer); + const reader = new InteractivePipeReader(buffer); + expect(writer.write(bytes("abcdef"))).toBe(6); + expect(text(reader.read(4))).toBe("abcd"); + expect(writer.write(bytes("ghijkl"))).toBe(6); + expect(reader.wait()).toBe(8); + expect(text(reader.read(100))).toBe("efghijkl"); + }); + + it("reports EOF only after buffered bytes are drained", () => { + const buffer = createInteractivePipe(16); + const writer = new InteractivePipeWriter(buffer); + const reader = new InteractivePipeReader(buffer); + writer.write(bytes("42\n")); + writer.close(); + expect(text(reader.read(2))).toBe("42"); + expect(text(reader.read(2))).toBe("\n"); + expect(reader.wait()).toBe(0); + expect(reader.read(2).byteLength).toBe(0); + }); + + it("returns bytes the writer publishes and closes with after the reader saw an empty buffer", () => { + const buffer = createInteractivePipe(16); + const writer = new InteractivePipeWriter(buffer); + class RacedReader extends InteractivePipeReader { + private raced = false; + + protected override buffered(): number { + const buffered = super.buffered(); + if (!this.raced) { + this.raced = true; + writer.write(bytes("42\n")); + writer.close(); + } + return buffered; + } + } + const reader = new RacedReader(buffer); + expect(text(reader.read(16))).toBe("42\n"); + expect(reader.read(16).byteLength).toBe(0); + }); + + it("polls without blocking", () => { + const buffer = createInteractivePipe(16); + const writer = new InteractivePipeWriter(buffer); + const reader = new InteractivePipeReader(buffer); + expect(reader.poll()).toBe(-1); + writer.write(bytes("ok")); + expect(reader.poll()).toBe(2); + expect(text(reader.read(2))).toBe("ok"); + writer.close(); + expect(reader.poll()).toBe(0); + }); + + it("fails writes once the reader has closed", () => { + const buffer = createInteractivePipe(16); + const writer = new InteractivePipeWriter(buffer); + new InteractivePipeReader(buffer).close(); + expect(writer.write(bytes("ignored"))).toBe(-1); + expect(writer.write(new Uint8Array())).toBe(-1); + }); + + it("fails a write the reader closes on partway, like the server's broken pipe", () => { + const buffer = createInteractivePipe(8); + const reader = new InteractivePipeReader(buffer); + class ClosedWhileBlocked extends InteractivePipeWriter { + protected override sleep(): void { + reader.read(4); + reader.close(); + } + } + expect(new ClosedWhileBlocked(buffer).write(bytes("0123456789ab"))).toBe(-1); + }); + + it("sizes rings to hold the writer's whole output budget", () => { + expect(interactivePipeCapacity(1024)).toBe(64 * 1024); + expect(interactivePipeCapacity(4 * 1024 * 1024)).toBe(4 * 1024 * 1024); + expect(interactivePipeCapacity(4 * 1024 * 1024 + 1)).toBe(8 * 1024 * 1024); + expect(interactivePipeCapacity(2 ** 40)).toBe(2 ** 30); + }); + + it("keeps positions consistent when the 32-bit counters wrap", () => { + const buffer = createInteractivePipe(8); + const header = new Int32Array(buffer, 0, 8); + header[0] = -3; + header[1] = -3; + const writer = new InteractivePipeWriter(buffer); + const reader = new InteractivePipeReader(buffer); + expect(writer.write(bytes("wrap!"))).toBe(5); + expect(text(reader.read(8))).toBe("wrap!"); + expect(header[0]).toBe(2); + }); + + it("rejects capacities that are not powers of two", () => { + expect(() => createInteractivePipe(12)).toThrow(/power of two/); + }); +}); diff --git a/src/runtime/interactive-pipe.ts b/src/runtime/interactive-pipe.ts new file mode 100644 index 0000000..1512990 --- /dev/null +++ b/src/runtime/interactive-pipe.ts @@ -0,0 +1,123 @@ +/** One direction of an interactive session: a single-producer, single-consumer byte ring in shared memory. */ +export const INTERACTIVE_PIPE_CAPACITY_BYTES = 64 * 1024; + +const READ = 0; +const WRITE = 1; +const WRITER_CLOSED = 2; +const READER_CLOSED = 3; +const SEQUENCE = 4; +const HEADER_BYTES = 32; + +/** + * The smallest ring that holds a writer's whole output budget. Only budgeted stdout bytes enter + * the pipe, so a write never blocks before the budget is spent, as on the server's unbounded pipe. + */ +export function interactivePipeCapacity(outputLimitBytes: number): number { + let capacity = INTERACTIVE_PIPE_CAPACITY_BYTES; + while (capacity < outputLimitBytes && capacity < 2 ** 30) capacity *= 2; + return capacity; +} + +export function createInteractivePipe(capacity = INTERACTIVE_PIPE_CAPACITY_BYTES): SharedArrayBuffer { + if (!Number.isSafeInteger(capacity) || capacity <= 0 || (capacity & (capacity - 1)) !== 0 || capacity > 2 ** 30) { + throw new Error("Interactive pipe capacity must be a power of two up to 1 GiB."); + } + return new SharedArrayBuffer(HEADER_BYTES + capacity); +} + +class InteractivePipeEnd { + protected readonly header: Int32Array; + protected readonly data: Uint8Array; + + constructor(buffer: SharedArrayBuffer) { + this.header = new Int32Array(buffer, 0, HEADER_BYTES / 4); + this.data = new Uint8Array(buffer, HEADER_BYTES); + } + + protected buffered(): number { + return (Atomics.load(this.header, WRITE) - Atomics.load(this.header, READ)) >>> 0; + } + + protected signal(index: number, value: number): void { + Atomics.store(this.header, index, value); + Atomics.add(this.header, SEQUENCE, 1); + Atomics.notify(this.header, SEQUENCE); + } + + protected sequence(): number { + return Atomics.load(this.header, SEQUENCE); + } + + protected sleep(sequence: number): void { + Atomics.wait(this.header, SEQUENCE, sequence); + } +} + +export class InteractivePipeReader extends InteractivePipeEnd { + /** Returns the buffered count, 0 at EOF, or -1 while the pipe is empty and the writer is open. */ + poll(): number { + const buffered = this.buffered(); + if (buffered > 0) return buffered; + if (Atomics.load(this.header, WRITER_CLOSED) === 0) return -1; + return this.buffered(); + } + + /** Blocks until bytes are buffered or the writer closed; returns the buffered count, or 0 at EOF. */ + wait(): number { + for (;;) { + const sequence = this.sequence(); + const available = this.poll(); + if (available >= 0) return available; + this.sleep(sequence); + } + } + + /** Blocks like `wait`, then consumes up to `maximum` bytes. An empty result means EOF. */ + read(maximum: number): Uint8Array { + const buffered = this.wait(); + const count = Math.min(buffered, maximum); + const output = new Uint8Array(count); + if (count === 0) return output; + const read = Atomics.load(this.header, READ); + const start = (read >>> 0) % this.data.length; + const first = Math.min(count, this.data.length - start); + output.set(this.data.subarray(start, start + first)); + output.set(this.data.subarray(0, count - first), first); + this.signal(READ, (read + count) | 0); + return output; + } + + close(): void { + this.signal(READER_CLOSED, 1); + } +} + +export class InteractivePipeWriter extends InteractivePipeEnd { + /** Blocks until every byte is buffered. Returns the count written, or -1 if the reader closed before it finished. */ + write(bytes: Uint8Array): number { + if (Atomics.load(this.header, READER_CLOSED) !== 0) return -1; + let offset = 0; + while (offset < bytes.length) { + const sequence = this.sequence(); + if (Atomics.load(this.header, READER_CLOSED) !== 0) return -1; + const free = this.data.length - this.buffered(); + if (free === 0) { + this.sleep(sequence); + continue; + } + const count = Math.min(free, bytes.length - offset); + const write = Atomics.load(this.header, WRITE); + const start = (write >>> 0) % this.data.length; + const first = Math.min(count, this.data.length - start); + this.data.set(bytes.subarray(offset, offset + first), start); + this.data.set(bytes.subarray(offset + first, offset + count), 0); + this.signal(WRITE, (write + count) | 0); + offset += count; + } + return offset; + } + + close(): void { + this.signal(WRITER_CLOSED, 1); + } +} diff --git a/src/runtime/interactive-side.worker.ts b/src/runtime/interactive-side.worker.ts new file mode 100644 index 0000000..3144257 --- /dev/null +++ b/src/runtime/interactive-side.worker.ts @@ -0,0 +1,55 @@ +/// + +import initRuntimeCore, { + run_interactive_side as runInteractiveSide, +} from "@/src/runner/generated/runtime-core.js"; +import { InteractivePipeReader, InteractivePipeWriter } from "./interactive-pipe"; + +export interface InteractiveSideStart { + runtimeCore: WebAssembly.Module; + request: unknown; + input: SharedArrayBuffer; + output: SharedArrayBuffer; +} + +export type InteractiveSideMessage = + | { type: "running" } + | { type: "result"; response: unknown } + | { type: "error"; message: string }; + +const scope: DedicatedWorkerGlobalScope = self as unknown as DedicatedWorkerGlobalScope; + +function post(message: InteractiveSideMessage): void { + scope.postMessage(message); +} + +scope.addEventListener("message", (event: MessageEvent) => { + void runSide(event.data); +}, { once: true }); + +async function runSide(start: InteractiveSideStart): Promise { + const input = new InteractivePipeReader(start.input); + const output = new InteractivePipeWriter(start.output); + let response: unknown; + try { + await initRuntimeCore({ module_or_path: start.runtimeCore }); + response = runInteractiveSide( + start.request, + (maximum: number) => input.read(maximum), + () => input.poll(), + () => input.wait(), + (bytes: Uint8Array) => output.write(bytes), + (fd: number) => (fd === 0 ? input : output).close(), + (running: boolean) => { + if (running) post({ type: "running" }); + }, + ); + } catch (error) { + post({ type: "error", message: error instanceof Error ? error.message : String(error) }); + return; + } finally { + output.close(); + input.close(); + } + post({ type: "result", response }); +} diff --git a/src/runtime/runner.worker.ts b/src/runtime/runner.worker.ts index e0c7cc8..2c32fde 100644 --- a/src/runtime/runner.worker.ts +++ b/src/runtime/runner.worker.ts @@ -53,16 +53,17 @@ import { withHandleLease, withWasmerCommand, } from "@/src/runner/package-handle-cache"; -import initRuntimeCore, { - interact_wasm_oj as interactWasmOjCore, - run_wasm_oj as runWasmOjCore, -} from "@/src/runner/generated/runtime-core.js"; +import initRuntimeCore, { run_wasm_oj as runWasmOjCore } from "@/src/runner/generated/runtime-core.js"; import runtimeCoreWasmUrl from "@/src/runner/generated/runtime-core_bg.wasm?url"; import { + createModuleWorker, createModuleWorkerBootstrap, type ModuleWorkerBootstrap, moduleWorkerBaseUrl, } from "./module-worker"; +import { createInteractivePipe, interactivePipeCapacity } from "./interactive-pipe"; +import type { InteractiveSideMessage, InteractiveSideStart } from "./interactive-side.worker"; +import interactiveSideWorkerUrl from "./interactive-side.worker?worker&url"; import wasmerThreadWorkerUrl from "./wasmer-thread.worker?worker&url"; import { loadBrowserRuntimeDriverPlugins } from "./browser-runtime-plugin"; import type { BrowserRuntimeDriverPlugin } from "@/src/core/types"; @@ -78,6 +79,7 @@ let wasmerThreadWorkerBootstrap: ModuleWorkerBootstrap | undefined; let runtimeDrivers: RuntimeDriverRegistry | undefined; let quickJsBytes: Promise | undefined; let toolchainSources: readonly BrowserToolchainSource[] | undefined; +let runtimeCoreModule: WebAssembly.Module | undefined; interface CoreRunResult { code: number; @@ -120,13 +122,11 @@ interface CoreInteractiveProcessResult { }; } -interface CoreInteractiveResponse { +interface CoreInteractiveSideResponse { ok: boolean; result?: { - contestant: CoreInteractiveProcessResult; - interactor: CoreInteractiveProcessResult; - contestantToInteractor: Uint8Array; - interactorToContestant: Uint8Array; + process: CoreInteractiveProcessResult; + protocol: Uint8Array; }; error?: { code: string; message: string }; } @@ -151,7 +151,8 @@ async function initializeRuntime( toolchainSources = snapshotBrowserToolchainSources(sources); progress(requestId, "initializing", "Starting deterministic Wasmer runner", 0.1); - await initRuntimeCore({ module_or_path: new URL(runtimeCoreWasmUrl, workerBaseUrl) }); + runtimeCoreModule = await compileRuntimeCore(); + await initRuntimeCore({ module_or_path: runtimeCoreModule }); runtimeDrivers = createDefaultRuntimeDrivers( createExtendedCostBaselineRegistry(additionalCostBaselines), ); @@ -164,6 +165,14 @@ async function initializeRuntime( progress(requestId, "initializing", "Deterministic Wasmer runner ready", 1); } +async function compileRuntimeCore(): Promise { + const response = await fetch(new URL(runtimeCoreWasmUrl, workerBaseUrl)); + if (!response.ok) throw new Error(`Unable to load the WASM-OJ runtime core (${response.status}).`); + return response.headers.get("Content-Type")?.startsWith("application/wasm") + ? WebAssembly.compileStreaming(response) + : WebAssembly.compile(await response.arrayBuffer()); +} + async function ensurePackageRuntime(): Promise { if (sdkRuntime) return sdkRuntime; sdkRuntimeInitialization ??= (async () => { @@ -430,26 +439,96 @@ async function interactArtifacts( runtimeDrivers, ), ]); - progress(request.requestId, "running", "Running interactive session with deterministic Wasmer", 0.25); - const response = await interactWasmOjCore({ - contestant: interactiveCoreProgram(contestant), - interactor: interactiveCoreProgram(interactor), - determinism: request.config.determinism, - }) as CoreInteractiveResponse; - if (!response.ok || !response.result) { - const error = response.error ?? { code: "RUNTIME_ERROR", message: "The interactive runtime returned no result." }; - throw Object.assign(new Error(error.message), { code: error.code }); - } + const [contestantSide, interactorSide] = await runInteractiveSides( + request.requestId, + { contestant: interactiveCoreProgram(contestant), interactor: interactiveCoreProgram(interactor) }, + request.config.determinism, + ); return { - contestant: interactiveProcessResult(response.result.contestant, contestant), - interactor: interactiveProcessResult(response.result.interactor, interactor), - contestantToInteractor: decoder.decode(response.result.contestantToInteractor), - interactorToContestant: decoder.decode(response.result.interactorToContestant), + contestant: interactiveProcessResult(contestantSide.process, contestant), + interactor: interactiveProcessResult(interactorSide.process, interactor), + contestantToInteractor: decoder.decode(contestantSide.protocol), + interactorToContestant: decoder.decode(interactorSide.protocol), durationMs: performance.now() - started, determinism: { ...request.config.determinism }, }; } +type InteractiveRole = "contestant" | "interactor"; +type CoreInteractiveSide = NonNullable; + +async function runInteractiveSides( + requestId: string, + programs: Record>, + determinism: InteractiveRunConfig["determinism"], +): Promise<[CoreInteractiveSide, CoreInteractiveSide]> { + if (!runtimeCoreModule) throw new Error("The WASM-OJ runtime core is not initialized."); + const contestantToInteractor = createInteractivePipe(interactivePipeCapacity(programs.contestant.resources.outputLimitBytes)); + const interactorToContestant = createInteractivePipe(interactivePipeCapacity(programs.interactor.resources.outputLimitBytes)); + const sides: ReturnType[] = []; + try { + sides.push(startInteractiveSide("contestant", { + runtimeCore: runtimeCoreModule, + request: { program: programs.contestant, determinism }, + input: interactorToContestant, + output: contestantToInteractor, + })); + sides.push(startInteractiveSide("interactor", { + runtimeCore: runtimeCoreModule, + request: { program: programs.interactor, determinism }, + input: contestantToInteractor, + output: interactorToContestant, + })); + void Promise.all(sides.map((side) => side.running)).then(() => { + progress(requestId, "running", "Running interactive session with deterministic Wasmer", 0.25); + }); + const results = await Promise.all(sides.map((side) => side.result)); + progress(requestId, "packaging", "Collecting interactive results", 0.9); + return [results[0], results[1]]; + } finally { + for (const side of sides) side.worker.terminate(); + } +} + +function startInteractiveSide(role: InteractiveRole, start: InteractiveSideStart) { + const worker = createModuleWorker(new URL(interactiveSideWorkerUrl, workerBaseUrl), { + name: `wasm-oj-interactive-${role}`, + }); + let markRunning!: () => void; + const running = new Promise((resolve) => { + markRunning = resolve; + }); + const result = new Promise((resolve, reject) => { + worker.addEventListener("message", (event: MessageEvent) => { + const message = event.data; + if (message.type === "running") { + markRunning(); + } else if (message.type === "error") { + reject(Object.assign(new Error(`Interactive ${role} failed: ${message.message}`), { code: "RUNTIME_ERROR" })); + } else { + const response = message.response as CoreInteractiveSideResponse; + if (response.ok && response.result) { + resolve(response.result); + } else { + const error = response.error ?? { code: "RUNTIME_ERROR", message: "the runtime returned no result." }; + reject(Object.assign(new Error(`Interactive ${role} failed: ${error.message}`), { code: error.code })); + } + } + }); + worker.addEventListener("error", (event) => { + event.preventDefault(); + reject(Object.assign(new Error(`The interactive ${role} Worker crashed${event.message ? `: ${event.message}` : "."}`), { code: "RUNTIME_ERROR" })); + }); + }); + try { + worker.postMessage(start); + } catch (error) { + worker.terminate(); + throw error; + } + return { worker, running, result }; +} + function interactiveRunConfig( program: InteractiveRunConfig["contestant"], determinism: InteractiveRunConfig["determinism"], @@ -473,6 +552,7 @@ function interactiveCoreProgram(prepared: Awaited