-
Notifications
You must be signed in to change notification settings - Fork 5k
ByteStream: hand off owned buffers from on_pull/on_data and skip the adapter pull view #36736
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
robobun
wants to merge
15
commits into
main
Choose a base branch
from
farm/0c9aa2da/bytestream-zerocopy-pull
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
15 commits
Select commit
Hold shift + click to select a range
3be437a
ByteStream: hand off owned buffers from on_pull/on_data and skip the …
robobun cb20309
Add test asserting request-body chunks are right-sized; drop dead Byt…
robobun 3a36bdb
[autofix.ci] apply automated fixes
autofix-ci[bot] ecb6ac6
Trim new code comments to one line each
robobun 3088550
native-readable: skip the pull Buffer.alloc entirely for ReadyOwned s…
robobun c306505
Drop ByteStream.pending_value and wire pullParkedP reject in the test
robobun 3bfb59f
Initialize kSourceOwnsChunks in constructNativeReadable; assert chunk…
robobun 5600b4d
Test: replace write-pacing sleep with a per-chunk ack from the handler
robobun 546bf74
Test: mark every ack/handler promise observed so a failure surfaces once
robobun 8537aa0
Drop the m_sourceOwnsChunks drain-path disjunct
robobun 2adce47
ci: retrigger
robobun fe6bc93
Keep the view-copy path for native-readable.ts; Owned handoff only fo…
robobun 41ebfb8
Trim code comments to satisfy comment-cop
robobun 6829d2f
One-line the on_pull dispatch comment
robobun e179877
on_cancel: gate pending.run on pending.state, not on a stored view
robobun File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
114 changes: 114 additions & 0 deletions
114
test/js/bun/http/serve-request-body-pipeline-memory.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,114 @@ | ||
| // `for await (chunk of req.body)` (what `pipeline(req.body, writable)` runs | ||
| // via `pumpToNode`) should yield right-sized chunks. A `Bun.serve` request | ||
| // body is a push source (`ByteStream`): bytes arrive via `on_data` and are | ||
| // handed to the reader as the owning allocation. Previously `on_pull`/`on_data` | ||
| // copied into the native-source adapter's ~256-512 KiB scratch view and | ||
| // surfaced each chunk as a subarray over that whole backing, so every in-flight | ||
| // request held a ~0.5 MB `ArrayBuffer` regardless of how little data was | ||
| // actually read. | ||
|
|
||
| import { expect, test } from "bun:test"; | ||
| import { connect } from "node:net"; | ||
|
|
||
| type Seen = { len: number; backing: number; off: number }; | ||
|
|
||
| async function runUpload( | ||
| handler: (req: Request, onParked: () => void, onChunk: () => void) => Promise<Seen[]>, | ||
| bodyBytes: number, | ||
| writes: number[], | ||
| ): Promise<Seen[]> { | ||
| const { promise: handlerP, resolve: handlerDone, reject: handlerFail } = Promise.withResolvers<Seen[]>(); | ||
| const { promise: pullParkedP, resolve: pullParked, reject: pullParkedFail } = Promise.withResolvers<void>(); | ||
| // One ack per client write so the next write only leaves after the server | ||
| // has observed the previous one as a distinct on_data chunk. | ||
| const acks = writes.map(() => Promise.withResolvers<void>()); | ||
| // `fail` rejects every promise, but only the one currently awaited has a | ||
| // handler at that point; mark the rest observed so a regression surfaces as | ||
| // one failure instead of one plus N unhandled-rejection noise entries. | ||
| void handlerP.catch(() => {}); | ||
| void pullParkedP.catch(() => {}); | ||
| for (const a of acks) void a.promise.catch(() => {}); | ||
| let ackIndex = 0; | ||
| const onChunk = () => acks[ackIndex++]?.resolve(); | ||
| const fail = (e: unknown) => { | ||
| handlerFail(e); | ||
| pullParkedFail(e); | ||
| for (const a of acks) a.reject(e); | ||
| }; | ||
|
robobun marked this conversation as resolved.
coderabbitai[bot] marked this conversation as resolved.
|
||
|
|
||
| await using server = Bun.serve({ | ||
| port: 0, | ||
| async fetch(req) { | ||
| try { | ||
| const seen = await handler(req, () => queueMicrotask(() => queueMicrotask(pullParked)), onChunk); | ||
| handlerDone(seen); | ||
| } catch (e) { | ||
| fail(e); | ||
| } | ||
| return new Response("ok"); | ||
| }, | ||
| }); | ||
|
|
||
| const sock = connect({ port: server.port, host: "127.0.0.1" }); | ||
| sock.on("error", fail); | ||
| await new Promise<void>((res, rej) => { | ||
| sock.once("connect", () => res()); | ||
| sock.once("error", rej); | ||
| }); | ||
|
|
||
| try { | ||
| sock.write(`POST / HTTP/1.1\r\nHost: x\r\nContent-Length: ${bodyBytes}\r\nConnection: close\r\n\r\n`); | ||
| await pullParkedP; | ||
| for (const [i, n] of writes.entries()) { | ||
| sock.write(Buffer.alloc(n, 0x61)); | ||
| await acks[i].promise; | ||
| } | ||
| sock.end(); | ||
| return await handlerP; | ||
| } finally { | ||
| sock.destroy(); | ||
| } | ||
| } | ||
|
|
||
| function checkRightSized(seen: Seen[], expectedTotal: number, minChunks: number) { | ||
| let total = 0; | ||
| for (const { len } of seen) total += len; | ||
| expect(total).toBe(expectedTotal); | ||
| expect(seen.length).toBeGreaterThanOrEqual(minChunks); | ||
| // Every chunk is its own allocation. On main each chunk was a subarray into | ||
| // a single ~516 KiB (for await) / ~64 KiB (fromWeb) scratch view, so | ||
| // `backing >> len` and `off` advanced per chunk. | ||
| for (const { len, backing, off } of seen) { | ||
| expect({ len, backing, off }).toEqual({ len, backing: len, off: 0 }); | ||
| } | ||
| } | ||
|
|
||
| test("for await (req.body) chunks are backed by right-sized buffers, not the adapter's scratch view", async () => { | ||
| // Content-Length large enough that on_start would have sized the pull view | ||
| // at its ~512 KiB ceiling on main. | ||
| const BODY = 2 * 1024 * 1024; | ||
| const writes = [8 * 1024, 8 * 1024, BODY - 16 * 1024]; | ||
| const seen = await runUpload( | ||
| async (req, onParked, onChunk) => { | ||
| const out: Seen[] = []; | ||
| // Let the client know the first pull has parked (no body bytes yet) so | ||
| // the first chunk is resolved from on_data, not from drain(). | ||
| onParked(); | ||
| for await (const chunk of req.body!) { | ||
| out.push({ len: chunk.byteLength, backing: chunk.buffer.byteLength, off: chunk.byteOffset }); | ||
| onChunk(); | ||
| } | ||
| return out; | ||
| }, | ||
| BODY, | ||
| writes, | ||
| ); | ||
| checkRightSized(seen, BODY, writes.length); | ||
| }); | ||
|
|
||
| // `Readable.fromWeb(req.body)` is deliberately left on the copy-into-view path: | ||
| // node:stream Readable calls `_read` ahead of downstream consumption, so | ||
| // metering via the pull view's size is what keeps the producer paused while | ||
| // a chunk is still sitting in the Readable's own buffer. The C++ adapter | ||
| // above pulls once per reader.read(), so it can hand off whole buffers | ||
| // without that hazard. | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.