Skip to content
Closed
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
37 changes: 37 additions & 0 deletions packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -472,6 +472,14 @@ struct HttpContext {
((HttpResponse<SSL> *) s)->resetTimeout();
}

/* Bun.serve: an async handler left the response pending, so any
* pipelined request in this TCP read must be buffered and replayed
* once this response completes (markDone). node:http dispatches
* pipelined requests immediately and queues their responses above. */
if constexpr (!IsNodeHttp) {
httpResponseData->deferPipeline = !((HttpResponse<SSL> *) s)->hasResponded();
}

/* Continue parsing */
return s;

Expand Down Expand Up @@ -543,6 +551,16 @@ struct HttpContext {
httpResponseData->inStream = nullptr;
}
}

/* Bun.serve: the handler may have responded synchronously inside
* the body fin callback (e.g. `await req.text()` drains microtasks
* and renders), so re-derive from HTTP_RESPONSE_PENDING here. */
if constexpr (!IsNodeHttp) {
if (fin) {
httpResponseData->deferPipeline =
(httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) != 0;
}
}
return user;
});

Expand Down Expand Up @@ -599,6 +617,16 @@ struct HttpContext {
}
}

/* Bun.serve async pipelining: pipelined request bytes were buffered
* because the current response is still pending. Pause reads so the
* buffer stays bounded by this single recv; replayPipelinedRequests()
* resumes and replays once the response completes. */
if constexpr (!IsNodeHttp) {
if (!httpResponseData->pipelinedBuffer.empty()) {
((HttpResponse<SSL> *) s)->pause();
}
}

/* Timeout on uncork failure */
auto [written, failed] = ((AsyncSocket<SSL> *) returnedData)->uncork();
if (written > 0 || failed) {
Expand Down Expand Up @@ -700,6 +728,15 @@ struct HttpContext {
* If write was never called, the developer should still return true so that we may drain. */
bool success = httpResponseData->callOnWritable(reinterpret_cast<HttpResponse<SSL> *>(asyncSocket), httpResponseData->offset);

/* The writable callback may have completed the response and replayed
* a buffered pipelined request whose dispatch closed or adopted this
* socket (parse error, Connection: close, WebSocket upgrade); every
* httpResponseData read below would then be on a destructed object.
* An in-place adopt leaves is_closed false, so also check kind. */
if (reinterpret_cast<HttpResponse<SSL> *>(s)->isNoLongerHttp()) {
return s;
}

if constexpr (!IsNodeHttp) {
/* Bun.serve: onEnd deferred close for a tryEnd tail (offset < total,
* nothing in AsyncSocketData::buffer). A retry that moves zero bytes
Expand Down
29 changes: 28 additions & 1 deletion packages/bun-uws/src/HttpParser.h
Original file line number Diff line number Diff line change
Expand Up @@ -589,6 +589,18 @@ struct HttpResponseData;
std::string fallback;
/* This guy really has only 30 bits since we reserve two highest bits to chunked encoding parsing state */
uint64_t remainingStreamingBytes = 0;

public:
/* Bun.serve async pipelining: while a response on this connection is
* still in flight (HTTP_RESPONSE_PENDING), the parse loop must stop
* BEFORE getHeaders mutates the next request's bytes; those bytes are
* held in pipelinedBuffer and replayed via replayPipelinedRequests()
* (from the uws_res_end* wrappers) once the
* in-flight response completes. node:http uses its own queue and
* never sets this. */
bool deferPipeline = false;
std::string pipelinedBuffer;
Comment thread
robobun marked this conversation as resolved.
private:
/* node:http compat: a completed request on this connection forbade keep-alive
* (Connection: close, or HTTP/1.0), so no further message may be dispatched
* (llhttp parses nothing after such a message: HPE_CLOSED_CONNECTION). */
Expand Down Expand Up @@ -1102,6 +1114,12 @@ struct HttpResponseData;
data[length + 1] = 'a'; /* Anything that is not \n, to trigger "invalid request" */
req->ancientHttp = false;
for (;length;) {
/* A response on this connection is still in flight (async handler).
* Stop before getHeaders mutates the next request's bytes so the
* caller can buffer them verbatim for replay. */
if (deferPipeline) {
break;
}
/* node:http server compat: an accepted Upgrade request whose body just
* finished parsing switched this connection into tunnel mode (the data
* handler set isConnectRequest when it saw the body fin). Everything
Expand Down Expand Up @@ -1536,7 +1554,16 @@ struct HttpResponseData;
length -= consumedBytes;

if (length) {
if (length < maxFallbackSize) {
if (deferPipeline) {
/* Pipelined request bytes held until the in-flight response
* completes; replayed via replayPipelinedRequests(). Reads are paused by the
* caller while this buffer is non-empty, so it is bounded by a
* single recv buffer's worth. Reserve post-padding so replay
* can hand this buffer straight back to getHeaders. */
pipelinedBuffer.reserve(pipelinedBuffer.length() + length
+ std::max<unsigned int>(MINIMUM_HTTP_POST_PADDING, sizeof(std::string)));
pipelinedBuffer.append(data, length);
} else if (length < maxFallbackSize) {
fallback.append(data, length);
} else {
return HttpParserResult::error(HTTP_ERROR_431_REQUEST_HEADER_FIELDS_TOO_LARGE, HTTP_PARSER_ERROR_REQUEST_HEADER_FIELDS_TOO_LARGE);
Expand Down
65 changes: 65 additions & 0 deletions packages/bun-uws/src/HttpResponse.h
Original file line number Diff line number Diff line change
Expand Up @@ -455,6 +455,71 @@ struct HttpResponse : public AsyncSocket<SSL> {
return this;
}

/* True once this socket is no longer a valid HttpResponse: closed, shut
* down, or adopted (in-place or relocated) into a WebSocket. An in-place
* adopt leaves is_closed false, so also check the kind byte. */
bool isNoLongerHttp() {
us_socket_t *s = (us_socket_t *) this;
return us_socket_is_closed(s) || us_socket_is_shut_down(s)
|| us_socket_kind(s) != HttpContext<SSL>::socketKind();
}

/* Dispatch pipelined request bytes that were buffered while the previous
* response was still in flight (Bun.serve's async-pipelining path). Called
* as the final action of each uws_res_end* C ABI wrapper: the replayed
* onData can synchronously close or adopt this socket (parse error,
* Connection: close, WebSocket upgrade), destructing HttpResponseData, so
* nothing may touch the response after this returns. node:http never
* buffers here (it dispatches immediately and queues responses), so replay
* is always the IsNodeHttp=false onData. */
void replayPipelinedRequests() {
/* The caller's own close gate (internalEnd's shouldCloseConnection
* branch) may have already destructed HttpResponseData before we run;
* isNoLongerHttp reads only us_socket_t flags, so check it first. */
if (isNoLongerHttp()) {
return;
}
HttpResponseData<SSL> *httpResponseData = getHttpResponseData();
if (httpResponseData->pipelinedBuffer.empty()) {
return;
}
/* Re-entering onData from inside onData would stomp the per-context
* isParsingHttp/upgradedWebSocket state; the pathological case is two
* connections whose handlers synchronously resolve each other. Leave
* the buffer and reads paused; the connection idles out rather than
* corrupting the outer parse. */
HttpContextData<SSL> *httpContextData = HttpContext<SSL>::getSocketContextDataS((us_socket_t *) this);
if (httpContextData->flags.isParsingHttp) {
return;
}
/* Flush the just-completed response before dispatching the next
* request: onData's parse-error path closes with uncorkWithoutSending,
* which would otherwise drop the corked bytes. */
Super::uncork();
/* Connection is being torn down after this response; the buffered
* request is discarded. On the corked async path internalEnd()'s own
* close gate was skipped (isCorked() was true, then its uncork()
* released the slot so HttpResponse::cork() early-returns), so close
* here rather than leave the socket open until client FIN/timeout. */
if (httpResponseData->shouldCloseConnection()) {
httpResponseData->pipelinedBuffer.clear();
if (((AsyncSocket<SSL> *) this)->hasFullyDrained()) {
((AsyncSocket<SSL> *) this)->shutdown();
((AsyncSocket<SSL> *) this)->close();
}
return;
}
Comment thread
claude[bot] marked this conversation as resolved.

std::string buffer = std::move(httpResponseData->pipelinedBuffer);
httpResponseData->pipelinedBuffer.clear();
/* Reads were paused when the buffer became non-empty. */
this->resume();
/* getHeaders writes post-padding bytes; the buffer reserved them at
* append time (see consumePostPadded). */
buffer.reserve(buffer.length() + MINIMUM_HTTP_POST_PADDING);
HttpContext<SSL>::template onData<false>((us_socket_t *) this, buffer.data(), (int) buffer.length());
}

/* Note: Headers are not checked in regards to timeout.
* We only check when you actively push data or end the request */

Expand Down
36 changes: 29 additions & 7 deletions packages/bun-uws/src/HttpResponseData.h
Original file line number Diff line number Diff line change
Expand Up @@ -57,26 +57,48 @@ struct HttpResponseData : AsyncSocketData<SSL>, HttpParser {

/* We are done with this request */
this->state &= ~HttpResponseData<SSL>::HTTP_RESPONSE_PENDING;
/* The parse loop may dispatch the next request on this connection. */
this->deferPipeline = false;

HttpResponseData<SSL> *httpResponseData = uwsRes->getHttpResponseData();
httpResponseData->isIdle = true;

/* Pipelined request bytes buffered while this response was in flight are
* NOT replayed here: every caller (internalEnd, uws_res_end_sendfile,
* uws_res_end_without_body) still reads our state afterwards, and a
* replay can close/adopt the socket and destruct this object. Callers
* invoke replayPipelinedRequests() themselves as their final action. */
}

/* Caller of onWritable. It is possible onWritable calls markDone so we need to borrow it. */
bool callOnWritable(uWS::HttpResponse<SSL>* response, uint64_t offset) {
/* Borrow real onWritable */
auto* borrowedOnWritable = std::move(onWritable);
auto* borrowedOnWritable = onWritable;

/* Set onWritable to placeholder */
onWritable = [](uWS::HttpResponse<SSL>*, uint64_t, void*) {return true;};
/* Set onWritable to a placeholder we can identify afterwards: the
* borrowed callback may reach markDone() and replay the next pipelined
* request, whose handler can install its OWN onWritable. Restoring by
* non-null would stomp that with this response's stale callback. */
static constexpr OnWritableCallback callOnWritablePlaceholder =
[](uWS::HttpResponse<SSL>*, uint64_t, void*) { return true; };
onWritable = callOnWritablePlaceholder;

/* Run borrowed onWritable */
bool ret = borrowedOnWritable(response, offset, writableUserData);

/* If we still have onWritable (the placeholder) then move back the real one */
if (onWritable) {
/* We haven't reset onWritable, so give it back */
onWritable = std::move(borrowedOnWritable);
/* The callback may have completed the response and replayed a buffered
* pipelined request whose dispatch closed or adopted this socket,
* destructing us; every field access below would then be on a dead
* object (or on WebSocketData after an in-place adopt). */
if (response->isNoLongerHttp()) {
return ret;
}

/* Only restore if onWritable is still OUR placeholder; anything else
* (null from markDone, or a new callback from a replayed request)
* belongs to whoever set it. */
if (onWritable == callOnWritablePlaceholder) {
onWritable = borrowedOnWritable;
}

return ret;
Expand Down
9 changes: 6 additions & 3 deletions src/runtime/server/FileRoute.rs
Original file line number Diff line number Diff line change
Expand Up @@ -604,9 +604,12 @@ impl FileRoute {
}

fn on_response_complete(this: *mut FileRoute, resp: AnyResponse) {
resp.clear_aborted();
resp.clear_on_writable();
resp.clear_timeout();
// See StaticRoute::on_response_complete.
if let AnyResponse::H3(_) = resp {
resp.clear_aborted();
resp.clear_on_writable();
resp.clear_timeout();
}
// SAFETY: `this` is live (ref held by caller); `deref()` may free it.
unsafe {
if let Some(mut server) = (*this).server.get() {
Expand Down
28 changes: 23 additions & 5 deletions src/runtime/server/RequestContext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1096,7 +1096,7 @@ where
};
if try_end_ok {
drop(bb);
self.detach_response();
self.detach_response_after_end();
self.end_request_streaming_and_drain();
self.finalize_without_deinit();
self.deref();
Expand Down Expand Up @@ -1556,7 +1556,7 @@ where
fn on_file_stream_complete(ctx: *mut c_void, _resp: uws::AnyResponse) {
// SAFETY: ctx is a *RequestContext registered with FileResponseStream
let this: &mut Self = unsafe { bun_ptr::callback_ctx::<Self>(ctx) };
this.detach_response();
this.detach_response_after_end();
this.end_request_streaming_and_drain();
this.deref();
}
Expand Down Expand Up @@ -1638,7 +1638,7 @@ where
let bytes = &bytes_[bytes_.len().min(write_offset)..];
// SAFETY: FFI handle
if resp.try_end(bytes, bytes_.len(), self.should_close_connection()) {
self.detach_response();
self.detach_response_after_end();
self.end_request_streaming_and_drain();
self.deref();
true
Expand Down Expand Up @@ -1673,7 +1673,7 @@ where
let done = resp.try_end(bytes, total_len, close_connection);
if done {
self.response_buf_owned.clear();
self.detach_response();
self.detach_response_after_end();
self.end_request_streaming_and_drain();
self.deref();
} else {
Expand Down Expand Up @@ -2389,6 +2389,24 @@ where
}
}

/// `detach_response` without the HTTP/1 `resp.clear_*` FFI calls: H1's
/// `markDone()` already nulled them, and the `uws_res_end*` wrapper then
/// replayed any buffered pipelined request, so those handlers now belong
/// to the next request. HTTP/3's `markDone()` deliberately leaves
/// `onAborted` armed for `on_stream_close`, so fall through to the full
/// detach there.
Comment thread
robobun marked this conversation as resolved.
fn detach_response_after_end(&mut self) {
if HTTP3 {
self.detach_response();
return;
}
self.request_body_buf = Vec::new();
self.resp.take();
self.flags.set_is_waiting_for_request_body(false);
self.flags.set_has_abort_handler(false);
self.flags.set_has_timeout_handler(false);
}

pub fn is_aborted_or_ended(&self) -> bool {
// resp == null or aborted or server.stop(true)
self.resp.is_none()
Expand Down Expand Up @@ -3807,7 +3825,7 @@ where
return;
}
}
self.detach_response();
self.detach_response_after_end();
self.end_request_streaming_and_drain();
self.deref();
}
Expand Down
10 changes: 8 additions & 2 deletions src/runtime/server/StaticRoute.rs
Original file line number Diff line number Diff line change
Expand Up @@ -402,11 +402,17 @@ impl StaticRoute {
/// `this` must be a live heap-allocated route with write provenance; may free
/// `*this` via `deref_` when the refcount reaches zero.
unsafe fn on_response_complete(this: *mut Self, resp: AnyResponse) {
// SAFETY: caller contract.
unsafe {
// HTTP/1: markDone() already nulled these and the wrapper then replayed
// any buffered pipelined request, so clearing would null THAT request's
// handlers. HTTP/3: Http3Response::markDone() leaves onAborted armed
// for on_stream_close, so clear it to avoid a second deref_ here.
Comment thread
robobun marked this conversation as resolved.
if let AnyResponse::H3(_) = resp {
resp.clear_aborted();
resp.clear_on_writable();
resp.clear_timeout();
}
// SAFETY: caller contract.
unsafe {
if let Some(mut server) = (*this).server.get() {
server.on_static_request_complete();
}
Expand Down
16 changes: 12 additions & 4 deletions src/runtime/webcore/streams.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1402,11 +1402,19 @@ impl<const SSL: bool, const HTTP3: bool> HTTPServerWritable<SSL, HTTP3> {
total_written = chunk_len as u64;

if self.requested_end {
if let Some(res) = self.any_res() {
res.clear_on_writable();
// HTTP/1: `send_readable` drained the parked `try_end`/`end`,
// so `markDone()` nulled our `onWritable` and the wrapper then
// replayed any buffered pipelined request; the socket now
// belongs to THAT request. HTTP/3 has no replay.
Comment thread
robobun marked this conversation as resolved.
if HTTP3 {
if let Some(res) = self.any_res() {
res.clear_on_writable();
}
} else {
// finalize()'s `if !self.done` block would clear_on_writable
// and end_stream() the replayed request's state.
Comment thread
robobun marked this conversation as resolved.
self.done = true;
}
// `send_readable` drained the parked `try_end`, so uWS has
// `markDone()`d the response and dropped its `onAborted`.
self.ended_response = true;
self.signal.close(None);
let _ = self.flush_promise(); // TODO: properly propagate exception upwards
Comment thread
robobun marked this conversation as resolved.
Expand Down
Loading
Loading