Skip to content
Merged

H2 fixes #14606

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
92 changes: 70 additions & 22 deletions src/bun.js/api/bun/h2_frame_parser.zig
Original file line number Diff line number Diff line change
Expand Up @@ -1141,6 +1141,7 @@ pub const H2FrameParser = struct {
this.signal = null;
signal.deinit();
}
JSC.VirtualMachine.get().eventLoop().processGCTimer();
}
};

Expand Down Expand Up @@ -1611,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
Expand Down Expand Up @@ -1757,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;
Expand All @@ -1776,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) {
Expand All @@ -1786,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: {
Expand Down Expand Up @@ -1817,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 {
Expand Down Expand Up @@ -1882,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();
Expand Down Expand Up @@ -2029,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;
Expand Down Expand Up @@ -2092,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) {
Expand Down Expand Up @@ -3253,7 +3265,26 @@ 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| {
// 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;
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());

Expand All @@ -3265,6 +3296,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);
Expand Down Expand Up @@ -3675,6 +3709,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 })) {
Expand Down Expand Up @@ -3859,17 +3894,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, it will erase stream info
pub fn detach(this: *H2FrameParser, comptime finalizing: bool) void {
this.flushCorked();
this.detachNativeSocket();
this.strong_ctx.deinit();
this.handlers.deinit();
Expand All @@ -3886,9 +3919,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(
Expand Down
8 changes: 8 additions & 0 deletions src/bun.js/api/h2.classes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,10 @@ export default [
fn: "flushFromJS",
length: 0,
},
detach: {
fn: "detachFromJS",
length: 0,
},
rstStream: {
fn: "rstStream",
length: 1,
Expand Down Expand Up @@ -93,6 +97,10 @@ export default [
fn: "emitErrorToAllStreams",
length: 1,
},
emitAbortToAllStreams: {
fn: "emitAbortToAllStreams",
length: 0,
},
getNextStream: {
fn: "getNextStream",
length: 0,
Expand Down
Loading