Skip to content
2 changes: 1 addition & 1 deletion packages/jco-transpile/test/perf.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ const ASYNC_G2G_CALL_LIMIT_NS = env.CI ? 70_000_000 : 40_000_000;
// test was introduced. These should narrow as call overhead improves.
const WASM_MODULE_CALL_OVERHEAD_RATIO_LIMIT = env.CI ? 12 : 10;
const SYNC_COMPONENT_CALL_OVERHEAD_RATIO_LIMIT = env.CI ? 1_500 : 1_000;
const ASYNC_COMPONENT_CALL_OVERHEAD_RATIO_LIMIT = env.CI ? 40_000 : 8_000;
const ASYNC_COMPONENT_CALL_OVERHEAD_RATIO_LIMIT = env.CI ? 60_000 : 8_000;
const ADDER_COMPONENT_PATH = fileURLToPath(new URL('./fixtures/components/adder.component.wasm', import.meta.url));
const ADDER_MODULE_PATH = fileURLToPath(new URL('./fixtures/runtime/adder.wat', import.meta.url));

Expand Down
206 changes: 104 additions & 102 deletions packages/preview2-shim/src/browser/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import type {
} from "../../types/http.js";
import type { Error as IoError } from "../../types/interfaces/wasi-io-error.js";
import type { Pollable } from "../../types/interfaces/wasi-io-poll.js";
import type { StreamError } from "../../types/interfaces/wasi-io-streams.js";
import type { InputStreamHandler } from "./io.js";
import { inputStreamCreate, ioErrorCreate, outputStreamCreate, pollableCreate } from "./io.js";

export { InMemoryHttpClient } from "./in-memory-http.js";
Expand Down Expand Up @@ -467,7 +469,7 @@ delete OutgoingRequest._handle;

class IncomingBody implements TypesNamespace.IncomingBody {
#finished = false;
#stream: any = undefined;
#stream: ReturnType<typeof inputStreamCreate> | null = null;

stream() {
if (!this.#stream) {
Expand Down Expand Up @@ -496,128 +498,123 @@ class IncomingBody implements TypesNamespace.IncomingBody {
let reader: ReadableStreamDefaultReader<Uint8Array> | null = null;
let readPromise: Promise<void> | null = null;
let readError: IoError | null = null;
let disposed = false;

function ensureReader() {
if (!reader && fetchResponse.body) {
reader = fetchResponse.body.getReader();
}
function ready(): boolean {
return done || (buffer !== null && bufferOffset < buffer.byteLength);
}

function startRead() {
if (readPromise || done) {
function startRead(): void {
if (readPromise || ready()) {
return;
}
ensureReader();
if (!reader) {
if (!fetchResponse.body) {
done = true;
return;
}
readPromise = reader.read().then(
(result) => {
readPromise = null;
if (result.done) {
done = true;
} else {
buffer = result.value;
bufferOffset = 0;
reader ??= fetchResponse.body.getReader();
const activeReader = reader;
readPromise = (async (): Promise<void> => {
try {
// Empty Fetch chunks do not make a WASI input stream readable.
// Keep a single read in flight and buffer at most one nonempty chunk.
while (!done) {
const result = await activeReader.read();
if (disposed) {
return;
}
if (result.done) {
done = true;
} else if (result.value.byteLength > 0) {
buffer = result.value;
bufferOffset = 0;
return;
}
}
},
(cause) => {
readPromise = null;
} catch (cause: unknown) {
done = true;
readError = ioErrorCreate(
cause instanceof Error ? cause.message : String(cause),
);
},
);
if (!disposed) {
readError = ioErrorCreate(
cause instanceof Error ? cause.message : String(cause),
);
}
} finally {
readPromise = null;
if (done) {
activeReader.releaseLock();
}
}
})();
}

function checkReadError() {
async function waitForReadable(): Promise<void> {
while (!ready()) {
startRead();
await readPromise;
}
}

function read(len: bigint): Uint8Array {
if (readError) {
throw { tag: "last-operation-failed", val: readError };
const error = readError;
readError = null;
throw { tag: "last-operation-failed", val: error } satisfies StreamError;
}
if (done && (buffer === null || bufferOffset >= buffer.byteLength)) {
throw { tag: "closed" } satisfies StreamError;
}
if (buffer === null || bufferOffset >= buffer.byteLength) {
startRead();
return new Uint8Array(0);
}
const toRead = Math.min(Number(len), buffer.byteLength - bufferOffset);
const slice = buffer.slice(bufferOffset, bufferOffset + toRead);
bufferOffset += toRead;
if (bufferOffset >= buffer.byteLength) {
buffer = null;
bufferOffset = 0;
startRead();
}
return slice;
}

function blockingRead(len: bigint): Uint8Array | Promise<Uint8Array> {
if (len === 0n || ready()) {
return read(len);
}
return waitForReadable().then(() => read(len));
}

function blockingSkip(len: bigint): bigint | Promise<bigint> {
const result = blockingRead(len);
return result instanceof Promise
? result.then((bytes) => BigInt(bytes.byteLength))
: BigInt(result.byteLength);
}

incomingBody.#stream = inputStreamCreate({
read(len: bigint) {
checkReadError();
if (done && (buffer === null || bufferOffset >= buffer.byteLength)) {
throw { tag: "closed" };
}
if (buffer !== null && bufferOffset < buffer.byteLength) {
const available = buffer.byteLength - bufferOffset;
const toRead = Math.min(Number(len), available);
const slice = buffer.slice(bufferOffset, bufferOffset + toRead);
bufferOffset += toRead;
if (bufferOffset >= buffer.byteLength) {
buffer = null;
bufferOffset = 0;
if (!done) {
startRead();
}
}
return slice;
}
throw { tag: "would-block" };
},
blockingRead(len: bigint): any {
checkReadError();
if (done && (buffer === null || bufferOffset >= buffer.byteLength)) {
throw { tag: "closed" };
}
if (buffer !== null && bufferOffset < buffer.byteLength) {
const available = buffer.byteLength - bufferOffset;
const toRead = Math.min(Number(len), available);
const slice = buffer.slice(bufferOffset, bufferOffset + toRead);
bufferOffset += toRead;
if (bufferOffset >= buffer.byteLength) {
buffer = null;
bufferOffset = 0;
if (!done) {
startRead();
}
}
return slice;
}
startRead();
const waitFor = readPromise || Promise.resolve();
return waitFor.then(() => {
checkReadError();
if (done && (buffer === null || bufferOffset >= buffer.byteLength)) {
throw { tag: "closed" };
}
if (buffer !== null && bufferOffset < buffer.byteLength) {
const available = buffer.byteLength - bufferOffset;
const toRead = Math.min(Number(len), available);
const slice = buffer.slice(bufferOffset, bufferOffset + toRead);
bufferOffset += toRead;
if (bufferOffset >= buffer.byteLength) {
buffer = null;
bufferOffset = 0;
if (!done) {
startRead();
}
}
return slice;
}
throw { tag: "closed" };
});
},
read,
// WIT declares synchronous signatures; JSPI awaits these host results
// for blocking imports. Preserve synchronous results for buffered bodies.
blockingRead: blockingRead as InputStreamHandler["blockingRead"],
blockingSkip: blockingSkip as InputStreamHandler["blockingSkip"],
subscribe() {
return pollableCreate({
ready: () =>
readError !== null ||
done ||
(buffer !== null && bufferOffset < buffer.byteLength),
wait: () => {
startRead();
return readPromise ?? Promise.resolve();
},
});
return pollableCreate({ ready, wait: waitForReadable });
},
drop() {
disposed = true;
done = true;
void reader?.cancel();
buffer = null;
readError = null;
if (reader) {
const activeReader = reader;
void activeReader
.cancel()
.catch(() => {})
.finally(() => {
activeReader.releaseLock();
});
}
},
});

Expand Down Expand Up @@ -885,6 +882,11 @@ class FutureIncomingResponse implements TypesNamespace.FutureIncomingResponse {
}
const result = this.#result;
this.#result = { tag: "err" };
// The returned response now owns the body. Dropping this consumed future
// must not abort a Fetch body that the guest is still streaming.
if (result.tag === "ok" && result.val.tag === "ok") {
this.#controller = null;
}
return result;
}

Expand Down
122 changes: 122 additions & 0 deletions packages/preview2-shim/test/browser/http-input-stream-e2e.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
import assert from "node:assert/strict";
import { mkdir, readFile, writeFile } from "node:fs/promises";
import { createServer } from "node:http";
import type { ServerResponse } from "node:http";
import { dirname, join } from "node:path";

import { componentize } from "@bytecodealliance/componentize-js";
import { transpile } from "@bytecodealliance/jco";
import { test } from "vitest";

import {
FIXTURES_WIT_DIR,
getTmpDir,
runBasicHarnessPageTest,
startTestServer,
} from "../common.js";

test("browser HTTP input stream crosses JSPI with delayed chunks and errors", async () => {
let body: ServerResponse | undefined;
let broken: ServerResponse | undefined;
const requests: string[] = [];
const api = createServer((request, response): void => {
requests.push(request.url ?? "");
response.setHeader("access-control-allow-origin", "*");
if (request.url === "/body" || request.url === "/broken") {
response.writeHead(200, { "content-type": "application/octet-stream" });
response.flushHeaders();
if (request.url === "/body") {
body = response;
} else {
broken = response;
}
return;
}
// Explicit guest-controlled gates, not timers: headers must be available
// before body bytes, and /second cannot arrive until the first chunk drains.
switch (request.url) {
case "/first":
body?.write("abcd");
break;
case "/second":
body?.write("ef");
break;
case "/finish":
body?.end();
break;
case "/abort":
broken?.destroy();
break;
default:
response.writeHead(404).end();
return;
}
response.writeHead(204).end();
});
await new Promise<void>((resolve, reject) => {
api.once("error", reject);
api.listen(0, "127.0.0.1", resolve);
});
let harness: Awaited<ReturnType<typeof startTestServer>> | undefined;
try {
const address = api.address();
assert.ok(address && typeof address !== "string");
const outDir = await getTmpDir();
harness = await startTestServer({ transpiledOutputDir: outDir });
const source = await readFile(
new URL("../fixtures/browser/http-input-stream/component.js", import.meta.url),
"utf8",
);
const sourcePath = join(outDir, "source.js");
await writeFile(sourcePath, source.replace("TEST_AUTHORITY", `127.0.0.1:${address.port}`));
const { component } = await componentize({
sourcePath,
witPath: FIXTURES_WIT_DIR,
worldName: "browser-http-fetch",
});
const { files } = await transpile(component, {
name: "component",
optimize: false,
asyncMode: "jspi",
// read and skip deliberately remain synchronous, as the WASI contract
// requires. Only polling and blocking operations may suspend.
asyncImports: [
"wasi:io/poll#[method]pollable.block",
"wasi:io/poll#poll",
"wasi:io/streams#[method]input-stream.blocking-read",
"wasi:io/streams#[method]input-stream.blocking-skip",
],
asyncExports: ["tests:p2-shim/test#run"],
outDir,
});
for (const [path, bytes] of Object.entries(files)) {
await mkdir(dirname(path), { recursive: true });
await writeFile(path, bytes);
}
const { statusJSON } = await runBasicHarnessPageTest({
browser: harness.browser,
url: `${harness.baseURL}/index.html#transpiled:component.js`,
});
assert.strictEqual(statusJSON.status, "success");
assert.strictEqual(
statusJSON.msg,
"empty reads, streaming, polling, EOF, and IO errors passed",
);
assert.deepStrictEqual(requests, [
"/body",
"/first",
"/second",
"/finish",
"/broken",
"/abort",
]);
} finally {
api.closeAllConnections();
await Promise.all([
harness?.cleanup(),
new Promise<void>((resolve, reject) =>
api.close((error) => (error ? reject(error) : resolve())),
),
]);
}
}, 120_000);
Loading
Loading