From 22b6182d8b9887a85b24e022e121cdbb8b8b0080 Mon Sep 17 00:00:00 2001 From: TakalaWang Date: Tue, 6 Oct 2026 20:40:29 +0800 Subject: [PATCH 1/2] fix(runtime): finish Python runtime-file export without waiting for EOF @wasmer/sdk 0.10.0 runs `Command.run()` in a dedicated thread-pool task that sends the exit code and then calls `thread_pool.close()` before the task drops the WASI runner, which owns the stdout/stderr pipe writers. The main-thread scheduler handles `Close` by dropping every `WorkerHandle`, and each drop terminates its Worker. When that termination lands before the task has dropped the runner, stdout and stderr never reach EOF, so `Instance.wait()` (which joins stdout EOF, stderr EOF and the exit code) never settles. The runtime-files stage then idles with no worker threads until the 300 s preparation deadline. Instrumented stalled stages showed all 10,695,683 archive bytes delivered on stdout, both streams still open with a pull pending, and the dedicated worker terminated before it reported idle. CPU contention from concurrent engines widens the window. The WOJFS002 archive is self-delimiting, so read stdout directly and finish once the archive terminator arrives; the parent still verifies the archive SHA-256 before use. A stdout EOF before completion reports stderr. The browser runner Worker used the same `Instance.wait()` export path and now shares the helper. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 4 ++ src/runner/runtime-files.test.ts | 59 ++++++++++++++++++++++++- src/runner/runtime-files.ts | 71 ++++++++++++++++++++++++++++++ src/runtime/runner.worker.ts | 19 ++++---- src/server/server-runner-stage.mjs | 12 ++--- 5 files changed, 151 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 30877b6..3d52fb0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,10 @@ 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 occasionally stalling until its 300 s timeout. The Wasmer + SDK can terminate its workers before stdout/stderr reach EOF even though the whole runtime + file archive has arrived, so server and browser preparation now finish once the + self-delimiting archive is complete instead of waiting for `Instance.wait()`. ## 0.2.3 - 2026-10-05 diff --git a/src/runner/runtime-files.test.ts b/src/runner/runtime-files.test.ts index 318731f..6fb7986 100644 --- a/src/runner/runtime-files.test.ts +++ b/src/runner/runtime-files.test.ts @@ -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(); @@ -59,3 +63,56 @@ describe("runtime file archives", () => { ); }); }); + +function outputStream(chunks: readonly Uint8Array[], end: boolean) { + const state = { cancelled: false }; + const stream = new ReadableStream({ + 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"); + }); +}); diff --git a/src/runner/runtime-files.ts b/src/runner/runtime-files.ts index 49551b9..7ac8680 100644 --- a/src/runner/runtime-files.ts +++ b/src/runner/runtime-files.ts @@ -96,6 +96,77 @@ export function decodeRuntimeFiles(archive: Uint8Array): Record; + readonly stderr: ReadableStream; +} + +/** + * Read an export until stdout holds a complete archive. `@wasmer/sdk` `Instance.wait()` also + * waits for stdout/stderr EOF, which its thread-pool teardown can drop after all output arrived. + */ +export async function readRuntimeFilesExport(exporter: RuntimeFilesExportProcess): Promise { + await exporter.stdin?.close().catch(() => undefined); + const stdout = exporter.stdout.getReader(); + const stderr = exporter.stderr.getReader(); + const diagnostics = readToEnd(stderr).catch(() => new Uint8Array()); + const chunks: Uint8Array[] = []; + let received = 0; + let required = MAGIC.byteLength + HEADER_BYTES; + try { + while (true) { + const { done, value } = await stdout.read(); + if (done) { + const detail = new TextDecoder().decode(await diagnostics); + throw new Error(`The runtime file export ended before its archive was complete: ${detail}`); + } + chunks.push(value); + received += value.byteLength; + if (received < required) continue; + const archive = concatenate(chunks); + chunks.splice(0, chunks.length, archive); + const missing = missingArchiveBytes(archive); + if (missing === 0) return archive; + required = received + missing; + } + } finally { + await Promise.allSettled([stdout.cancel(), stderr.cancel()]); + } +} + +function missingArchiveBytes(archive: Uint8Array): number { + const view = new DataView(archive.buffer, archive.byteOffset, archive.byteLength); + let offset = MAGIC.byteLength; + while (offset + HEADER_BYTES <= archive.byteLength) { + const pathLength = view.getUint32(offset, true); + const dataLength = Number(view.getBigUint64(offset + 4, true)); + offset += HEADER_BYTES; + if (pathLength === 0 && dataLength === 0) return 0; + offset += pathLength + dataLength; + } + return offset + HEADER_BYTES - archive.byteLength; +} + +async function readToEnd(reader: ReadableStreamDefaultReader): Promise { + const chunks: Uint8Array[] = []; + while (true) { + const { done, value } = await reader.read(); + if (done) return concatenate(chunks); + chunks.push(value); + } +} + +function concatenate(chunks: readonly Uint8Array[]): Uint8Array { + const output = new Uint8Array(chunks.reduce((total, chunk) => total + chunk.byteLength, 0)); + let offset = 0; + for (const chunk of chunks) { + output.set(chunk, offset); + offset += chunk.byteLength; + } + return output; +} + export async function verifyAndDecodeRuntimeFiles( archive: Uint8Array, expectedSha256: string, diff --git a/src/runtime/runner.worker.ts b/src/runtime/runner.worker.ts index ebeb2f8..e0c7cc8 100644 --- a/src/runtime/runner.worker.ts +++ b/src/runtime/runner.worker.ts @@ -46,6 +46,7 @@ import { openOptionalRuntimeFilesCache, restoreOrExportRuntimeFiles, } from "@/src/runtime/runtime-files-cache"; +import { readRuntimeFilesExport } from "@/src/runner/runtime-files"; import { PackageHandleCache, WasmerPackageHandle, @@ -275,7 +276,7 @@ async function exportPackageFileSystem( }, async () => { const lease = await acquirePackage(request.packageSpecifier); - const output = await withHandleLease( + return withHandleLease( lease, (pkg) => withWasmerCommand(pkg, request.command, async (command) => { const instance = await command.run({ @@ -286,15 +287,17 @@ async function exportPackageFileSystem( PYTHONDONTWRITEBYTECODE: "1", }, }); - return instance.wait(); + try { + return await readRuntimeFilesExport(instance); + } catch (error) { + throw new Error( + `Unable to export runtime files from ${request.packageSpecifier}: ${error instanceof Error ? error.message : String(error)}`, + ); + } finally { + instance.free(); + } }), ); - if (!output.ok) { - throw new Error( - `Unable to export runtime files from ${request.packageSpecifier}: exit ${output.code}: ${output.stderr}`, - ); - } - return output.stdoutBytes.slice(); }, ); })(); diff --git a/src/server/server-runner-stage.mjs b/src/server/server-runner-stage.mjs index a3eb1ff..57d45b8 100644 --- a/src/server/server-runner-stage.mjs +++ b/src/server/server-runner-stage.mjs @@ -11,6 +11,7 @@ import { PYTHON_PACKAGE, PYTHON_PACKAGE_SHA256, } from "../core/toolchains.ts"; +import { readRuntimeFilesExport } from "../runner/runtime-files.ts"; const MAX_INPUT_BYTES = 1024 * 1024; const MAX_RESULT_BYTES = 256 * 1024 * 1024; @@ -59,14 +60,15 @@ try { PYTHONDONTWRITEBYTECODE: "1", }, })); - const output = await withProcessKeepalive(instance.wait()); - if (!output.ok) { + try { + bytes = await withProcessKeepalive(readRuntimeFilesExport(instance)); + } catch (error) { throw new Error( - `Unable to export runtime files from ${input.request.packageSpecifier}: ` - + `exit ${output.code}: ${boundedText(output.stderr)}`, + `Unable to export runtime files from ${input.request.packageSpecifier}: ${boundedText(errorText(error))}`, ); + } finally { + instance.free(); } - bytes = output.stdoutBytes.slice(); } if (bytes.byteLength > MAX_RESULT_BYTES) { From d6e02b91bd0d78a6e320397b1a5bd5fa4c37955d Mon Sep 17 00:00:00 2001 From: JacobLinCool Date: Fri, 9 Oct 2026 06:33:47 +0800 Subject: [PATCH 2/2] fix(runtime): fail fast on incomplete exports and stop waiting for TypeScript EOF Reading the runtime-file export until the archive completes left the failure path waiting on stdout EOF, which the same SDK 0.10 teardown race drops. A guest that died mid-archive then sat until the 300 s preparation deadline and lost its traceback. Either stream's EOF means the guest has exited, so the other stream now gets a 2 s idle grace (reset by new bytes) before the read stops and the export fails with stderr. The returned archive is also trimmed to its terminator so chunking cannot change it. TypeScript compilation still awaited Instance.wait() and could stall until the 120 s build timeout (1 of 88 runs at 12x concurrency). It now reads stdout until it holds the driver's one JSON response, then collects stderr until EOF or the same grace. The driver writes that response last, so a complete response stands in for the exit code wait() would report. Both readers share readProcessMessage in src/core/process-output.ts. --- CHANGELOG.md | 10 +-- src/compiler/wasmer-engine.ts | 27 ++++--- src/core/process-output.test.ts | 81 +++++++++++++++++++++ src/core/process-output.ts | 119 +++++++++++++++++++++++++++++++ src/runner/runtime-files.test.ts | 39 ++++++++++ src/runner/runtime-files.ts | 74 +++++-------------- 6 files changed, 274 insertions(+), 76 deletions(-) create mode 100644 src/core/process-output.test.ts create mode 100644 src/core/process-output.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 3d52fb0..12e6838 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,10 +16,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 occasionally stalling until its 300 s timeout. The Wasmer - SDK can terminate its workers before stdout/stderr reach EOF even though the whole runtime - file archive has arrived, so server and browser preparation now finish once the - self-delimiting archive is complete instead of waiting for `Instance.wait()`. +- 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 diff --git a/src/compiler/wasmer-engine.ts b/src/compiler/wasmer-engine.ts index ca662d7..f330317 100644 --- a/src/compiler/wasmer-engine.ts +++ b/src/compiler/wasmer-engine.ts @@ -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, @@ -307,7 +307,7 @@ interface TypeScriptWasiResponse { files: Record; } -async function transpileScriptProject(project: Project, requestId: string): Promise<{ files: Record; output: Output; response?: TypeScriptWasiResponse }> { +async function transpileScriptProject(project: Project, requestId: string): Promise<{ files: Record; stderr: string; response?: TypeScriptWasiResponse }> { const scriptFiles = scriptSourceFiles(project); const emittedFiles = emittedSourceFiles(project); const dependencyFiles = npmDependencyFiles(project); @@ -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( + instance, + completeJsonObject, + { collectStderr: true }, + ).finally(() => instance.free()); + const response = output.message; const files: Record = {}; if (response) { for (const outputPath of outputPaths) { @@ -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 { @@ -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, { diff --git a/src/core/process-output.test.ts b/src/core/process-output.test.ts new file mode 100644 index 0000000..2a812f2 --- /dev/null +++ b/src/core/process-output.test.ts @@ -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; + const stream = new ReadableStream({ 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: "" }); + }); +}); diff --git a/src/core/process-output.ts b/src/core/process-output.ts new file mode 100644 index 0000000..a15ae37 --- /dev/null +++ b/src/core/process-output.ts @@ -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; + readonly stderr: ReadableStream; +} + +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 { + /** 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( + child: StreamingProcess, + parse: (stdout: Uint8Array) => T | undefined, + options: ProcessMessageOptions = {}, +): Promise> { + 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 | undefined = nextStdout(); + let stderrRead: Promise | 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 => 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 | 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(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), + }; +} diff --git a/src/runner/runtime-files.test.ts b/src/runner/runtime-files.test.ts index 6fb7986..20ba179 100644 --- a/src/runner/runtime-files.test.ts +++ b/src/runner/runtime-files.test.ts @@ -115,4 +115,43 @@ describe("runtime file exports", () => { 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({ + 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); + } + }); }); diff --git a/src/runner/runtime-files.ts b/src/runner/runtime-files.ts index 7ac8680..6134759 100644 --- a/src/runner/runtime-files.ts +++ b/src/runner/runtime-files.ts @@ -1,5 +1,6 @@ import { WASM_OJ_STORAGE } from "../core/contract.ts"; import { sha256Hex } from "../core/hash.ts"; +import { readProcessMessage, type StreamingProcess } from "../core/process-output.ts"; const MAGIC = new TextEncoder().encode("WOJFS002"); const HEADER_BYTES = 12; @@ -96,75 +97,32 @@ export function decodeRuntimeFiles(archive: Uint8Array): Record; - readonly stderr: ReadableStream; -} - /** - * Read an export until stdout holds a complete archive. `@wasmer/sdk` `Instance.wait()` also - * waits for stdout/stderr EOF, which its thread-pool teardown can drop after all output arrived. + * Read an export until stdout holds a complete archive and return exactly its framed bytes. + * `@wasmer/sdk` `Instance.wait()` can hang after all output arrived; see readProcessMessage. */ -export async function readRuntimeFilesExport(exporter: RuntimeFilesExportProcess): Promise { - await exporter.stdin?.close().catch(() => undefined); - const stdout = exporter.stdout.getReader(); - const stderr = exporter.stderr.getReader(); - const diagnostics = readToEnd(stderr).catch(() => new Uint8Array()); - const chunks: Uint8Array[] = []; - let received = 0; - let required = MAGIC.byteLength + HEADER_BYTES; - try { - while (true) { - const { done, value } = await stdout.read(); - if (done) { - const detail = new TextDecoder().decode(await diagnostics); - throw new Error(`The runtime file export ended before its archive was complete: ${detail}`); - } - chunks.push(value); - received += value.byteLength; - if (received < required) continue; - const archive = concatenate(chunks); - chunks.splice(0, chunks.length, archive); - const missing = missingArchiveBytes(archive); - if (missing === 0) return archive; - required = received + missing; - } - } finally { - await Promise.allSettled([stdout.cancel(), stderr.cancel()]); +export async function readRuntimeFilesExport( + exporter: StreamingProcess, + idleGraceMs?: number, +): Promise { + const { message, stderr } = await readProcessMessage(exporter, completeArchive, { idleGraceMs }); + if (message === undefined) { + throw new Error(`The runtime file export ended before its archive was complete: ${stderr}`); } + return message; } -function missingArchiveBytes(archive: Uint8Array): number { - const view = new DataView(archive.buffer, archive.byteOffset, archive.byteLength); +function completeArchive(bytes: Uint8Array): Uint8Array | undefined { + const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength); let offset = MAGIC.byteLength; - while (offset + HEADER_BYTES <= archive.byteLength) { + while (offset + HEADER_BYTES <= bytes.byteLength) { const pathLength = view.getUint32(offset, true); const dataLength = Number(view.getBigUint64(offset + 4, true)); offset += HEADER_BYTES; - if (pathLength === 0 && dataLength === 0) return 0; + if (pathLength === 0 && dataLength === 0) return bytes.slice(0, offset); offset += pathLength + dataLength; } - return offset + HEADER_BYTES - archive.byteLength; -} - -async function readToEnd(reader: ReadableStreamDefaultReader): Promise { - const chunks: Uint8Array[] = []; - while (true) { - const { done, value } = await reader.read(); - if (done) return concatenate(chunks); - chunks.push(value); - } -} - -function concatenate(chunks: readonly Uint8Array[]): Uint8Array { - const output = new Uint8Array(chunks.reduce((total, chunk) => total + chunk.byteLength, 0)); - let offset = 0; - for (const chunk of chunks) { - output.set(chunk, offset); - offset += chunk.byteLength; - } - return output; + return undefined; } export async function verifyAndDecodeRuntimeFiles(