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
25 changes: 15 additions & 10 deletions src/js/node/http2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2060,6 +2060,8 @@ enum StreamState {
// callback). Until then no 'error' listener can exist, so stream errors must not be emitted:
// node never constructs the JS stream object before a complete header block arrives.
Delivered = 1 << 8, // 100000000 = 256
// Native dispatched streamError or aborted: it closed the stream and wrote the RST_STREAM if one was due.
NativeReset = 1 << 9, // 1000000000 = 512
}
// native.writeStream() return-value flag (mirrors WRITE_FLUSHED_WITHOUT_CALLBACK in
// h2_frame_parser.rs): the chunk was handed to the socket without queueing and the engine did
Expand Down Expand Up @@ -2226,8 +2228,9 @@ function markStreamClosed(stream: Http2Stream) {
markWritableDone(stream);
}
}
function rstNextTick(id: number, rstCode: number) {
const session = this as Http2Session;
function rstNextTick(this: Http2Stream, session: Http2Session, id: number, rstCode: number) {
// Native drops this call only via the stream's table entry: evicted on the next read, absent for a pushed stream.
if ((this[bunHTTP2StreamStatus] & (StreamState.NativeClosed | StreamState.NativeReset)) !== 0) return;
session[bunHTTP2Native]?.rstStream(id, rstCode);
}
// node streamOnPause/streamOnResume (lib/internal/http2/core.js): the readable's flow state
Expand All @@ -2247,7 +2250,7 @@ function streamOnResume(this: Http2Stream) {
// A close() on a stream that has not been submitted yet (no id): the RST_STREAM has to follow the
// HEADERS frame, which is sent when the queued request becomes ready (node's finishCloseStream).
function sendRstOnReady(this: Http2Stream, session: Http2Session, code: number) {
setImmediate(rstNextTick.bind(session, this.id, code));
setImmediate(rstNextTick.bind(this, session, this.id, code));
}
function uncorkNT(stream: Http2Stream) {
stream.uncork();
Expand Down Expand Up @@ -2591,9 +2594,9 @@ class Http2Stream extends (Duplex as Http2StreamBase) {
// RST_STREAM has to be sent after the HEADERS frame, once the id is assigned.
this.once("ready", sendRstOnReady.bind(this, session, code));
} else if (this.writableFinished || code) {
setImmediate(rstNextTick.bind(session, this.#id, code));
setImmediate(rstNextTick.bind(this, session, this.#id, code));
} else {
this.once("finish", rstNextTick.bind(session, this.#id, code));
this.once("finish", rstNextTick.bind(this, session, this.#id, code));
}
// node destroys the stream once both halves have finished; without this a stream closed
// while idle never emits 'close'.
Expand Down Expand Up @@ -2680,11 +2683,10 @@ class Http2Stream extends (Duplex as Http2StreamBase) {
session &&
typeof this.#id === "number" &&
!this[kNeverAnnounced] &&
// A cleanly closed stream the native side already freed has nothing to send:
// the deferred rstStream would be a guaranteed no-op host call per request.
(rstCode !== 0 || (this[bunHTTP2StreamStatus] & StreamState.NativeClosed) === 0)
// Native already closed or reset the stream: nothing is left to send, whatever the rstCode.
(this[bunHTTP2StreamStatus] & (StreamState.NativeClosed | StreamState.NativeReset)) === 0
) {
setImmediate(rstNextTick.bind(session, this.#id, rstCode));
setImmediate(rstNextTick.bind(this, session, this.#id, rstCode));
}

// Diagnostics channels: published after the stream is closed and destroyed, with the same error
Expand Down Expand Up @@ -4071,6 +4073,7 @@ class ServerHttp2Session extends Http2Session {
},
aborted(self: ServerHttp2Session, stream: ServerHttp2Stream, error: any, old_state: number) {
if (!self || typeof stream !== "object") return;
stream[bunHTTP2StreamStatus] |= StreamState.NativeReset;
stream.rstCode = constants.NGHTTP2_CANCEL;
// if writable and not closed emit aborted
if (old_state != 5 && old_state != 7) {
Expand All @@ -4083,6 +4086,7 @@ class ServerHttp2Session extends Http2Session {
},
streamError(self: ServerHttp2Session, stream: ServerHttp2Stream, error: number) {
if (!self || typeof stream !== "object") return;
stream[bunHTTP2StreamStatus] |= StreamState.NativeReset;
self.#connections--;
if (stream.id % 2 === 1) self.#peerInitiatedStreams--;
process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed);
Expand Down Expand Up @@ -5077,6 +5081,7 @@ class ClientHttp2Session extends Http2Session {
),
aborted: withStreamFrame((self: ClientHttp2Session, stream: ClientHttp2Stream, error: any, old_state: number) => {
if (!self || typeof stream !== "object") return;
stream[bunHTTP2StreamStatus] |= StreamState.NativeReset;
stream.rstCode = constants.NGHTTP2_CANCEL;
// if writable and not closed emit aborted
if (old_state != 5 && old_state != 7) {
Expand All @@ -5088,7 +5093,7 @@ class ClientHttp2Session extends Http2Session {
}),
streamError: withStreamFrame((self: ClientHttp2Session, stream: ClientHttp2Stream, error: number) => {
if (!self || typeof stream !== "object") return;

stream[bunHTTP2StreamStatus] |= StreamState.NativeReset;
self.#connections--;
process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed);
}),
Expand Down
23 changes: 22 additions & 1 deletion src/runtime/api/bun/h2/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1377,6 +1377,10 @@ impl Connection {
discard = true;
}
Some(st) => {
if st.state == State::ReservedRemote {
self.data_on_reserved_stream(sink, hdr.stream_id);
return StreamedDataStart::Fatal;
}
if !stream::can_receive_data(st.state) {
self.send_rst_stream(sink, hdr.stream_id, ErrorCode::StreamClosed);
if let Some(st2) = self.streams.get_mut(&hdr.stream_id) {
Expand Down Expand Up @@ -1494,6 +1498,7 @@ impl Connection {
// below don't alias the streams map.
enum DataDecision {
Rst(ErrorCode),
ReservedStream,
FlowControlViolation,
Deliver(u32),
}
Expand All @@ -1513,7 +1518,9 @@ impl Connection {
// §5.1: DATA for an unknown/closed stream is a STREAM_CLOSED error.
None => DataDecision::Rst(ErrorCode::StreamClosed),
Some(s) => {
if !stream::can_receive_data(s.state) {
if s.state == State::ReservedRemote {
DataDecision::ReservedStream
} else if !stream::can_receive_data(s.state) {
DataDecision::Rst(ErrorCode::StreamClosed)
} else {
s.recv_window.on_data(consumed);
Expand All @@ -1536,6 +1543,10 @@ impl Connection {
sink.on_stream_reset(hdr.stream_id, code.as_u32());
return false;
}
DataDecision::ReservedStream => {
self.data_on_reserved_stream(sink, hdr.stream_id);
return true;
Comment thread
robobun marked this conversation as resolved.
}
DataDecision::FlowControlViolation => {
// nghttp2 (nghttp2_session_update_recv_stream_window_size): a stream flow-control
// violation terminates the whole session with FLOW_CONTROL_ERROR; node surfaces
Expand Down Expand Up @@ -1580,6 +1591,16 @@ impl Connection {
false
}

/// RFC 9113 §5.1 reserved (remote): DATA is a connection PROTOCOL_ERROR, as in nghttp2.
fn data_on_reserved_stream(&mut self, sink: &impl Sink, stream_id: u32) {
if let Some(s) = self.streams.get_mut(&stream_id) {
s.state = State::Closed;
}
// Session teardown in the embedder misses a promised stream: report this one first.
sink.on_stream_reset(stream_id, ErrorCode::InternalError.as_u32());
self.send_go_away(sink, ErrorCode::ProtocolError, b"DATA: stream in reserved");
}

/// RFC 9113 §8.1.1: once END_STREAM arrives, a request whose received DATA total contradicts
/// its declared `content-length` is malformed. Resets the stream with PROTOCOL_ERROR instead
/// of signalling end-of-stream and returns true if it did so.
Expand Down
Loading
Loading