Skip to content
Open
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
10 changes: 4 additions & 6 deletions src/io/PipeReader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -781,7 +781,7 @@ impl PosixBufferedReader {
} else if streaming {
vtable.on_read_chunk(
Chunk::Scratch(&scratch[..filled]),
Self::read_state(stop.as_ref(), received_hup),
Self::read_state(stop.as_ref()),
)
} else {
// SAFETY: caller contract; borrow ends at `;`.
Expand All @@ -802,7 +802,7 @@ impl PosixBufferedReader {
// Moved out so a re-entrant read cannot alias or reallocate it under the consumer.
// SAFETY: caller contract; borrow ends at `;`.
let mut buffer = unsafe { mem::take(&mut (*this)._buffer) };
let state = Self::read_state(stop.as_ref(), received_hup);
let state = Self::read_state(stop.as_ref());
if matches!(stop, Some(Stop::Eof | Stop::OverBudget | Stop::Error(_))) {
vtable.on_read_chunk(Chunk::Owned(buffer), state)
} else {
Expand Down Expand Up @@ -895,13 +895,11 @@ impl PosixBufferedReader {
}
}

fn read_state(stop: Option<&Stop>, received_hup: bool) -> ReadState {
fn read_state(stop: Option<&Stop>) -> ReadState {
match stop {
Some(Stop::Eof | Stop::OverBudget) => ReadState::Eof,
Some(Stop::WouldBlock) => ReadState::Drained,
Some(Stop::Error(_)) => ReadState::Progress,
None if received_hup => ReadState::Eof,
None => ReadState::Progress,
Some(Stop::Error(_)) | None => ReadState::Progress,
Comment thread
claude[bot] marked this conversation as resolved.
}
}

Expand Down
180 changes: 178 additions & 2 deletions test/js/bun/http/bun-serve-file.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,23 @@
import type { Server } from "bun";
import type { Server, SocketHandler } from "bun";
import { afterAll, beforeAll, describe, expect, it, mock, test } from "bun:test";
import { bunEnv, bunExe, isASAN, isLinux, isMacOS, isWindows, rmScope, rss, tempDir, tempDirWithFiles } from "harness";
import {
bunEnv,
bunExe,
isASAN,
isLinux,
isMacOS,
isWindows,
libcPathForDlopen,
rmScope,
rss,
tempDir,
tempDirWithFiles,
} from "harness";
import { mkfifo } from "mkfifo";
import { closeSync, openSync, unlinkSync, writeSync } from "node:fs";
import { open as fsOpen } from "node:fs/promises";
import { join } from "node:path";
import { positionDependentBytes, unixSockets } from "socketpair";

const LARGE_SIZE = 1024 * 1024 * 8;
const files = {
Expand Down Expand Up @@ -1149,6 +1162,169 @@ process.exit(0);
);
}

// One wakeup of the server can report the last bytes of a pollable body and
// the hangup of its writer together. The response must still carry every byte
// and end as a complete message, and a read error behind the bytes must not
// end it as one. Each row queues its bytes on one end of a socketpair and
// closes that end, in one turn of this thread, while the server waits on its
// poll for the other end. A host whose limit for a socket buffer is too low for
// the largest row skips them all.
const sockets = isLinux || isMacOS ? unixSockets(libcPathForDlopen()) : undefined;
const skipHungUpSockets = !sockets || sockets.limitIsBelow(400_000);
describe.skipIf(skipHungUpSockets)("a file response whose socket hangs up with bytes unread", () => {
const first = Buffer.from("first");

function parse(wire: Buffer) {
const headEnd = wire.indexOf("\r\n\r\n");
if (headEnd === -1) return { status: "", body: Buffer.alloc(0), complete: false };
const head = wire.subarray(0, headEnd).toString("latin1");
const status = head.slice(0, head.indexOf("\r\n"));
const rest = wire.subarray(headEnd + 4);
const declared = /^content-length:\s*(\d+)/im.exec(head);
if (declared) return { status, body: rest, complete: rest.length >= Number(declared[1]) };
const chunks: Buffer[] = [];
let at = 0;
for (;;) {
const lineEnd = rest.indexOf("\r\n", at);
const size = lineEnd === -1 ? NaN : parseInt(rest.subarray(at, lineEnd).toString("latin1"), 16);
if (Number.isNaN(size)) return { status, body: Buffer.concat(chunks), complete: false };
if (size === 0) return { status, body: Buffer.concat(chunks), complete: rest.length >= lineEnd + 4 };
chunks.push(rest.subarray(lineEnd + 2, lineEnd + 2 + size));
at = lineEnd + 2 + size + 2;
}
}

// One GET from a raw client on this thread. It ends when the message is complete or the connection is gone.
// `whenBodyStarts` runs once the first bytes of the body are in.
async function get(server: Server<undefined>, unix: string | undefined, whenBodyStarts?: () => void) {
const { promise, resolve } = Promise.withResolvers<void>();
let wire = Buffer.alloc(0);
let closed: string | undefined;
const socket: SocketHandler<undefined> = {
open(socket) {
socket.write("GET / HTTP/1.1\r\nHost: x\r\n\r\n");
},
data(_socket, chunk) {
wire = Buffer.concat([wire, chunk]);
const { body, complete } = parse(wire);
if (body.length >= first.length) {
whenBodyStarts?.();
whenBodyStarts = undefined;
}
if (complete) resolve();
},
close(_socket, err) {
closed = (err as { code?: string } | undefined)?.code ?? "end";
resolve();
},
error(_socket, err) {
closed = (err as { code?: string }).code;
resolve();
},
};
await using _ = await (unix
? Bun.connect({ unix, socket })
: Bun.connect({ hostname: "127.0.0.1", port: server.port!, socket }));
await promise;
return { ...parse(wire), closed };
}

// Runs `fn` inside the read loop of another reader, which holds the 256 KiB buffer that the read loops of one
// thread share. A read loop that runs in there reads into a buffer of its own and delivers at 128 KiB.
async function insideAnotherReadLoop(fn: () => unknown) {
using other = sockets!.pair();
const reader = Bun.file(other.source).stream().getReader();
const read = reader.read();
writeSync(other.peer, "x");
await read;
try {
await fn();
} finally {
await reader.cancel();
}
}

test.each([
// One read takes 256 KiB at most, so the first read does not reach the end.
[270_000, "after its first bytes", "tcp", "the event loop"],
[270_000, "on an empty source", "tcp", "the event loop"],
// A read that fills half of its buffer stops there. The next read finds the end and no bytes.
[200_000, "after its first bytes", "tcp", "the event loop"],
[200_000, "on an empty source", "tcp", "the event loop"],
// A unix socket takes less than one read of 256 KiB, and the client is on this thread, so the server writes
// the second delivery to a response that pushes back.
[400_000, "after its first bytes", "unix", "the event loop"],
[131_072, "on an empty source", "tcp", "the callback of another reader"],
[131_073, "on an empty source", "tcp", "the callback of another reader"],
[200_000, "on an empty source", "tcp", "the callback of another reader"],
] as const)("%d bytes, response started %s, over %s, hangup handled by %s", async (length, started, over, by) => {
using dir = tempDir("serve-socket-hangup", {});
using pair = sockets!.pair();
const unix = over === "unix" ? join(String(dir), "server.sock") : undefined;
const bytes = positionDependentBytes(length);
let queued = -1;
const { promise: hungUp, resolve, reject } = Promise.withResolvers<void>();
// Runs when the server has started the response and waits on its poll.
const hangUp = () => {
if (by === "the event loop") {
queued = pair.hangUp(bytes);
return resolve();
}
// `resolves` runs the event loop until the promise settles, so the server handles the wakeup from in here.
insideAnotherReadLoop(async () => {
queued = pair.hangUp(bytes);
await expect(response).resolves.toBeDefined();
}).then(resolve, reject);
};

if (started === "after its first bytes") writeSync(pair.peer, first);
await using server = Bun.serve({
...(unix ? { unix } : { port: 0, hostname: "127.0.0.1" }),
fetch() {
if (started === "on an empty source") setImmediate(hangUp);
return new Response(Bun.file(pair.source));
},
});
const response = get(server, unix, started === "after its first bytes" ? hangUp : undefined);
await hungUp;
const { status, body, complete, closed } = await response;

const expected = started === "after its first bytes" ? Buffer.concat([first, bytes]) : bytes;
expect({ queued, status, complete, closed, received: body.length, intact: body.equals(expected) }).toEqual({
queued: length,
status: "HTTP/1.1 200 OK",
complete: true,
closed: undefined,
received: expected.length,
intact: true,
});
});

// On Linux a unix socket whose peer closes with unread input is reset: the read after the peer's last bytes
// fails with ECONNRESET. The bytes are enough to end one read before that error (half of the read buffer).
// The server answers a failed body with a reset, so the client must not get a complete message.
test.skipIf(!isLinux)("a read error behind the bytes resets the connection", async () => {
using pair = sockets!.pair();
writeSync(pair.source, "input that the peer never reads");
let queued = -1;
await using server = Bun.serve({
port: 0,
hostname: "127.0.0.1",
fetch() {
setImmediate(() => (queued = pair.hangUp(positionDependentBytes(160_000))));
return new Response(Bun.file(pair.source));
},
});
const { status, complete, closed } = await get(server, undefined);
expect({ queued, status, complete, closed }).toEqual({
queued: 160_000,
status: "HTTP/1.1 200 OK",
complete: false,
closed: "ECONNRESET",
});
});
});

// A FIFO's stat size is 0, but the body length is unknown until EOF. Writing
// Content-Length from the stat size and then streaming the pipe to EOF puts
// body bytes on the wire past the declared length; on a keep-alive connection
Expand Down
31 changes: 31 additions & 0 deletions test/js/bun/spawn/spawn-stdout-unread-at-hangup-child-fixture.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

64 changes: 64 additions & 0 deletions test/js/bun/spawn/spawn-stdout-unread-at-hangup-fixture.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

47 changes: 47 additions & 0 deletions test/js/bun/spawn/spawn.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,15 +10,19 @@ import {
isBroken,
isDebug,
isLinux,
isMacOS,
isPosix,
isWindows,
libcPathForDlopen,
shellExe,
tempDir,
tmpdirSync,
withoutAggressiveGC,
} from "harness";
import { mkfifo } from "mkfifo";
import { closeSync, fstatSync, openSync, readFileSync, readSync, realpathSync, rmSync, writeFileSync } from "node:fs";
import path, { join } from "path";
import { unixSockets } from "socketpair";

let tmp: string;

Expand Down Expand Up @@ -935,6 +939,49 @@ describe.skipIf(isWindows)("stdout reader of an unref'd child and process lifeti
});
});

// One wakeup of the parent can report the last bytes of a child's stdout and
// the hangup together. The consumer of that stdout must still get every byte,
// also when more bytes are unread than one read takes (256 KiB). The child
// queues the bytes on its stdout socket and closes it while the parent is
// blocked in a synchronous read, so the parent polls the socket only after
// the hangup. A host whose limit for a socket buffer is too low skips the rows.
const skipHungUpStdout = !(isLinux || isMacOS) || unixSockets(libcPathForDlopen()).limitIsBelow(270_000);
describe.skipIf(skipHungUpStdout)("stdout bytes that are still unread when the child hangs up", () => {
it.concurrent.each([
["Bun.write(file, proc.stdout)", "write"],
["Bun.write(file, new Response(proc.stdout))", "write-response"],
["new HTMLRewriter().transform(new Response(proc.stdout))", "rewriter"],
// Guards: a wrong end-of-stream label does not cut these short, or not in every run.
["Bun.spawn({ stdin: proc.stdout })", "stdin"],
["fetch(url, { body: proc.stdout })", "fetch"],
["a shell capture", "shell"],
])("%s gets all of them", async (_, consumer) => {
Comment thread
robobun marked this conversation as resolved.
const length = 270_000;
using dir = tempDir("spawn-stdout-unread-at-hangup", {});
const fifo = join(String(dir), "hung-up.fifo");
mkfifo(fifo);

await using proc = spawn({
cmd: [
bunExe(),
join(import.meta.dir, "spawn-stdout-unread-at-hangup-fixture.ts"),
consumer,
fifo,
join(String(dir), "out.bin"),
String(length),
],
env: { ...bunEnv, LIBC_PATH: libcPathForDlopen() },
cwd: String(dir),
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(stderr).toBe("");
expect(JSON.parse(stdout)).toEqual({ queued: length, received: length, intact: true });
expect(exitCode).toBe(0);
});
});

describe("unref() + .exited with nothing else ref'd (Windows)", () => {
// Windows: with only an unref'd uv_process_t left, uv_run() used to skip its
// body and never dequeue the IOCP exit packet, so these children busy-spun
Expand Down
Loading
Loading