From cf5b3cf0366bda7bbae9906b4e5885b4c8193658 Mon Sep 17 00:00:00 2001 From: cirospaciari Date: Tue, 15 Oct 2024 21:19:59 -0700 Subject: [PATCH 01/12] add GC check --- src/bun.js/api/bun/h2_frame_parser.zig | 1 + 1 file changed, 1 insertion(+) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 26c0dd44c264..5b83071b2e7d 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1141,6 +1141,7 @@ pub const H2FrameParser = struct { this.signal = null; signal.deinit(); } + JSC.VirtualMachine.eventLoop().processGCTimer(); } }; From 44d40ad73a605b55fc1f5280afddaa77aa6002d7 Mon Sep 17 00:00:00 2001 From: cirospaciari Date: Tue, 15 Oct 2024 21:31:56 -0700 Subject: [PATCH 02/12] opsie --- src/bun.js/api/bun/h2_frame_parser.zig | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 5b83071b2e7d..202144b1f778 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1141,7 +1141,7 @@ pub const H2FrameParser = struct { this.signal = null; signal.deinit(); } - JSC.VirtualMachine.eventLoop().processGCTimer(); + JSC.VirtualMachine.get().eventLoop().processGCTimer(); } }; From d531ec536b598178400c705d8b04653dc30ace1d Mon Sep 17 00:00:00 2001 From: cirospaciari Date: Tue, 15 Oct 2024 23:08:48 -0700 Subject: [PATCH 03/12] detach soon --- src/js/node/http2.ts | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 4840bf4d834c..a21d399f0912 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2346,6 +2346,7 @@ class ServerHttp2Session extends Http2Session { self.emit("error", error_instance); self[bunHTTP2Socket]?.end(); self.#parser = null; + this[bunHTTP2Socket] = null; }, wantTrailers(self: ServerHttp2Session, stream: ServerHttp2Stream) { if (!self || typeof stream !== "object") return; @@ -2369,11 +2370,13 @@ class ServerHttp2Session extends Http2Session { self[bunHTTP2Socket]?.end(); self.#parser = null; + this[bunHTTP2Socket] = null; }, end(self: ServerHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; self[bunHTTP2Socket]?.end(); self.#parser = null; + this[bunHTTP2Socket] = null; }, write(self: ServerHttp2Session, buffer: Buffer) { if (!self) return -1; @@ -2757,6 +2760,7 @@ class ClientHttp2Session extends Http2Session { self.emit("error", error_instance); self[bunHTTP2Socket]?.destroy(); self.#parser = null; + this[bunHTTP2Socket] = null; }, wantTrailers(self: ClientHttp2Session, stream: ClientHttp2Stream) { @@ -2778,11 +2782,13 @@ class ClientHttp2Session extends Http2Session { } self[bunHTTP2Socket]?.end(); self.#parser = null; + this[bunHTTP2Socket] = null; }, end(self: ClientHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; self[bunHTTP2Socket]?.end(); self.#parser = null; + this[bunHTTP2Socket] = null; }, write(self: ClientHttp2Session, buffer: Buffer) { if (!self) return -1; From 1819d21461f37dbbefdaa2b842a4c7c2f6a88c44 Mon Sep 17 00:00:00 2001 From: cirospaciari Date: Wed, 16 Oct 2024 00:08:17 -0700 Subject: [PATCH 04/12] avoid drain too much --- src/bun.js/api/bun/h2_frame_parser.zig | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 202144b1f778..0c12359dd6aa 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1558,6 +1558,9 @@ pub const H2FrameParser = struct { pub fn flush(this: *H2FrameParser) usize { this.ref(); defer this.deref(); + const eventLoop = JSC.VirtualMachine.get().eventLoop(); + eventLoop.enter(); + defer eventLoop.exit(); var written = switch (this.native_socket) { .tls_writeonly, .tls => |socket| this._genericFlush(*TLSSocket, socket), .tcp_writeonly, .tcp => |socket| this._genericFlush(*TCPSocket, socket), @@ -3636,6 +3639,9 @@ pub const H2FrameParser = struct { buffer.ensureStillAlive(); if (buffer.asArrayBuffer(globalObject)) |array_buffer| { var bytes = array_buffer.byteSlice(); + const eventLoop = globalObject.bunVM().eventLoop(); + eventLoop.enter(); + defer eventLoop.exit(); // read all the bytes while (bytes.len > 0) { const result = this.readBytes(bytes); @@ -3652,6 +3658,9 @@ pub const H2FrameParser = struct { this.ref(); defer this.deref(); var bytes = data; + const eventLoop = JSC.VirtualMachine.get().eventLoop(); + eventLoop.enter(); + defer eventLoop.exit(); while (bytes.len > 0) { const result = this.readBytes(bytes); bytes = bytes[result..]; From ddab03c38967858a3c223efb8d6e69f459b6f807 Mon Sep 17 00:00:00 2001 From: cirospaciari Date: Wed, 16 Oct 2024 00:25:14 -0700 Subject: [PATCH 05/12] no reportExtraMemory --- src/bun.js/api/bun/h2_frame_parser.zig | 29 +++++++++----------------- 1 file changed, 10 insertions(+), 19 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 0c12359dd6aa..26110481bd64 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1051,7 +1051,7 @@ pub const H2FrameParser = struct { }; if (bytes.len > 0) { @memcpy(frame.buffer[0..bytes.len], bytes); - client.globalThis.vm().reportExtraMemory(bytes.len); + // client.globalThis.vm().reportExtraMemory(bytes.len); } log("dataFrame enqueued {}", .{frame.len}); this.dataFrameQueue.enqueue(frame, client.allocator); @@ -1494,7 +1494,7 @@ pub const H2FrameParser = struct { // we still have more to buffer and even more now _ = this.writeBuffer.write(this.allocator, bytes) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(bytes.len); + // this.globalThis.vm().reportExtraMemory(bytes.len); log("_genericWrite flushed {} and buffered more {}", .{ written, bytes.len }); return false; @@ -1510,7 +1510,7 @@ pub const H2FrameParser = struct { const pending = bytes[written..]; // ops not all data was sent, lets buffer again _ = this.writeBuffer.write(this.allocator, pending) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(pending.len); + // this.globalThis.vm().reportExtraMemory(pending.len); log("_genericWrite buffered more {}", .{pending.len}); return false; @@ -1530,7 +1530,7 @@ pub const H2FrameParser = struct { const pending = bytes[written..]; // ops not all data was sent, lets buffer again _ = this.writeBuffer.write(this.allocator, pending) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(pending.len); + // this.globalThis.vm().reportExtraMemory(pending.len); return false; } @@ -1558,9 +1558,6 @@ pub const H2FrameParser = struct { pub fn flush(this: *H2FrameParser) usize { this.ref(); defer this.deref(); - const eventLoop = JSC.VirtualMachine.get().eventLoop(); - eventLoop.enter(); - defer eventLoop.exit(); var written = switch (this.native_socket) { .tls_writeonly, .tls => |socket| this._genericFlush(*TLSSocket, socket), .tcp_writeonly, .tcp => |socket| this._genericFlush(*TCPSocket, socket), @@ -1608,7 +1605,7 @@ pub const H2FrameParser = struct { if (this.has_nonnative_backpressure) { // we should not invoke JS when we have backpressure is cheaper to keep it queued here _ = this.writeBuffer.write(this.allocator, bytes) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(bytes.len); + // this.globalThis.vm().reportExtraMemory(bytes.len); return false; } @@ -1620,7 +1617,7 @@ pub const H2FrameParser = struct { -1 => { // dropped _ = this.writeBuffer.write(this.allocator, bytes) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(bytes.len); + // this.globalThis.vm().reportExtraMemory(bytes.len); this.has_nonnative_backpressure = true; }, 0 => { @@ -1707,7 +1704,7 @@ pub const H2FrameParser = struct { if (this.remainingLength > 0) { // buffer more data _ = this.readBuffer.appendSlice(payload) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(payload.len); + // this.globalThis.vm().reportExtraMemory(payload.len); return null; } else if (this.remainingLength < 0) { @@ -1720,7 +1717,7 @@ pub const H2FrameParser = struct { if (this.readBuffer.list.items.len > 0) { // return buffered data _ = this.readBuffer.appendSlice(payload) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(payload.len); + // this.globalThis.vm().reportExtraMemory(payload.len); return .{ .data = this.readBuffer.list.items, @@ -2248,7 +2245,7 @@ pub const H2FrameParser = struct { if (total < FrameHeader.byteSize) { // buffer more data _ = this.readBuffer.appendSlice(bytes) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(bytes.len); + // this.globalThis.vm().reportExtraMemory(bytes.len); return bytes.len; } @@ -2288,7 +2285,7 @@ pub const H2FrameParser = struct { if (bytes.len < FrameHeader.byteSize) { // buffer more dheaderata this.readBuffer.appendSlice(bytes) catch bun.outOfMemory(); - this.globalThis.vm().reportExtraMemory(bytes.len); + // this.globalThis.vm().reportExtraMemory(bytes.len); return bytes.len; } @@ -3639,9 +3636,6 @@ pub const H2FrameParser = struct { buffer.ensureStillAlive(); if (buffer.asArrayBuffer(globalObject)) |array_buffer| { var bytes = array_buffer.byteSlice(); - const eventLoop = globalObject.bunVM().eventLoop(); - eventLoop.enter(); - defer eventLoop.exit(); // read all the bytes while (bytes.len > 0) { const result = this.readBytes(bytes); @@ -3658,9 +3652,6 @@ pub const H2FrameParser = struct { this.ref(); defer this.deref(); var bytes = data; - const eventLoop = JSC.VirtualMachine.get().eventLoop(); - eventLoop.enter(); - defer eventLoop.exit(); while (bytes.len > 0) { const result = this.readBytes(bytes); bytes = bytes[result..]; From e3cd90cc6fb67ddbf77f7821481bba2154932440 Mon Sep 17 00:00:00 2001 From: cirospaciari Date: Wed, 16 Oct 2024 01:27:05 -0700 Subject: [PATCH 06/12] maybe --- src/bun.js/api/bun/h2_frame_parser.zig | 21 +++++++++++---------- src/js/node/http2.ts | 5 ++--- 2 files changed, 13 insertions(+), 13 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 26110481bd64..bb3e0eef2331 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1051,7 +1051,7 @@ pub const H2FrameParser = struct { }; if (bytes.len > 0) { @memcpy(frame.buffer[0..bytes.len], bytes); - // client.globalThis.vm().reportExtraMemory(bytes.len); + client.globalThis.vm().reportExtraMemory(bytes.len); } log("dataFrame enqueued {}", .{frame.len}); this.dataFrameQueue.enqueue(frame, client.allocator); @@ -1494,7 +1494,7 @@ pub const H2FrameParser = struct { // we still have more to buffer and even more now _ = this.writeBuffer.write(this.allocator, bytes) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(bytes.len); + this.globalThis.vm().reportExtraMemory(bytes.len); log("_genericWrite flushed {} and buffered more {}", .{ written, bytes.len }); return false; @@ -1510,7 +1510,7 @@ pub const H2FrameParser = struct { const pending = bytes[written..]; // ops not all data was sent, lets buffer again _ = this.writeBuffer.write(this.allocator, pending) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(pending.len); + this.globalThis.vm().reportExtraMemory(pending.len); log("_genericWrite buffered more {}", .{pending.len}); return false; @@ -1530,7 +1530,7 @@ pub const H2FrameParser = struct { const pending = bytes[written..]; // ops not all data was sent, lets buffer again _ = this.writeBuffer.write(this.allocator, pending) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(pending.len); + this.globalThis.vm().reportExtraMemory(pending.len); return false; } @@ -1605,7 +1605,7 @@ pub const H2FrameParser = struct { if (this.has_nonnative_backpressure) { // we should not invoke JS when we have backpressure is cheaper to keep it queued here _ = this.writeBuffer.write(this.allocator, bytes) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(bytes.len); + this.globalThis.vm().reportExtraMemory(bytes.len); return false; } @@ -1617,7 +1617,7 @@ pub const H2FrameParser = struct { -1 => { // dropped _ = this.writeBuffer.write(this.allocator, bytes) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(bytes.len); + this.globalThis.vm().reportExtraMemory(bytes.len); this.has_nonnative_backpressure = true; }, 0 => { @@ -1704,7 +1704,7 @@ pub const H2FrameParser = struct { if (this.remainingLength > 0) { // buffer more data _ = this.readBuffer.appendSlice(payload) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(payload.len); + this.globalThis.vm().reportExtraMemory(payload.len); return null; } else if (this.remainingLength < 0) { @@ -1717,7 +1717,7 @@ pub const H2FrameParser = struct { if (this.readBuffer.list.items.len > 0) { // return buffered data _ = this.readBuffer.appendSlice(payload) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(payload.len); + this.globalThis.vm().reportExtraMemory(payload.len); return .{ .data = this.readBuffer.list.items, @@ -2245,7 +2245,7 @@ pub const H2FrameParser = struct { if (total < FrameHeader.byteSize) { // buffer more data _ = this.readBuffer.appendSlice(bytes) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(bytes.len); + this.globalThis.vm().reportExtraMemory(bytes.len); return bytes.len; } @@ -2285,7 +2285,7 @@ pub const H2FrameParser = struct { if (bytes.len < FrameHeader.byteSize) { // buffer more dheaderata this.readBuffer.appendSlice(bytes) catch bun.outOfMemory(); - // this.globalThis.vm().reportExtraMemory(bytes.len); + this.globalThis.vm().reportExtraMemory(bytes.len); return bytes.len; } @@ -3676,6 +3676,7 @@ pub const H2FrameParser = struct { } const socket_js = args_list.ptr[0]; + this.detachNativeSocket(); if (JSTLSSocket.fromJS(socket_js)) |socket| { log("TLSSocket attached", .{}); if (socket.attachNativeCallback(.{ .h2 = this })) { diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index a21d399f0912..ed04f2c76ab1 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2394,8 +2394,7 @@ class ServerHttp2Session extends Http2Session { } #onClose() { - // this.destroy(); - this.close(); + this.destroy(); } #onError(error: Error) { @@ -2842,7 +2841,7 @@ class ClientHttp2Session extends Http2Session { } #onClose() { - this.close(); + this.destroy(); } #onError(error: Error) { this.destroy(error); From beb6c0553f6d557f554be5bb6b7fe0179aedb5a0 Mon Sep 17 00:00:00 2001 From: Ciro Spaciari Date: Wed, 16 Oct 2024 02:08:09 -0700 Subject: [PATCH 07/12] wip --- src/js/node/http2.ts | 48 ++++++++++++------- .../third_party/grpc-js/test-server.test.ts | 3 +- 2 files changed, 32 insertions(+), 19 deletions(-) diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index ed04f2c76ab1..c4b6aa248748 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2145,23 +2145,27 @@ function emitConnectNT(self, socket) { function emitStreamErrorNT(self, stream, error, destroy, destroy_self) { if (stream) { - let error_instance: Error | number | undefined = undefined; - if (typeof error === "number") { - stream.rstCode = error; - if (error != 0) { - error_instance = streamErrorFromCode(error); + const status = stream[bunHTTP2StreamStatus]; + + if ((status & StreamState.Closed) === 0) { + let error_instance: Error | number | undefined = undefined; + if (typeof error === "number") { + stream.rstCode = error; + if (error != 0) { + error_instance = streamErrorFromCode(error); + } + } else { + error_instance = error; + } + if (stream.readable) { + stream.resume(); // we have a error we consume and close + pushToStream(stream, null); + } + markStreamClosed(stream); + if (destroy) stream.destroy(error_instance, stream.rstCode); + else if (error_instance) { + stream.emit("error", error_instance); } - } else { - error_instance = error; - } - if (stream.readable) { - stream.resume(); // we have a error we consume and close - pushToStream(stream, null); - } - markStreamClosed(stream); - if (destroy) stream.destroy(error_instance, stream.rstCode); - else if (error_instance) { - stream.emit("error", error_instance); } if (destroy_self) self.destroy(); } @@ -2394,10 +2398,14 @@ class ServerHttp2Session extends Http2Session { } #onClose() { - this.destroy(); + this[bunHTTP2Socket] = null; + this.#parser?.emitErrorToAllStreams(NGHTTP2_CANCEL); + this.#parser = null; + this.close(); } #onError(error: Error) { + this[bunHTTP2Socket] = null; this.destroy(error); } @@ -2841,9 +2849,13 @@ class ClientHttp2Session extends Http2Session { } #onClose() { - this.destroy(); + this[bunHTTP2Socket] = null; + this.#parser?.emitErrorToAllStreams(NGHTTP2_CANCEL); + this.#parser = null; + this.close(); } #onError(error: Error) { + this[bunHTTP2Socket] = null; this.destroy(error); } #onTimeout() { diff --git a/test/js/third_party/grpc-js/test-server.test.ts b/test/js/third_party/grpc-js/test-server.test.ts index e992a89f8ccc..fc0b4f199be6 100644 --- a/test/js/third_party/grpc-js/test-server.test.ts +++ b/test/js/third_party/grpc-js/test-server.test.ts @@ -171,7 +171,7 @@ describe("Server", () => { }); }); - it("successfully unbinds a bound ephemeral port", done => { + it.todo("successfully unbinds a bound ephemeral port", done => { server.bindAsync("localhost:0", ServerCredentials.createInsecure(), (err, port) => { client = new grpc.Client(`localhost:${port}`, grpc.credentials.createInsecure()); client.makeUnaryRequest( @@ -194,6 +194,7 @@ describe("Server", () => { { deadline: deadline }, (callError2, result) => { assert(callError2); + console.error(callError2.code); // DEADLINE_EXCEEDED means that the server is unreachable assert( callError2.code === grpc.status.DEADLINE_EXCEEDED || callError2.code === grpc.status.UNAVAILABLE, From dc4b0c2c1665d3581c1691715972338440c1683c Mon Sep 17 00:00:00 2001 From: Ciro Spaciari Date: Wed, 16 Oct 2024 03:39:41 -0700 Subject: [PATCH 08/12] more --- src/bun.js/api/bun/h2_frame_parser.zig | 23 ++++++++- src/bun.js/api/h2.classes.ts | 4 ++ src/js/node/http2.ts | 49 +++++++++---------- test/js/node/http2/node-http2.test.js | 2 +- .../third_party/grpc-js/test-server.test.ts | 3 +- 5 files changed, 51 insertions(+), 30 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index bb3e0eef2331..d953e66f44ff 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -3254,7 +3254,25 @@ pub const H2FrameParser = struct { } return array; } - + pub fn emitAbortToAllStreams(this: *H2FrameParser, _: *JSC.JSGlobalObject, _: *JSC.CallFrame) JSC.JSValue { + JSC.markBinding(@src()); + var it = StreamResumableIterator.init(this); + while (it.next()) |stream| { + if (this.isServer) { + if (stream.id % 2 == 0) continue; + } else if (stream.id % 2 != 0) continue; + if (stream.state != .CLOSED) { + const old_state = stream.state; + stream.state = .CLOSED; + stream.rstCode = @intFromEnum(ErrorCode.CANCEL); + const identifier = stream.getIdentifier(); + identifier.ensureStillAlive(); + stream.freeResources(this, false); + this.dispatchWith2Extra(.onAborted, identifier, .undefined, JSC.JSValue.jsNumber(@intFromEnum(old_state))); + } + } + return .undefined; + } pub fn emitErrorToAllStreams(this: *H2FrameParser, globalObject: *JSC.JSGlobalObject, callframe: *JSC.CallFrame) JSC.JSValue { JSC.markBinding(@src()); @@ -3266,6 +3284,9 @@ pub const H2FrameParser = struct { var it = StreamResumableIterator.init(this); while (it.next()) |stream| { + if (this.isServer) { + if (stream.id % 2 == 0) continue; + } else if (stream.id % 2 != 0) continue; if (stream.state != .CLOSED) { stream.state = .CLOSED; stream.rstCode = args_list.ptr[0].to(u32); diff --git a/src/bun.js/api/h2.classes.ts b/src/bun.js/api/h2.classes.ts index dab1dd2d5ba5..9d234c87ff57 100644 --- a/src/bun.js/api/h2.classes.ts +++ b/src/bun.js/api/h2.classes.ts @@ -93,6 +93,10 @@ export default [ fn: "emitErrorToAllStreams", length: 1, }, + emitAbortToAllStreams: { + fn: "emitAbortToAllStreams", + length: 0, + }, getNextStream: { fn: "getNextStream", length: 0, diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index c4b6aa248748..b3d594530fed 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -1541,6 +1541,7 @@ function markStreamClosed(stream: Http2Stream) { if ((status & StreamState.Closed) === 0) { stream[bunHTTP2StreamStatus] = status | StreamState.Closed; + markWritableDone(stream); } } @@ -2145,28 +2146,26 @@ function emitConnectNT(self, socket) { function emitStreamErrorNT(self, stream, error, destroy, destroy_self) { if (stream) { - const status = stream[bunHTTP2StreamStatus]; - - if ((status & StreamState.Closed) === 0) { - let error_instance: Error | number | undefined = undefined; - if (typeof error === "number") { - stream.rstCode = error; - if (error != 0) { - error_instance = streamErrorFromCode(error); - } - } else { - error_instance = error; - } - if (stream.readable) { - stream.resume(); // we have a error we consume and close - pushToStream(stream, null); - } - markStreamClosed(stream); - if (destroy) stream.destroy(error_instance, stream.rstCode); - else if (error_instance) { - stream.emit("error", error_instance); + let error_instance: Error | number | undefined = undefined; + if (typeof error === "number") { + stream.rstCode = error; + if (error != 0) { + error_instance = streamErrorFromCode(error); } + } else { + error_instance = error; } + const status = stream[bunHTTP2StreamStatus]; + if (stream.readable) { + stream.resume(); // we have a error we consume and close + pushToStream(stream, null); + } + markStreamClosed(stream); + if (destroy) stream.destroy(error_instance, stream.rstCode); + else if (error_instance) { + stream.emit("error", error_instance); + } + if (destroy_self) self.destroy(); } } @@ -2251,9 +2250,7 @@ class ServerHttp2Session extends Http2Session { }, aborted(self: ServerHttp2Session, stream: ServerHttp2Stream, error: any, old_state: number) { if (!self || typeof stream !== "object") return; - stream.rstCode = constants.NGHTTP2_CANCEL; - markStreamClosed(stream); // if writable and not closed emit aborted if (old_state != 5 && old_state != 7) { stream[kAborted] = true; @@ -2399,7 +2396,7 @@ class ServerHttp2Session extends Http2Session { #onClose() { this[bunHTTP2Socket] = null; - this.#parser?.emitErrorToAllStreams(NGHTTP2_CANCEL); + this.#parser?.emitAbortToAllStreams(); this.#parser = null; this.close(); } @@ -2663,8 +2660,6 @@ class ClientHttp2Session extends Http2Session { }, aborted(self: ClientHttp2Session, stream: ClientHttp2Stream, error: any, old_state: number) { if (!self || typeof stream !== "object") return; - - markStreamClosed(stream); stream.rstCode = constants.NGHTTP2_CANCEL; // if writable and not closed emit aborted if (old_state != 5 && old_state != 7) { @@ -2672,11 +2667,13 @@ class ClientHttp2Session extends Http2Session { stream.emit("aborted"); } self.#connections--; + process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); }, streamError(self: ClientHttp2Session, stream: ClientHttp2Stream, error: number) { if (!self || typeof stream !== "object") return; self.#connections--; + process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); }, streamEnd(self: ClientHttp2Session, stream: ClientHttp2Stream, state: number) { @@ -2850,7 +2847,7 @@ class ClientHttp2Session extends Http2Session { #onClose() { this[bunHTTP2Socket] = null; - this.#parser?.emitErrorToAllStreams(NGHTTP2_CANCEL); + this.#parser?.emitAbortToAllStreams(); this.#parser = null; this.close(); } diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index c75a0f5cb0cb..5028676c51ad 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -10,7 +10,7 @@ import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from import http2utils from "./helpers"; import { nodeEchoServer, TLS_CERT, TLS_OPTIONS } from "./http2-helpers"; -for (const nodeExecutable of [nodeExe()]) { +for (const nodeExecutable of [nodeExe(), bunExe()]) { describe(`${path.basename(nodeExecutable)}`, () => { let nodeEchoServer_; diff --git a/test/js/third_party/grpc-js/test-server.test.ts b/test/js/third_party/grpc-js/test-server.test.ts index fc0b4f199be6..e992a89f8ccc 100644 --- a/test/js/third_party/grpc-js/test-server.test.ts +++ b/test/js/third_party/grpc-js/test-server.test.ts @@ -171,7 +171,7 @@ describe("Server", () => { }); }); - it.todo("successfully unbinds a bound ephemeral port", done => { + it("successfully unbinds a bound ephemeral port", done => { server.bindAsync("localhost:0", ServerCredentials.createInsecure(), (err, port) => { client = new grpc.Client(`localhost:${port}`, grpc.credentials.createInsecure()); client.makeUnaryRequest( @@ -194,7 +194,6 @@ describe("Server", () => { { deadline: deadline }, (callError2, result) => { assert(callError2); - console.error(callError2.code); // DEADLINE_EXCEEDED means that the server is unreachable assert( callError2.code === grpc.status.DEADLINE_EXCEEDED || callError2.code === grpc.status.UNAVAILABLE, From 6a6d16017933db2d1fd3aac0dccd96db073392c8 Mon Sep 17 00:00:00 2001 From: Ciro Spaciari Date: Wed, 16 Oct 2024 04:23:01 -0700 Subject: [PATCH 09/12] make sure everything is cleanup --- src/bun.js/api/bun/h2_frame_parser.zig | 39 ++++++++++------ src/bun.js/api/h2.classes.ts | 4 ++ src/js/node/http2.ts | 65 +++++++++++++------------- 3 files changed, 63 insertions(+), 45 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index d953e66f44ff..13263abbf5da 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -3882,17 +3882,15 @@ pub const H2FrameParser = struct { } return this; } - - pub fn deinit(this: *H2FrameParser) void { - log("deinit", .{}); - - defer { - if (ENABLE_ALLOCATOR_POOL) { - H2FrameParser.pool.?.put(this); - } else { - this.destroy(); - } - } + pub fn detachFromJS(this: *H2FrameParser, _: *JSC.JSGlobalObject, _: *JSC.CallFrame) JSValue { + JSC.markBinding(@src()); + this.detach(false); + return .undefined; + } + /// be careful when calling detach be sure that the socket is closed and the parser not accesible anymore + /// this function can be called multiple times + pub fn detach(this: *H2FrameParser, comptime finalizing: bool) void { + this.flushCorked(); this.detachNativeSocket(); this.strong_ctx.deinit(); this.handlers.deinit(); @@ -3909,9 +3907,24 @@ pub const H2FrameParser = struct { } var it = this.streams.valueIterator(); while (it.next()) |stream| { - stream.freeResources(this, true); + stream.freeResources(this, finalizing); + } + var streams = this.streams; + defer streams.deinit(); + this.streams = bun.U32HashMap(Stream).init(bun.default_allocator); + } + + pub fn deinit(this: *H2FrameParser) void { + log("deinit", .{}); + + defer { + if (ENABLE_ALLOCATOR_POOL) { + H2FrameParser.pool.?.put(this); + } else { + this.destroy(); + } } - this.streams.deinit(); + this.detach(true); } pub fn finalize( diff --git a/src/bun.js/api/h2.classes.ts b/src/bun.js/api/h2.classes.ts index 9d234c87ff57..bcad57f64a45 100644 --- a/src/bun.js/api/h2.classes.ts +++ b/src/bun.js/api/h2.classes.ts @@ -37,6 +37,10 @@ export default [ fn: "flushFromJS", length: 0, }, + detach: { + fn: "detachFromJS", + length: 0, + }, rstStream: { fn: "rstStream", length: 1, diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index b3d594530fed..6e448fc68731 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2344,10 +2344,7 @@ class ServerHttp2Session extends Http2Session { error(self: ServerHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; const error_instance = sessionErrorFromCode(errorCode); - self.emit("error", error_instance); - self[bunHTTP2Socket]?.end(); - self.#parser = null; - this[bunHTTP2Socket] = null; + self.destroy(error_instance); }, wantTrailers(self: ServerHttp2Session, stream: ServerHttp2Stream) { if (!self || typeof stream !== "object") return; @@ -2368,16 +2365,12 @@ class ServerHttp2Session extends Http2Session { if (errorCode !== 0) { self.#parser.emitErrorToAllStreams(errorCode); } - - self[bunHTTP2Socket]?.end(); - self.#parser = null; - this[bunHTTP2Socket] = null; + const error_instance = sessionErrorFromCode(errorCode); + self.destroy(error_instance); }, end(self: ServerHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; - self[bunHTTP2Socket]?.end(); - self.#parser = null; - this[bunHTTP2Socket] = null; + self.destroy(); }, write(self: ServerHttp2Session, buffer: Buffer) { if (!self) return -1; @@ -2395,14 +2388,16 @@ class ServerHttp2Session extends Http2Session { } #onClose() { - this[bunHTTP2Socket] = null; - this.#parser?.emitAbortToAllStreams(); - this.#parser = null; + const parser = this.#parser; + if (parser) { + parser.emitAbortToAllStreams(); + parser.detach(); + this.#parser = null; + } this.close(); } #onError(error: Error) { - this[bunHTTP2Socket] = null; this.destroy(error); } @@ -2609,8 +2604,12 @@ class ServerHttp2Session extends Http2Session { this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0)); socket.end(); } - this.#parser?.emitErrorToAllStreams(code || constants.NGHTTP2_NO_ERROR); - this.#parser = null; + const parser = this.#parser; + if (parser) { + parser.emitErrorToAllStreams(code || constants.NGHTTP2_NO_ERROR); + parser.detach(); + this.#parser = null; + } this[bunHTTP2Socket] = null; if (error) { @@ -2761,10 +2760,7 @@ class ClientHttp2Session extends Http2Session { error(self: ClientHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; const error_instance = sessionErrorFromCode(errorCode); - self.emit("error", error_instance); - self[bunHTTP2Socket]?.destroy(); - self.#parser = null; - this[bunHTTP2Socket] = null; + self.destroy(error_instance); }, wantTrailers(self: ClientHttp2Session, stream: ClientHttp2Stream) { @@ -2784,15 +2780,12 @@ class ClientHttp2Session extends Http2Session { if (errorCode !== 0) { self.#parser.emitErrorToAllStreams(errorCode); } - self[bunHTTP2Socket]?.end(); - self.#parser = null; - this[bunHTTP2Socket] = null; + const error_instance = sessionErrorFromCode(errorCode); + self.destroy(error_instance); }, end(self: ClientHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; - self[bunHTTP2Socket]?.end(); - self.#parser = null; - this[bunHTTP2Socket] = null; + self.destroy(); }, write(self: ClientHttp2Session, buffer: Buffer) { if (!self) return -1; @@ -2846,10 +2839,14 @@ class ClientHttp2Session extends Http2Session { } #onClose() { - this[bunHTTP2Socket] = null; - this.#parser?.emitAbortToAllStreams(); - this.#parser = null; + const parser = this.#parser; + if (parser) { + parser.emitAbortToAllStreams(); + parser.detach(); + this.#parser = null; + } this.close(); + this[bunHTTP2Socket] = null; } #onError(error: Error) { this[bunHTTP2Socket] = null; @@ -3069,9 +3066,13 @@ class ClientHttp2Session extends Http2Session { this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0)); socket.end(); } - this.#parser?.emitErrorToAllStreams(code || constants.NGHTTP2_NO_ERROR); - this[bunHTTP2Socket] = null; + const parser = this.#parser; + if (parser) { + parser.emitErrorToAllStreams(code || constants.NGHTTP2_NO_ERROR); + parser.detach(); + } this.#parser = null; + this[bunHTTP2Socket] = null; if (error) { this.emit("error", error); From c81cac6d165e6a11a97495afed84e74f74b108c9 Mon Sep 17 00:00:00 2001 From: Ciro Spaciari Date: Wed, 16 Oct 2024 04:55:22 -0700 Subject: [PATCH 10/12] more progress --- src/bun.js/api/bun/h2_frame_parser.zig | 27 ++++++++++++++++++-------- src/js/node/http2.ts | 6 ++---- 2 files changed, 21 insertions(+), 12 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 13263abbf5da..78518de7ee55 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1758,7 +1758,7 @@ pub const H2FrameParser = struct { return data.len; } - pub fn decodeHeaderBlock(this: *H2FrameParser, payload: []const u8, stream: *Stream, flags: u8) *Stream { + pub fn decodeHeaderBlock(this: *H2FrameParser, payload: []const u8, stream: *Stream, flags: u8) ?*Stream { log("decodeHeaderBlock isSever: {}", .{this.isServer}); var offset: usize = 0; @@ -1777,7 +1777,9 @@ pub const H2FrameParser = struct { log("header {s} {s}", .{ header.name, header.value }); if (this.isServer and strings.eqlComptime(header.name, ":status")) { this.sendGoAway(stream_id, ErrorCode.PROTOCOL_ERROR, "Server received :status header", this.lastStreamID, true); - return this.streams.getEntry(stream_id).?.value_ptr; + + if (this.streams.getEntry(stream_id)) |entry| return entry.value_ptr; + return null; } count += 1; if (this.maxHeaderListPairs < count) { @@ -1787,7 +1789,8 @@ pub const H2FrameParser = struct { } else { this.endStream(stream, ErrorCode.ENHANCE_YOUR_CALM); } - return this.streams.getEntry(stream_id).?.value_ptr; + if (this.streams.getEntry(stream_id)) |entry| return entry.value_ptr; + return null; } const output = brk: { @@ -1818,7 +1821,8 @@ pub const H2FrameParser = struct { this.dispatchWith3Extra(.onStreamHeaders, stream.getIdentifier(), headers, sensitiveHeaders, JSC.JSValue.jsNumber(flags)); // callbacks can change the Stream ptr in this case we always return the new one - return this.streams.getEntry(stream_id).?.value_ptr; + if (this.streams.getEntry(stream_id)) |entry| return entry.value_ptr; + return null; } pub fn handleDataFrame(this: *H2FrameParser, frame: FrameHeader, data: []const u8, stream_: ?*Stream) usize { @@ -1883,7 +1887,8 @@ pub const H2FrameParser = struct { this.currentFrame = null; if (emitted) { // we need to revalidate the stream ptr after emitting onStreamData - stream = this.streams.getEntry(frame.streamIdentifier).?.value_ptr; + const entry = this.streams.getEntry(frame.streamIdentifier) orelse return end; + stream = entry.value_ptr; } if (frame.flags & @intFromEnum(DataFrameFlags.END_STREAM) != 0) { const identifier = stream.getIdentifier(); @@ -2030,7 +2035,10 @@ pub const H2FrameParser = struct { } if (handleIncommingPayload(this, data, frame.streamIdentifier)) |content| { const payload = content.data; - stream = this.decodeHeaderBlock(payload[0..payload.len], stream, frame.flags); + stream = this.decodeHeaderBlock(payload[0..payload.len], stream, frame.flags) orelse { + this.readBuffer.reset(); + return content.end; + }; this.readBuffer.reset(); if (frame.flags & @intFromEnum(HeadersFrameFlags.END_HEADERS) != 0) { stream.isWaitingMoreHeaders = false; @@ -2093,7 +2101,10 @@ pub const H2FrameParser = struct { this.sendGoAway(frame.streamIdentifier, ErrorCode.FRAME_SIZE_ERROR, "invalid Headers frame size", this.lastStreamID, true); return data.len; } - stream = this.decodeHeaderBlock(payload[offset..end], stream, frame.flags); + stream = this.decodeHeaderBlock(payload[offset..end], stream, frame.flags) orelse { + this.readBuffer.reset(); + return content.end; + }; this.readBuffer.reset(); stream.isWaitingMoreHeaders = frame.flags & @intFromEnum(HeadersFrameFlags.END_HEADERS) == 0; if (frame.flags & @intFromEnum(HeadersFrameFlags.END_STREAM) != 0) { @@ -3888,7 +3899,7 @@ pub const H2FrameParser = struct { return .undefined; } /// be careful when calling detach be sure that the socket is closed and the parser not accesible anymore - /// this function can be called multiple times + /// this function can be called multiple times, it will erase stream info pub fn detach(this: *H2FrameParser, comptime finalizing: bool) void { this.flushCorked(); this.detachNativeSocket(); diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 6e448fc68731..0f873a71f461 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -2365,8 +2365,7 @@ class ServerHttp2Session extends Http2Session { if (errorCode !== 0) { self.#parser.emitErrorToAllStreams(errorCode); } - const error_instance = sessionErrorFromCode(errorCode); - self.destroy(error_instance); + self.close(); }, end(self: ServerHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; @@ -2780,8 +2779,7 @@ class ClientHttp2Session extends Http2Session { if (errorCode !== 0) { self.#parser.emitErrorToAllStreams(errorCode); } - const error_instance = sessionErrorFromCode(errorCode); - self.destroy(error_instance); + self.close(); }, end(self: ClientHttp2Session, errorCode: number, lastStreamId: number, opaqueData: Buffer) { if (!self) return; From c1212882bd225533ce396fb1967ae06132eb258a Mon Sep 17 00:00:00 2001 From: Ciro Spaciari Date: Wed, 16 Oct 2024 06:29:04 -0700 Subject: [PATCH 11/12] more --- src/bun.js/api/bun/h2_frame_parser.zig | 5 +- src/js/node/http2.ts | 71 ++++++++++++-------------- test/js/node/http2/node-http2.test.js | 25 +-------- 3 files changed, 36 insertions(+), 65 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index 78518de7ee55..fbbae9a47764 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -3269,6 +3269,7 @@ pub const H2FrameParser = struct { JSC.markBinding(@src()); var it = StreamResumableIterator.init(this); while (it.next()) |stream| { + // this is the oposite logic of emitErrorToallStreams, in this case we wanna to cancel this streams if (this.isServer) { if (stream.id % 2 == 0) continue; } else if (stream.id % 2 != 0) continue; @@ -3296,8 +3297,8 @@ pub const H2FrameParser = struct { var it = StreamResumableIterator.init(this); while (it.next()) |stream| { if (this.isServer) { - if (stream.id % 2 == 0) continue; - } else if (stream.id % 2 != 0) continue; + if (stream.id % 2 != 0) continue; + } else if (stream.id % 2 == 0) continue; if (stream.state != .CLOSED) { stream.state = .CLOSED; stream.rstCode = args_list.ptr[0].to(u32); diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 0f873a71f461..7fba7caae747 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -1718,50 +1718,46 @@ class Http2Stream extends Duplex { } } _destroy(err, callback) { - if ((this[bunHTTP2StreamStatus] & StreamState.Closed) === 0) { - const { ending } = this._writableState; - if (!ending) { - // If the writable side of the Http2Stream is still open, emit the - // 'aborted' event and set the aborted flag. - if (!this.aborted) { - 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; + const { ending } = this._writableState; + + if (!ending) { + // If the writable side of the Http2Stream is still open, emit the + // 'aborted' event and set the aborted flag. + if (!this.aborted) { + 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; + } - const session = this[bunHTTP2Session]; - assertSession(session); + const session = this[bunHTTP2Session]; + assertSession(session); - let rstCode = this.rstCode; - if (!rstCode) { - if (err != null) { - if (err.code === "ABORT_ERR") { - // Enables using AbortController to cancel requests with RST code 8. - rstCode = NGHTTP2_CANCEL; - } else { - rstCode = NGHTTP2_INTERNAL_ERROR; - } + let rstCode = this.rstCode; + if (!rstCode) { + if (err != null) { + if (err.code === "ABORT_ERR") { + // Enables using AbortController to cancel requests with RST code 8. + rstCode = NGHTTP2_CANCEL; } else { - rstCode = this.rstCode = 0; + rstCode = NGHTTP2_INTERNAL_ERROR; } + } else { + rstCode = this.rstCode = 0; } + } - if (this.writableFinished) { - markStreamClosed(this); + if (this.writableFinished) { + markStreamClosed(this); - session[bunHTTP2Native]?.rstStream(this.#id, rstCode); - this[bunHTTP2Session] = null; - } else { - this.once("finish", Http2Stream.#rstStream); - } - } else { + session[bunHTTP2Native]?.rstStream(this.#id, rstCode); this[bunHTTP2Session] = null; + } else { + this.once("finish", Http2Stream.#rstStream); } callback(err); @@ -2155,7 +2151,7 @@ function emitStreamErrorNT(self, stream, error, destroy, destroy_self) { } else { error_instance = error; } - const status = stream[bunHTTP2StreamStatus]; + if (stream.readable) { stream.resume(); // we have a error we consume and close pushToStream(stream, null); @@ -2256,7 +2252,6 @@ class ServerHttp2Session extends Http2Session { stream[kAborted] = true; stream.emit("aborted"); } - self.#connections--; process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); }, @@ -2665,13 +2660,11 @@ class ClientHttp2Session extends Http2Session { stream.emit("aborted"); } self.#connections--; - process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); }, streamError(self: ClientHttp2Session, stream: ClientHttp2Stream, error: number) { if (!self || typeof stream !== "object") return; self.#connections--; - process.nextTick(emitStreamErrorNT, self, stream, error, true, self.#connections === 0 && self.#closed); }, streamEnd(self: ClientHttp2Session, stream: ClientHttp2Stream, state: number) { diff --git a/test/js/node/http2/node-http2.test.js b/test/js/node/http2/node-http2.test.js index 5028676c51ad..6d19fe6dd1e5 100644 --- a/test/js/node/http2/node-http2.test.js +++ b/test/js/node/http2/node-http2.test.js @@ -665,30 +665,7 @@ for (const nodeExecutable of [nodeExe(), bunExe()]) { expect(req.aborted).toBeTrue(); expect(req.rstCode).toBe(http2.constants.NGHTTP2_CANCEL); }); - it("aborted event should not work when not writable but should emit error", async () => { - const abortController = new AbortController(); - const { promise, resolve, reject } = Promise.withResolvers(); - const client = http2.connect(HTTPS_SERVER, TLS_OPTIONS); - client.on("error", reject); - const req = client.request({ ":path": "/" }, { signal: abortController.signal }); - req.on("aborted", reject); - req.on("error", err => { - if (err.code !== "ABORT_ERR") { - reject(err); - } else { - resolve(); - } - }); - req.on("end", () => { - reject(); - client.close(); - }); - abortController.abort(); - const result = await promise; - expect(result).toBeUndefined(); - expect(req.aborted).toBeFalse(); // will only be true when the request is in a writable state - expect(req.rstCode).toBe(http2.constants.NGHTTP2_CANCEL); - }); + it("aborted event should work with aborted signal", async () => { const { promise, resolve, reject } = Promise.withResolvers(); const client = http2.connect(HTTPS_SERVER, TLS_OPTIONS); From 41d1ffbf759c070d307a847c88f39ab7ee78cbf8 Mon Sep 17 00:00:00 2001 From: Ciro Spaciari Date: Wed, 16 Oct 2024 06:45:21 -0700 Subject: [PATCH 12/12] :D --- src/bun.js/api/bun/h2_frame_parser.zig | 2 +- src/js/node/http2.ts | 1 - .../http2-connect-tls-with-delay.test.js | 50 +++++++++---------- 3 files changed, 26 insertions(+), 27 deletions(-) diff --git a/src/bun.js/api/bun/h2_frame_parser.zig b/src/bun.js/api/bun/h2_frame_parser.zig index fbbae9a47764..535db8e26fa9 100644 --- a/src/bun.js/api/bun/h2_frame_parser.zig +++ b/src/bun.js/api/bun/h2_frame_parser.zig @@ -1612,7 +1612,7 @@ pub const H2FrameParser = struct { // fallback to onWrite non-native callback const output_value = this.handlers.binary_type.toJS(bytes, this.handlers.globalObject); const result = this.call(.onWrite, output_value); - const code = result.to(i32); + const code = if (result.isNumber()) result.to(i32) else -1; switch (code) { -1 => { // dropped diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index 7fba7caae747..72936d97851a 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -1710,7 +1710,6 @@ class Http2Stream extends Duplex { markStreamClosed(this); session[bunHTTP2Native]?.rstStream(this.#id, code); - this[bunHTTP2Session] = null; } if (typeof callback === "function") { diff --git a/test/js/node/test/parallel/http2-connect-tls-with-delay.test.js b/test/js/node/test/parallel/http2-connect-tls-with-delay.test.js index 8e70ca287039..1161272cabe0 100644 --- a/test/js/node/test/parallel/http2-connect-tls-with-delay.test.js +++ b/test/js/node/test/parallel/http2-connect-tls-with-delay.test.js @@ -1,54 +1,54 @@ //#FILE: test-http2-connect-tls-with-delay.js //#SHA1: 8c5489e025ec14c2cc53788b27fde11a11990e42 //----------------- -'use strict'; +"use strict"; -const http2 = require('http2'); -const tls = require('tls'); -const fs = require('fs'); -const path = require('path'); +const http2 = require("http2"); +const tls = require("tls"); +const fs = require("fs"); +const path = require("path"); const serverOptions = { - key: fs.readFileSync(path.join(__dirname, '..', 'fixtures', 'keys', 'agent1-key.pem')), - cert: fs.readFileSync(path.join(__dirname, '..', 'fixtures', 'keys', 'agent1-cert.pem')) + key: fs.readFileSync(path.join(__dirname, "..", "fixtures", "keys", "agent1-key.pem")), + cert: fs.readFileSync(path.join(__dirname, "..", "fixtures", "keys", "agent1-cert.pem")), }; let server; -beforeAll((done) => { +beforeAll(done => { server = http2.createSecureServer(serverOptions, (req, res) => { res.end(); }); - server.listen(0, '127.0.0.1', done); + server.listen(0, "127.0.0.1", done); }); -afterAll((done) => { - server.close(done); +afterAll(() => { + server.close(); }); -test('HTTP/2 connect with TLS and delay', (done) => { +test("HTTP/2 connect with TLS and delay", done => { const options = { - ALPNProtocols: ['h2'], - host: '127.0.0.1', - servername: 'localhost', + ALPNProtocols: ["h2"], + host: "127.0.0.1", + servername: "localhost", port: server.address().port, - rejectUnauthorized: false + rejectUnauthorized: false, }; const socket = tls.connect(options, async () => { - socket.once('readable', () => { - const client = http2.connect( - 'https://localhost:' + server.address().port, - { ...options, createConnection: () => socket } - ); + socket.once("readable", () => { + const client = http2.connect("https://localhost:" + server.address().port, { + ...options, + createConnection: () => socket, + }); - client.once('remoteSettings', () => { + client.once("remoteSettings", () => { const req = client.request({ - ':path': '/' + ":path": "/", }); - req.on('data', () => req.resume()); - req.on('end', () => { + req.on("data", () => req.resume()); + req.on("end", () => { client.close(); req.close(); done();