Bun.serve: defer pipelined HTTP/1.1 requests behind an async response instead of closing the socket - #33664
Bun.serve: defer pipelined HTTP/1.1 requests behind an async response instead of closing the socket#33664robobun wants to merge 9 commits into
Conversation
|
Updated 11:15 PM PT - Jul 25th, 2026
❌ @robobun, your commit e8de2ab has some failures in 🧪 To try this PR locally: bunx bun-pr 33664That installs a local version of the PR into your bun-33664 --bun |
|
Warning Review limit reached
Next review available in: 8 minutes Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available. How can I continue?After more reviews become available, a review can be triggered using the To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews. How do review limits work?CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability. For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window. Please refer docs for additional details. Review details⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (10)
WalkthroughThis PR adds deferred HTTP/1.1 pipeline parsing to uWebSockets integration in Bun. When an async response is pending, subsequent pipelined request bytes are buffered instead of dispatched immediately, then drained and re-parsed once the response completes. Tests validate ordering behavior. ChangesDeferred pipeline parsing and draining
Possibly related PRs
Suggested reviewers: Comment |
|
This PR may be a duplicate of:
🤖 Generated with Claude Code |
|
Acknowledged re #32868: the Related section at the bottom of the description covers it. Short version: #32868 delivers the first response and closes (both Bun.serve and node:http); this PR buffers and serves every pipelined request in order for Bun.serve only (RFC 9112 9.3.2), leaving node:http to #32868 / #32488 so |
|
CI status on e8de2ab (build #81833): 120 passed, 68 expired, 1 timed_out ( Eight review rounds (986f600 .. e8de2ab) addressed the same class: Rust/C++ frames reading or writing
One open thread (round nine): Not retriggering CI (mass job expiry is an agent-capacity issue a retrigger would hit again, and the open thread needs a decision first). |
|
Confirming this also fixes the sendfile(2) variant: a pipelined request arriving mid-body of an 8 MB import * as fs from "node:fs"; import * as net from "node:net";
fs.writeFileSync("/tmp/sf-big.bin", Buffer.alloc(8 << 20, 0x41));
const srv = Bun.serve({ port: 0, hostname: "127.0.0.1",
routes: { "/big": Bun.file("/tmp/sf-big.bin"), "/m": new Response("MARKER") },
fetch: () => new Response("x", { status: 404 }) });
// raw client: GET /big, then mid-body GET /m
// main: got=2691072 of cl=8388608, closed=true, marker=false
// this branch: got=8388608 of cl=8388608, closed=false, marker=true
That branch also carries the minimal "finish then close" alternative (set Related: #32868 (minimal variant), #33698 (sendfile |
…of closing the socket When a pipelined HTTP/1.1 request arrives while the previous response on the connection is still pending (async handler), uWS was closing the socket the moment it saw the second request head. The in-flight response was discarded (onAborted fired, resp null), and the pipelined request was never dispatched: both requests produced zero bytes on the wire. The canonical trigger is a relay handler that forwards req.body into an outgoing fetch(): the handler is genuinely async, so by the time the parse loop reaches the pipelined GET behind the CL-delimited body, HTTP_RESPONSE_PENDING is still set and the !IsNodeHttp branch closed the socket. But any async handler exhibits this, e.g. two pipelined GETs with an `await Bun.sleep()` in the handler. Fix: while a Bun.serve response is pending, break the parse loop before getHeaders mutates the next request's bytes, buffer them verbatim in a per-connection pipelinedBuffer, and pause reads. markDone() clears the defer flag and replays the buffer through onData so the next request is dispatched once the response is written. deferPipeline is re-derived from HTTP_RESPONSE_PENDING after the body fin callback so a handler that responds synchronously from inside the body (await req.text()) keeps the existing sync-pipelining fast path. node:http has its own pipelined-response queue and is unchanged (every new path is gated on !IsNodeHttp).
9649db6 to
a6c2a29
Compare
|
Rebased onto current main (a6c2a29). Same buffer-and-replay approach as before with two refinements: |
…estructed HttpResponseData Two findings from review, both reachable from network input: 1. markDone() -> replayPipelinedRequests() -> onData can synchronously close or adopt the socket (parse error on the buffered bytes, Connection: close on the replayed request, WebSocket upgrade), which runs ~HttpResponseData(). Every caller of markDone() (internalEnd both paths, uws_res_end_sendfile, uws_res_end_without_body) still reads the response data afterwards (shouldCloseConnection(), hasFullyDrained(), resetTimeout()). Fix: markDone() only clears state and deferPipeline. Replay is invoked by each caller as its final action, after all reads. Replay now also uncorks first so the just-completed response reaches the wire even when the replayed bytes are a parse error (onData's error path closes with uncorkWithoutSending). 2. callOnWritable() restores the borrowed callback whenever onWritable is non-null afterwards. A replayed request's handler can install its own onWritable (backpressure), which the restore then overwrites with the completed response's stale callback while writableUserData already points at the new request's context. Fix: compare against the placeholder function pointer, not non-null. Also: HttpContext::onWritable bails on us_socket_is_closed after callOnWritable, since replay inside the borrowed callback can close the socket before the state reads that follow. New test: async request pipelined with a malformed request (HTTP/9.9); the first response must reach the wire before the 505 and close.
… in-place adopt via socket kind
The previous commit moved replay to the end of internalEnd(), but tryEnd()
reads hasResponded() after internalEnd() returns and uws_res_try_end then
calls clearOnWritableAndAborted(), both on a response the replay may have
destructed or adopted. Move replay out of internalEnd() entirely and invoke
it only as the final action of each C ABI wrapper (uws_res_end,
uws_res_end_stream, uws_res_try_end, uws_res_end_sendfile,
uws_res_end_without_body), after every read and write those wrappers make.
For uws_res_try_end, replay is gated on pair.second (hasResponded), which
tryEnd() computed before replay ran.
An in-place WebSocket adopt (us_poll_resize returns the same pointer when
the new ext is not larger, which it is not: HttpResponseData carries three
std::strings) leaves is_closed false, so us_socket_is_closed alone does not
detect it. Add HttpResponse::isNoLongerHttp() that also checks
us_socket_kind() against HttpContext::socketKind(), and use it in
replayPipelinedRequests(), callOnWritable() (before touching this->onWritable
after the borrowed callback), and HttpContext::onWritable (after
callOnWritable returns).
Test fixes: add error listeners to readPipeline and the separate-write
test's socket so a regressed build fails on the toEqual diff rather than an
unhandled ECONNRESET; slice the parse-error test's 505 check from the exact
end of the first response's body instead of raw.indexOf("A"), which would
match inside Date: ... Apr/Aug ... headers.
…lers clearing the replayed request's handlers Two further findings: 1. replayPipelinedRequests() read pipelinedBuffer.empty() before checking isNoLongerHttp(), but internalEnd()'s own shouldCloseConnection gate can close and destruct HttpResponseData before the C ABI wrapper calls replay. Hoist isNoLongerHttp() to the top (it reads only us_socket_t flags) and just return; the destructor already freed the buffer. 2. After uws_res_try_end replays and the next request's handler installs its own onAborted/inStream/onWritable, the Rust callers' detach_response() / on_response_complete() called resp.clear_aborted()/clear_on_data()/ clear_on_writable()/clear_timeout(), nulling the next request's handlers. markDone()+clearOnWritableAndAborted() inside uws_res_try_end already nulled ours, so those FFI clears were always redundant after a successful end; now they are actively wrong. Add detach_response_after_end() that drops only the local handle and flags, and use it at the four RequestContext try_end-success sites. Drop the clears from StaticRoute::on_response_complete and FileRoute::on_response_complete for the same reason.
…'replayed from markDone' comments - on_file_stream_complete (the Response(Bun.file) completion path) now uses detach_response_after_end() so it does not null the replayed request's onAborted. - HTTPServerWritable::on_writable no longer calls res.clear_on_writable() after send_readable() reached try_end/end: that runs from callOnWritable (isParsingHttp false), so the wrapper replays and the clear would null the replayed request's onWritable (or write into a destructed HttpResponseData when the replayed bytes parse-error or upgrade). - Three comments in HttpParser.h / HttpContext.h that still said 'replayed from markDone()' now point at replayPipelinedRequests().
Http3Response::markDone() deliberately leaves onAborted armed so on_stream_close can notify the holder when the QUIC stream is freed; removing resp.clear_aborted() from StaticRoute/FileRoute:: on_response_complete meant on_stream_close fired it again, double- derefing the route (serve-http3.test.ts core dump on alpine aarch64). HTTP/3 has no pipelining replay, so the clears are always safe there. Gate the post-end clears on the response variant: H3 keeps them, HTTP/1 skips them (markDone already nulled ours and the wrapper replayed the next pipelined request). detach_response_after_end() falls through to detach_response() for HTTP3=true; streams.rs restores res.clear_on_writable() under the HTTP3 const.
On the corked async small-body path, internalEnd()'s own close gate is skipped (isCorked() is true, then its uncork() releases the slot so HttpResponse::cork() early-returns without reaching its close gate), so the resume() here left a Connection: close socket open until the client FINs or idle timeout. shutdown()+close() when fully drained instead. The non-pipelined corked case has the same pre-existing gap; this makes the pipelined case strictly better.
…kips end_stream() on the replayed request
| if (pair.second) { | ||
| uwsRes->replayPipelinedRequests(); | ||
| } |
There was a problem hiding this comment.
🔴 🔴 Same-class site the 7943912/860c0c21 fixes did not reach: when an async handler resolves to new Response(directStream) and pull(c){c.write(...);c.end()} drains synchronously, assign_to_stream → send_readable → res.try_end() → replayPipelinedRequests() dispatches B (this is a promise-resolution microtask, so isParsingHttp is false and pipelinedBuffer is populated). Back at RequestContext.rs:2097 resp.has_responded() now reads B's HTTP_RESPONSE_PENDING and returns false, so the 'done' branch is skipped and do_render_stream falls through to this.render_missing() (line 2284) → ctx.end(b"") → detach_response() (A's has_abort_handler from to_async → resp.clear_aborted() nulls B's onAborted) → resp.end(b"") writes a bogus HTTP/1.1 200 OK…Content-Length: 0 on B's cleared state and markDone()s B. The Pending arm (line 2155) has the same shape: resp.on_writable(on_writable_response_stream, this) overwrites B's onWritable with A's stale ctx. The 7943912 comment dismissed HTTPServerWritable's try_end as "the sync-resolved-stream path inside the initial dispatch, where pipelinedBuffer is empty" — that only holds when the handler is sync; the same send_readable runs from do_render_stream in a microtask.
Extended reasoning...
What the bug is
7943912/860c0c21 fixed the "Rust callers touch resp after the wrapper replayed" class at every site where a Rust frame calls resp.clear_*() after a uws_res_end* wrapper returns. do_render_stream (RequestContext.rs:1998) has the same shape but with a different guard: instead of a post-end resp.clear_*(), it reads resp.has_responded() (line 2097) after assign_to_stream() returns. Before this PR that predicate meant "did A respond?"; now that uws_res_try_end/uws_res_end synchronously replay the next pipelined request, it means "is the connection idle?", and B's resetResponseState() re-sets HTTP_RESPONSE_PENDING — so the answer flips to false and do_render_stream proceeds to write on B's state as if it were still A's.
The 7943912 timeline comment dismissed HTTPServerWritable's try_end/end (streams.rs:1296/1317/1781, reached from end_from_js) as "the sync-resolved-stream path inside the initial dispatch, where pipelinedBuffer is empty and replay is a no-op". That holds only when the handler is sync. When the handler is async and resolves to new Response(directStream), do_render_stream runs from on_response (a promise-resolution microtask, isParsingHttp == false) with pipelinedBuffer populated, and the same send_readable → try_end path replays.
Step-by-step proof — the render_missing() fall-through
- Client sends, in one TCP write:
GET /a HTTP/1.1\r\nHost: x\r\n\r\nGET /b HTTP/1.1\r\nHost: x\r\n\r\nagainstBun.serve({ async fetch() { await 1; // go async → deferPipeline=true, B lands in pipelinedBuffer return new Response(new ReadableStream({ type: 'direct', pull(c) { c.write('A'); c.end(); } })); }});
- A's promise resolves from a microtask →
render()→ RequestContext.rs:3200resp.run_corked_with_type(do_render_stream, ...). Inside,render_metadata()writes A's status/headers (setshas_written_status), thenassign_to_stream()→readDirectStream(BunStreamSource.cpp:936) synchronously callspull(controller). c.write('A')buffers intosink.buffer;c.end()→HTTPServerWritable::end_from_js(streams.rs:1751):requested_end=true,readable_len>0→send_readable(0)→ line 1293 (requested_end && !HTTP_WRITE_CALLED) →res.try_end(buffer, end_len, false)at streams.rs:1296 →uws_res_try_end.- In
uws_res_try_end:internalEnd()→markDone()clears A'sHTTP_RESPONSE_PENDING; corked, so the else-branchuncork()releases the slot.pair.second == true→replayPipelinedRequests().isNoLongerHttp()false,pipelinedBuffernon-empty,isParsingHttpfalse (microtask, notonData),shouldCloseConnection()false →onData<false>(B). B'srequestHandlerrunsresetResponseState()(clearsHTTP_STATUS_CALLED | HTTP_WRITE_CALLED | HTTP_END_CALLED, setsHTTP_RESPONSE_PENDING); B's async handler installsonAbortedviato_async. - Control returns through
end_from_js:mark_done()setssink.done=true;signal.close(None)→readDirectStreamCloseImpl→stream->m_state = Closed;finalize()'s!doneblock is skipped.readDirectStreamline 944 (m_state == Readable) is false → returnsjsUndefined(). - Back in
do_render_stream, line 2097resp.has_responded()=!(state & HTTP_RESPONSE_PENDING)= false (B set it) → the 'done' branch that would have safely calledthis.end_stream()(which seesis_response_pending()==trueand… also mis-fires — see below — but pre-PR would have seenfalseand skipped the FFI end) is skipped. effective_resultstays undefined (sink.pending_flushis None). Line 2126 skipped. Line 2224is_aborted_or_ended()= false (this.respstillSome, not aborted).is_in_progress= true (sink.wrote > 0fromhandle_wrote). Falls through to line 2284this.render_missing().render_missing_corked(line 1004, sincehas_written_statuswas set at step 2) →ctx.end(b"", …)→ line 1137self.detach_response()— the old one. A'sflags.has_abort_handler()is true (set byto_async) →resp.clear_aborted()writesnullptrintohttpResponseData->onAborted, nulling B's freshly-installed handler.- Line 1139
resp.end(b"", …)→uws_res_end. B'sresetResponseState()clearedHTTP_STATUS_CALLED/HTTP_END_CALLED, sowriteStatus(HTTP_200_OK)writesHTTP/1.1 200 OK\r\n,internalEndwritesContent-Length: 0\r\n\r\n,markDone()clears B'sHTTP_RESPONSE_PENDING, andreplayPipelinedRequests()runs again. B's real response, when it eventually resolves, findsthis.resp.has_responded() == trueand is silently dropped (or spliced after the bogus bytes if a third request was buffered).
The Pending-promise arm (line 2155)
When readDirectStream returns the closePromise (stream still Readable after pull() returns — e.g. pull returns a promise), or on the readStreamIntoSink non-direct path which always returns a promise, and A's try_end/end inside assign_to_stream replayed B: promise.unwrap → Pending → line 2155 resp.on_writable(Self::on_writable_response_stream, this) writes A's callback + A's RequestContext* into httpResponseData->{onWritable, writableUserData}, overwriting whatever B installed. On the next writable event, A's on_writable_response_stream runs against A's already-torn-down sink, or B's backpressured body stalls because its onWritable was overwritten.
Why existing code doesn't prevent it
replayPipelinedRequests()'sisParsingHttpguard doesn't help: this runs from a microtask, notonData.callOnWritable's placeholder-restore fix and theisNoLongerHttp()guard inonWritabledon't help: this is not the writable path.- The open comment at streams.rs:1416 (860c0c2's
on_writable→finalize()) is a different entry point (callOnWritable); this one isdo_render_stream's post-assign_to_streamstate reads. end_from_js's ownmark_done()/finalize()correctly skip the!doneblock, so the sink itself is fine — it's the caller (do_render_stream) that then reads the wrong state.
Impact
Wire corruption on any async handler returning a synchronously-drained direct ReadableStream with a pipelined request behind it: B's client receives a bogus empty 200 (or A's stale onWritable stomps B's), B's real response is dropped, and B's onAborted/req.signal never fire on disconnect. If B's replay parse-errors or WS-upgrades (destructing HttpResponseData), the subsequent resp.has_responded() / resp.on_writable() / resp.end() read/write a destructed ext block.
How to fix
do_render_stream needs to detect "A's response completed inside assign_to_stream" without reading the shared per-socket state after that call. Options:
- Snapshot
resp.state()(orhas_responded()) beforeassign_to_stream()and compare, or checkresponse_stream.sink.done/sink.ended_response(set byend_from_js'smark_done()) instead ofresp.has_responded()at line 2097. When true, take the 'done' branch but calldetach_response_after_end()+end_request_streaming_and_drain()+deref()directly instead ofthis.end_stream()(whoseresp.state().is_response_pending()guard at line 1157 is now also inverted). - Gate line 2155's
resp.on_writable(...)(and any other post-assign_to_streamresp.write) on the same local "A completed" signal.
Per REVIEW.md ("Fix the whole class in the same PR — same-class sites are ONE concern"), this belongs with the other post-replay-resp-access fixes already in this PR.
There was a problem hiding this comment.
Valid, and the ninth same-class site across eight review rounds. The pattern is clear: synchronous replay from a uws_res_end* wrapper means every Rust frame that reads/writes resp after a call that can reach a wrapper is a hazard, and do_render_stream is one more (it reads resp.has_responded() and writes resp.on_writable(...) after assign_to_stream() which can reach try_end).
The local fix is what this comment suggests: check response_stream.sink.ended_response (set by end_from_js's mark_done()) instead of resp.has_responded() at RequestContext.rs:2097, take the done branch with detach_response_after_end() + end_request_streaming_and_drain() + deref() directly, and gate line 2155's resp.on_writable(...) on the same sink flag.
But each round has found a deeper frame, and do_render_stream is large enough that there may be more. The structural alternative is to make replay asynchronous (defer the onData call to the next loop tick via a per-HttpContext pending-replay socket list that onClose prunes), so no caller frame can observe the socket after replay. That's a larger change I have not attempted.
Leaving this open for a maintainer to decide between the local fix and the structural one; happy to implement either.
…d it HttpContext::onData pauses the socket while request bytes are parked behind a pending response. If that response turns out to be a WebSocket upgrade, upgrade() destructs the HttpResponseData (and the parked bytes with it, as for bytes trailing a synchronous upgrade) but us_socket_adopt carries the paused flag over, so the WebSocket sent its 101 and then never read a frame. Drop the parked bytes and resume before internalEnd(), so markDone() does not arm a replay dispatch for them either. Tests folded in from the earlier pipelining PRs (#33664, #35036, #32868) for the cases bun-serve-pipelining.test.ts did not cover yet: a Bun.file() body from the handler and a Bun.file() route over each transport, a held request that carries a body with a Connection: close request behind it, held bytes that are a parse error, and the upgrade case above.
|
Closing in favour of #38128, which fixes the same close (the The four cases from this PR's |
…d it HttpContext::onData pauses the socket while request bytes are parked behind a pending response. If that response turns out to be a WebSocket upgrade, upgrade() destructs the HttpResponseData (and the parked bytes with it, as for bytes trailing a synchronous upgrade) but us_socket_adopt carries the paused flag over, so the WebSocket sent its 101 and then never read a frame. Drop the parked bytes and resume before internalEnd(), so markDone() does not arm a replay dispatch for them either. Tests folded in from the earlier pipelining PRs (#33664, #35036, #32868) for the cases bun-serve-pipelining.test.ts did not cover yet: a Bun.file() body from the handler and a Bun.file() route over each transport, a held request that carries a body with a Connection: close request behind it, held bytes that are a parse error, and the upgrade case above.
…d it HttpContext::onData pauses the socket while request bytes are parked behind a pending response. If that response turns out to be a WebSocket upgrade, upgrade() destructs the HttpResponseData (and the parked bytes with it, as for bytes trailing a synchronous upgrade) but us_socket_adopt carries the paused flag over, so the WebSocket sent its 101 and then never read a frame. Drop the parked bytes and resume before internalEnd(), so markDone() does not arm a replay dispatch for them either. Tests folded in from the earlier pipelining PRs (#33664, #35036, #32868) for the cases bun-serve-pipelining.test.ts did not cover yet: a Bun.file() body from the handler and a Bun.file() route over each transport, a held request that carries a body with a Connection: close request behind it, held bytes that are a parse error, and the upgrade case above.
…d it HttpContext::onData pauses the socket while request bytes are parked behind a pending response. If that response turns out to be a WebSocket upgrade, upgrade() destructs the HttpResponseData (and the parked bytes with it, as for bytes trailing a synchronous upgrade) but us_socket_adopt carries the paused flag over, so the WebSocket sent its 101 and then never read a frame. Drop the parked bytes and resume before internalEnd(), so markDone() does not arm a replay dispatch for them either. Tests folded in from the earlier pipelining PRs (#33664, #35036, #32868) for the cases bun-serve-pipelining.test.ts did not cover yet: a Bun.file() body from the handler and a Bun.file() route over each transport, a held request that carries a body with a Connection: close request behind it, held bytes that are a parse error, and the upgrade case above.
When a pipelined HTTP/1.1 request arrives while the previous response on the connection is still pending (async handler), uWS was closing the socket the moment it saw the second request head with
HTTP_RESPONSE_PENDINGset. The in-flight response was discarded (onAbortedfires,respnulled) and the pipelined request was never dispatched: both requests produced zero bytes on the wire.Reproduction
The canonical trigger is a relay handler forwarding
req.bodyinto an outgoingfetch(), but any async handler exhibits it (e.g. two pipelined GETs with anawait Bun.sleep()in the handler).Cause
fenceAndConsumePostPaddedparses every complete request in a TCP read in one loop. After dispatching A and delivering its body, it loops to B's head.requestHandlerfor B seesHTTP_RESPONSE_PENDINGstill set (A's handler is async) and, for!IsNodeHttp, callsus_socket_close()(HttpContext.h:376). The comment there says "denying async pipelining until, if ever, we want to support it"; this PR supports it.Fix
While a Bun.serve response is pending, the parse loop breaks before
getHeadersmutates the next request's bytes (deferPipelineinHttpParser). The unconsumed bytes go to a per-connectionpipelinedBufferand reads are paused.markDone()clears the flag and replays the buffer throughonData, so the next request is dispatched once the in-flight response is written.deferPipelineis re-derived fromHTTP_RESPONSE_PENDINGafter the body fin callback, so a handler that responds synchronously from inside the body data callback (await req.text()then return) keeps the existing sync-pipelining fast path and never buffers.Replay is skipped when
isParsingHttpis already set (the pathological case of two connections whose handlers synchronously resolve each other); the buffer stays and the connection idles out rather than re-enteringonDataand stomping per-context parse state.shouldCloseConnection()(request or response asked forConnection: close) discards the buffer and proceeds with the existing close-after-drain path.node:httpdispatches pipelined requests immediately and queues responses on the JS side; every new path here is gated on!IsNodeHttpand that mechanism is unchanged.Verification
Three new tests in
test/js/bun/http/serve.test.tsunder "dispatches a pipelined request after the previous async response completes": thereq.bodyforwarded tofetch()case, three async GETs in one write, and a pipelined request sent in a separate TCP write while the first is pending. All three fail onmain(zero responses) and pass with this change.Related
Supersedes #35036 (which delivers A with
Connection: closeand drops B) and #32868. Likely addresses #6961, #22174, #26406, #11228 (pipelined request arriving while aBun.file()response is mid-sendfile).no test proof · iteration 5 · Platform-specific test(s) that do not run on this machine. Deferring to CI, which covers all platforms: test/js/bun/http/serve.test.ts