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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,12 @@ All notable changes to WASM-OJ are recorded here. Releases follow
- Accept runtime-bundle interactors (for example CPython) in `Runner.interact` on the server
and in the browser runner Worker. Either side of an interactive session may now be a
standalone Wasm module or a runtime bundle that provides streaming fd 0.
- Fixed Python runtime preparation and TypeScript compilation occasionally stalling until their
300 s and 120 s timeouts. The Wasmer SDK can terminate its workers before stdout/stderr reach
EOF even though all output has arrived, so the server and browser now finish once the
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.

## 0.2.3 - 2026-10-05

Expand Down
27 changes: 13 additions & 14 deletions src/compiler/wasmer-engine.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
import {
Runtime,
Wasmer,
type Output,
} from "@wasmer/sdk";
import { WASM_OJ_CONTRACT_VERSION } from "../core/contract.ts";
import { completeJsonObject, readProcessMessage } from "../core/process-output.ts";
import {
canonicalRuntimeBundleFiles,
createRuntimeBundleManifest,
Expand Down Expand Up @@ -307,7 +307,7 @@ interface TypeScriptWasiResponse {
files: Record<string, string>;
}

async function transpileScriptProject(project: Project, requestId: string): Promise<{ files: Record<string, string | Uint8Array>; output: Output; response?: TypeScriptWasiResponse }> {
async function transpileScriptProject(project: Project, requestId: string): Promise<{ files: Record<string, string | Uint8Array>; stderr: string; response?: TypeScriptWasiResponse }> {
const scriptFiles = scriptSourceFiles(project);
const emittedFiles = emittedSourceFiles(project);
const dependencyFiles = npmDependencyFiles(project);
Expand Down Expand Up @@ -337,15 +337,14 @@ async function transpileScriptProject(project: Project, requestId: string): Prom
outputs: outputPaths.map((path) => `/project/build/${path}`),
}),
});
const output = await instance.wait();
let response: TypeScriptWasiResponse | undefined;
if (output.ok) {
try {
response = JSON.parse(output.stdout) as TypeScriptWasiResponse;
} catch {
response = undefined;
}
}
// Instance.wait() can hang after all output arrived (see readProcessMessage). It is also the only
// way to read the exit code, so a complete response, the driver's last write, stands in for exit 0.
const output = await readProcessMessage<TypeScriptWasiResponse>(
instance,
completeJsonObject,
{ collectStderr: true },
).finally(() => instance.free());
const response = output.message;
const files: Record<string, string | Uint8Array> = {};
if (response) {
for (const outputPath of outputPaths) {
Expand All @@ -354,7 +353,7 @@ async function transpileScriptProject(project: Project, requestId: string): Prom
}
}
Object.assign(files, dependencyFiles);
return { files, output, response };
return { files, stderr: output.stderr, response };
}

async function buildScript(project: Project, cacheKey: string, requestId: string): Promise<BuildResult> {
Expand Down Expand Up @@ -388,11 +387,11 @@ async function buildScript(project: Project, cacheKey: string, requestId: string
progress(requestId, "compiling", "Compiling TypeScript with TypeScript/WASI", 0.5);
const transpiled = await transpileScriptProject(project, requestId);
files = transpiled.files;
stderr = transpiled.output.stderr;
stderr = transpiled.stderr;
diagnostics = parseTypeScriptDiagnostics(transpiled.response?.diagnostics ?? "");
const emittedOutputsPresent = emittedSourceFiles(project)
.every((file) => Object.hasOwn(files, emittedScriptPath(file.path)));
if (!transpiled.output.ok || !transpiled.response || transpiled.response.status !== 0 || !emittedOutputsPresent || diagnostics.some((diagnostic) => diagnostic.severity === "error")) {
if (!transpiled.response || transpiled.response.status !== 0 || !emittedOutputsPresent || diagnostics.some((diagnostic) => diagnostic.severity === "error")) {
return {
success: false,
diagnostics: ensureFailureDiagnostic(diagnostics, {
Expand Down
81 changes: 81 additions & 0 deletions src/core/process-output.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
import { describe, expect, it } from "vitest";
import { completeJsonObject, readProcessMessage } from "./process-output";

const encoder = new TextEncoder();

function pipe() {
let controller!: ReadableStreamDefaultController<Uint8Array>;
const stream = new ReadableStream<Uint8Array>({ start: (value) => { controller = value; } });
return {
stream,
write: (text: string) => controller.enqueue(encoder.encode(text)),
close: () => controller.close(),
};
}

const delay = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));

describe("process output messages", () => {
it("completes a JSON response without stdout or stderr EOF and keeps late stderr", async () => {
const stdout = pipe();
const stderr = pipe();
const result = readProcessMessage<{ status: number; text: string }>(
{ stdout: stdout.stream, stderr: stderr.stream },
completeJsonObject,
{ collectStderr: true, idleGraceMs: 50 },
);
// A chunk boundary right after a `}` inside a string must not complete the response.
stdout.write('{"status":0,"text":"f() {}');
stdout.write('"}\n');
await delay(10);
stderr.write("warning: late");
await expect(result).resolves.toEqual({
message: { status: 0, text: "f() {}" },
stderr: "warning: late",
});
});

it("returns as soon as stderr ends after the response completes", async () => {
const stdout = pipe();
const stderr = pipe();
const started = performance.now();
const result = readProcessMessage(
{ stdout: stdout.stream, stderr: stderr.stream },
completeJsonObject,
{ collectStderr: true, idleGraceMs: 60_000 },
);
stdout.write('{"status":0}');
stderr.close();
await expect(result).resolves.toEqual({ message: { status: 0 }, stderr: "" });
expect(performance.now() - started).toBeLessThan(1_000);
});

it("gives up after the idle grace once stdout ended without a complete response", async () => {
const stdout = pipe();
const stderr = pipe();
const result = readProcessMessage(
{ stdout: stdout.stream, stderr: stderr.stream },
completeJsonObject,
{ idleGraceMs: 30 },
);
stdout.write('{"status":');
stderr.write("panic: out of memory");
stdout.close();
await expect(result).resolves.toEqual({ message: undefined, stderr: "panic: out of memory" });
});

it("keeps waiting while both streams are open", async () => {
const stdout = pipe();
const stderr = pipe();
const result = readProcessMessage(
{ stdout: stdout.stream, stderr: stderr.stream },
completeJsonObject,
{ idleGraceMs: 10 },
);
stdout.write('{"status":');
const early = await Promise.race([result, delay(100).then(() => "pending")]);
expect(early).toBe("pending");
stdout.write("0}");
await expect(result).resolves.toEqual({ message: { status: 0 }, stderr: "" });
});
});
119 changes: 119 additions & 0 deletions src/core/process-output.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
/** How long a read waits for new bytes once the guest has exited or the message is complete. */
export const PROCESS_OUTPUT_IDLE_GRACE_MS = 2_000;

export interface StreamingProcess {
readonly stdin?: WritableStream | undefined;
readonly stdout: ReadableStream<Uint8Array>;
readonly stderr: ReadableStream<Uint8Array>;
}

export interface ProcessMessageOptions {
readonly idleGraceMs?: number | undefined;
/** Keep reading stderr after the message completes, until EOF or the idle grace. */
readonly collectStderr?: boolean | undefined;
}

export interface ProcessMessage<T> {
/** The parsed stdout message, or undefined when the process stopped before completing it. */
readonly message: T | undefined;
readonly stderr: string;
}

type ReadEvent = { readonly source: "stdout" | "stderr"; readonly chunk?: Uint8Array } | { readonly source: "idle" };

/**
* Read stdout until `parse` recognizes one complete self-delimiting message.
*
* `@wasmer/sdk` 0.10 `Instance.wait()` joins stdout EOF, stderr EOF and the exit code, and its
* thread-pool teardown can drop either EOF after all output arrived, so `wait()` never settles.
* Completion therefore comes from the message itself. Either EOF means the guest exited, so the
* other stream then gets `idleGraceMs` without new bytes before the read stops. If both EOFs are
* lost before the message completes, the caller's own deadline still bounds the read.
*/
export async function readProcessMessage<T>(
child: StreamingProcess,
parse: (stdout: Uint8Array) => T | undefined,
options: ProcessMessageOptions = {},
): Promise<ProcessMessage<T>> {
const idleGraceMs = options.idleGraceMs ?? PROCESS_OUTPUT_IDLE_GRACE_MS;
await child.stdin?.close().catch(() => undefined);
const stdoutReader = child.stdout.getReader();
const stderrReader = child.stderr.getReader();
const nextStdout = () => stdoutReader.read()
.then(({ done, value }): ReadEvent => ({ source: "stdout", chunk: done ? undefined : value }));
const nextStderr = () => stderrReader.read()
.then(({ done, value }): ReadEvent => ({ source: "stderr", chunk: done ? undefined : value }))
.catch((): ReadEvent => ({ source: "stderr" }));
const stdout = byteBuffer();
const stderr = byteBuffer();
let stdoutRead: Promise<ReadEvent> | undefined = nextStdout();
let stderrRead: Promise<ReadEvent> | undefined = nextStderr();
let message: T | undefined;
try {
while (true) {
const complete = message !== undefined;
const reads = complete ? [options.collectStderr ? stderrRead : undefined] : [stdoutRead, stderrRead];
const pending = reads.filter((read): read is Promise<ReadEvent> => read !== undefined);
if (pending.length === 0) break;
// While both streams are open and the message is incomplete, the guest may still be running.
let timer: ReturnType<typeof setTimeout> | undefined;
if (complete || pending.length === 1) {
pending.push(new Promise((resolve) => {
timer = setTimeout(() => resolve({ source: "idle" }), idleGraceMs);
}));
}
const event = await Promise.race(pending).finally(() => clearTimeout(timer));
if (event.source === "idle") break;
if (event.source === "stdout") {
if (!event.chunk) {
stdoutRead = undefined;
continue;
}
stdout.append(event.chunk);
message = parse(stdout.bytes());
stdoutRead = message === undefined ? nextStdout() : undefined;
} else if (!event.chunk) {
stderrRead = undefined;
} else {
stderr.append(event.chunk);
stderrRead = nextStderr();
}
}
} finally {
await Promise.allSettled([stdoutReader.cancel(), stderrReader.cancel()]);
}
return { message, stderr: new TextDecoder().decode(stderr.bytes()) };
}

const JSON_WHITESPACE = new Set([0x09, 0x0a, 0x0d, 0x20]);

/** Parse stdout once it holds one complete JSON object, optionally followed by whitespace. */
export function completeJsonObject<T>(stdout: Uint8Array): T | undefined {
let end = stdout.byteLength;
while (end > 0 && JSON_WHITESPACE.has(stdout[end - 1])) end -= 1;
// Only output ending in `}` can hold the whole object, which keeps parse attempts rare.
if (end === 0 || stdout[end - 1] !== 0x7d) return undefined;
try {
return JSON.parse(new TextDecoder().decode(stdout.subarray(0, end))) as T;
} catch {
return undefined;
}
}

/** An append-only byte buffer that grows geometrically, so appending chunks stays linear. */
function byteBuffer() {
let buffer = new Uint8Array(0);
let length = 0;
return {
append(chunk: Uint8Array): void {
if (length + chunk.byteLength > buffer.byteLength) {
const grown = new Uint8Array(Math.max(length + chunk.byteLength, buffer.byteLength * 2));
grown.set(buffer.subarray(0, length));
buffer = grown;
}
buffer.set(chunk, length);
length += chunk.byteLength;
},
bytes: (): Uint8Array => buffer.subarray(0, length),
};
}
98 changes: 97 additions & 1 deletion src/runner/runtime-files.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
import { describe, expect, it } from "vitest";
import { sha256Hex } from "../core/hash";
import { decodeRuntimeFiles, verifyAndDecodeRuntimeFiles } from "./runtime-files";
import {
decodeRuntimeFiles,
readRuntimeFilesExport,
verifyAndDecodeRuntimeFiles,
} from "./runtime-files";

const encoder = new TextEncoder();

Expand Down Expand Up @@ -59,3 +63,95 @@ describe("runtime file archives", () => {
);
});
});

function outputStream(chunks: readonly Uint8Array[], end: boolean) {
const state = { cancelled: false };
const stream = new ReadableStream<Uint8Array>({
start(controller) {
for (const chunk of chunks) controller.enqueue(chunk);
if (end) controller.close();
},
cancel() {
state.cancelled = true;
},
});
return { stream, state };
}

function split(bytes: Uint8Array, size: number): Uint8Array[] {
const chunks: Uint8Array[] = [];
for (let offset = 0; offset < bytes.byteLength; offset += size) {
chunks.push(bytes.slice(offset, offset + size));
}
return chunks;
}

describe("runtime file exports", () => {
it("returns a complete archive even when stdout and stderr never reach EOF", async () => {
const expected = archive([
["/cpython/lib/python314.zip", "stdlib".repeat(100)],
["/cpython/lib/python3.14/os.py", "os"],
]);
const stdout = outputStream(split(expected, 7), false);
const stderr = outputStream([], false);
let stdinClosed = false;

const bytes = await readRuntimeFilesExport({
stdin: new WritableStream({ close: () => { stdinClosed = true; } }),
stdout: stdout.stream,
stderr: stderr.stream,
});

expect(bytes).toEqual(expected);
expect(stdinClosed).toBe(true);
expect(stdout.state.cancelled).toBe(true);
expect(stderr.state.cancelled).toBe(true);
});

it("reports stderr when stdout ends before the archive is complete", async () => {
const complete = archive([["/cpython/lib/python314.zip", "stdlib"]]);
await expect(readRuntimeFilesExport({
stdout: outputStream([complete.subarray(0, complete.byteLength - 1)], true).stream,
stderr: outputStream([encoder.encode("Traceback: boom")], true).stream,
})).rejects.toThrow("Traceback: boom");
});

it("reports stderr after the idle grace when only stderr reaches EOF", async () => {
const complete = archive([["/cpython/lib/python314.zip", "stdlib"]]);
const stdout = outputStream([complete.subarray(0, 20)], false);
await expect(readRuntimeFilesExport({
stdout: stdout.stream,
stderr: outputStream([encoder.encode("MemoryError")], true).stream,
}, 30)).rejects.toThrow("ended before its archive was complete: MemoryError");
expect(stdout.state.cancelled).toBe(true);
});

it("keeps reading while stdout still delivers after stderr reached EOF", async () => {
const expected = archive([["/cpython/lib/python314.zip", "stdlib".repeat(100)]]);
const chunks = split(expected, 32);
const stdout = new ReadableStream<Uint8Array>({
async pull(controller) {
await new Promise((resolve) => setTimeout(resolve, 15));
const chunk = chunks.shift();
if (chunk) controller.enqueue(chunk);
},
});
// The 100 ms grace is shorter than the whole delivery but longer than each gap between chunks.
await expect(readRuntimeFilesExport({
stdout,
stderr: outputStream([], true).stream,
}, 100)).resolves.toEqual(expected);
});

it("returns exactly the framed archive whatever trails it in the same chunk", async () => {
const expected = archive([["/cpython/lib/python314.zip", "stdlib"]]);
const trailing = new Uint8Array([...expected, ...encoder.encode("TRAILING")]);
for (const chunks of [[trailing], split(trailing, 5)]) {
const bytes = await readRuntimeFilesExport({
stdout: outputStream(chunks, false).stream,
stderr: outputStream([], false).stream,
});
expect(bytes).toEqual(expected);
}
});
});
Loading