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
37 changes: 29 additions & 8 deletions src/js/node/http2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2044,6 +2044,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
// Destroyed by session.destroy(): like node's _destroy, the writable ends without _final.
SessionDestroyed = 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 @@ -2280,8 +2282,15 @@ function destroyStreamForSessionDestroy(error: Error | undefined, rstCode: numbe
// listener would otherwise turn session.destroy(code) into an uncaught
// exception (e.g. grpc-js forceShutdown destroying sessions with
// NGHTTP2_CANCEL while unread UNIMPLEMENTED streams are still around).
stream[bunHTTP2StreamStatus] |= StreamState.SessionDestroyed;
stream.destroy(error !== undefined && stream.listenerCount("error") > 0 ? error : undefined);
}
// Client counterpart, through emitStreamErrorNT for the deferred path's error and rstCode.
function cancelStreamForSessionDestroy(session: ClientHttp2Session, rstCode: number, stream: Http2Stream) {
if (stream.destroyed || stream.closed) return;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
stream[bunHTTP2StreamStatus] |= StreamState.SessionDestroyed;
emitStreamErrorNT(session, stream, rstCode, true, false);
}
class Http2Stream extends Duplex {
#id: number;
[bunHTTP2Session]: ClientHttp2Session | ServerHttp2Session | null = null;
Expand Down Expand Up @@ -2588,11 +2597,16 @@ class Http2Stream extends Duplex {
this[kAborted] = true;
this.emit("aborted");
}
// at this state destroyed will be true but we need to close the writable side
this._writableState.destroyed = false;
this.end();
// we now restore the destroyed flag
this._writableState.destroyed = true;
if ((this[bunHTTP2StreamStatus] & StreamState.SessionDestroyed) !== 0) {
// destroyed stays set, so end() marks the writable ended and _final does not run.
this.end();
} else {
// at this state destroyed will be true but we need to close the writable side
this._writableState.destroyed = false;
this.end();
// we now restore the destroyed flag
this._writableState.destroyed = true;
}
}

const session = this[bunHTTP2Session];
Expand Down Expand Up @@ -4289,7 +4303,6 @@ class ServerHttp2Session extends Http2Session {
// Windows agents the frame deterministically arrived first).
self.destroy();
} else {
self.#parser?.emitErrorToAllStreams(errorCode);
// Like Node, destroy with an error but send our own goaway with
// NGHTTP2_NO_ERROR since this side had no error.
self.destroy(sessionErrorFromCode(errorCode), constants.NGHTTP2_NO_ERROR);
Expand Down Expand Up @@ -4317,7 +4330,8 @@ class ServerHttp2Session extends Http2Session {
#onClose() {
const parser = this.#parser;
if (parser) {
parser.emitAbortToAllStreams();
// Node's socketOnClose: close(NGHTTP2_CANCEL) every stream, then destroy it.
parser.forEachStream(streamCancel);
parser.forEachStream(streamSocketClosed);
parser.detach();
this.#parser = null;
Expand Down Expand Up @@ -5870,9 +5884,16 @@ class ClientHttp2Session extends Http2Session {
}
// Like Node's Http2Stream._destroy: a received GOAWAY's code takes
// precedence over the destroy code when streams are torn down.
const streamRstCode = this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL);
// The native sweep throws on a non-numeric code: the retry must still find the streams.
if (typeof streamRstCode === "number") {
parser.forEachStream(
FunctionPrototypeBind.$call(cancelStreamForSessionDestroy, undefined, this, streamRstCode),
);
}
this[bunHTTP2SessionTeardownFrame] = $getInternalField($asyncContext, 0);
try {
parser.emitErrorToAllStreams(this[kGoawayCode] || (code !== undefined ? code : constants.NGHTTP2_CANCEL));
parser.emitErrorToAllStreams(streamRstCode);
} finally {
this[bunHTTP2SessionTeardownFrame] = kNoSessionTeardown;
}
Expand Down
39 changes: 0 additions & 39 deletions src/runtime/api/bun/h2_frame_parser.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6519,45 +6519,6 @@ impl H2FrameParser {
Ok(JSValue::UNDEFINED)
}

#[bun_jsc::host_fn(method)]
pub(crate) fn emit_abort_to_all_streams(
this: &Self,
_global_object: &JSGlobalObject,
_callframe: &CallFrame,
) -> JsResult<JSValue> {
// R-2: StreamResumableIterator stores a `ParentRef`; `streams` is `JsCell`-backed,
// so the loop body can keep using `this` (`&Self`) directly.
let mut it = StreamResumableIterator::init(this);
while let Some(stream_ptr) = it.next() {
// SAFETY: stream_ptr is a *mut Stream stored in self.streams (heap::alloc); valid for
// the lifetime of the entry. Separate heap allocation from `this`, so no aliasing.
let stream = unsafe { &mut *stream_ptr };
// this is the oposite logic of emitErrorToallStreams, in this case we wanna to cancel this streams
if this.is_server.get() {
if stream.id % 2 == 0 {
continue;
}
} else if stream.id % 2 != 0 {
continue;
}
if stream.state != StreamState::CLOSED {
let old_state = stream.state;
stream.state = StreamState::CLOSED;
stream.rst_code = ErrorCode::CANCEL.0;
let identifier = stream.get_identifier();
identifier.ensure_still_alive();
stream.free_resources::<false>(this);
this.dispatch_with_2_extra(
JSH2FrameParser::Gc::onAborted,
identifier,
JSValue::UNDEFINED,
JSValue::js_number(old_state as u8 as f64),
);
}
}
Ok(JSValue::UNDEFINED)
}

#[bun_jsc::host_fn(method)]
pub(crate) fn emit_error_to_all_streams(
this: &Self,
Expand Down
4 changes: 0 additions & 4 deletions src/runtime/api/h2.classes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -121,10 +121,6 @@ export default [
fn: "emitErrorToAllStreams",
length: 1,
},
emitAbortToAllStreams: {
fn: "emitAbortToAllStreams",
length: 0,
},
getNextStream: {
fn: "getNextStream",
length: 0,
Expand Down
Loading
Loading