Skip to content
Merged
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
77 changes: 66 additions & 11 deletions lib/protocol/http2/connection.rb
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,9 @@ def initialize(framer, local_stream_id)

@local_window = LocalWindow.new
@remote_window = Window.new

# The lowest Last-Stream-ID received in a GOAWAY frame, or nil if no GOAWAY frame has been received:
@goaway_stream_id = nil
end

# The connection stream ID (always 0 for connection-level operations).
Expand Down Expand Up @@ -94,16 +97,41 @@ def maximum_concurrent_streams
# The highest stream_id that has been successfully accepted by this connection.
attr :remote_stream_id

# The lowest Last-Stream-ID received in a GOAWAY frame.
attr :goaway_stream_id

# Whether the connection is effectively or actually closed.
def closed?
@state == :closed || @framer.nil?
end

# Whether the remote peer has sent us a GOAWAY frame. We must not initiate any new streams on this connection, but existing streams may still be in progress.
# @returns [Boolean] True if a GOAWAY frame has been received.
def goaway_received?
!@goaway_stream_id.nil?
end

# Transition the connection into the closed state if a graceful GOAWAY was received and there is nothing left to drain.
#
# As with {close!}, this is a state transition only: the owner of the connection is responsible for closing the underlying framer.
def close_if_drained!
if self.goaway_received? && @streams.empty?
self.close!
end
end

# Remove a stream from the active streams collection.
#
# If the remote peer has sent a graceful GOAWAY frame, the connection is only kept open in order to drain the streams it accepted, so when the last one completes there is nothing left to read.
#
# @parameter id [Integer] The stream ID to remove.
# @returns [Stream | Nil] The removed stream, or nil if not found.
def delete(id)
@streams.delete(id)
stream = @streams.delete(id)

self.close_if_drained!

return stream
end

# Close the underlying framer and all streams.
Expand Down Expand Up @@ -229,25 +257,39 @@ def send_goaway(error_code = 0, message = "")
end

# Process a GOAWAY frame from the remote peer.
#
# A GOAWAY frame with a zero error code is a graceful shutdown: the remote peer will not accept any new streams, but it is still processing the streams at or below `last_stream_id` and will send their responses (RFC 9113 §6.8). We must keep reading until those streams complete, otherwise requests which the remote peer has already processed - and whose side effects have already happened - fail locally. The connection is closed once the last of those streams completes, or the remote peer closes it.
#
# A GOAWAY frame with a non-zero error code is a connection error: the connection transitions into the closed state and {GoawayError} is raised.
#
# @parameter frame [GoawayFrame] The GOAWAY frame to process.
# @raises [GoawayError] If the frame indicates a connection error.
def receive_goaway(frame)
# We capture the last stream that was processed.
@remote_stream_id, error_code, message = frame.unpack
# We capture the last locally-initiated stream that may have been processed by the peer.
goaway_stream_id, error_code, message = frame.unpack

self.close!
# A peer can send an initial GOAWAY with a high stream ID, followed by another GOAWAY with a lower stream ID. The effective cutoff can only decrease (RFC 9113 §6.8).
if @goaway_stream_id.nil? || goaway_stream_id < @goaway_stream_id
@goaway_stream_id = goaway_stream_id
end

# Streams above the last stream ID were not processed by the remote peer and are safe to retry (RFC 9113 §6.8).
error = ::Protocol::HTTP::RefusedError.new("GOAWAY: request not processed.")
# Locally-initiated streams above the last stream ID were not processed by the remote peer and are safe to retry (RFC 9113 §6.8). They are removed from the connection before being closed, both so that what remains is exactly the set of streams we are waiting on, and so that closing them cannot mutate the collection while we are traversing it.
refused_streams = @streams.select{|id, stream| local_stream_id?(id) && id > @goaway_stream_id}
refused_streams.each_key{|id| @streams.delete(id)}

@streams.each_value do |stream|
if stream.id > @remote_stream_id
stream.close(error)
end
# The state of the connection is decided before any stream is closed, so that it cannot be left undecided by a `closed` hook which raises, and cannot be influenced by one which creates a stream.
if error_code != 0
self.close!
else
self.close_if_drained!
end

unless refused_streams.empty?
error = ::Protocol::HTTP::RefusedError.new("GOAWAY: request not processed.")
refused_streams.each_value{|stream| stream.close(error)}
end

if error_code != 0
# Shut down immediately.
raise GoawayError.new(message, error_code)
end
end
Expand Down Expand Up @@ -404,6 +446,14 @@ def valid_remote_stream_id?(stream_id)
false
end

# Check if the given stream ID represents a locally-initiated stream.
# This method should be overridden by client/server implementations.
# @parameter id [Integer] The stream ID to check.
# @returns [Boolean] True if the stream ID is locally-initiated.
def local_stream_id?(id)
false
end

# Accept an incoming stream from the other side of the connnection.
# On the server side, we accept requests.
def accept_stream(stream_id, &block)
Expand All @@ -425,6 +475,11 @@ def accept_push_promise_stream(stream_id, &block)
# On the client side, we create requests.
# @return [Stream] the created stream.
def create_stream(id = next_stream_id, &block)
if self.goaway_received? and local_stream_id?(id)
# Receivers of a GOAWAY frame MUST NOT open additional streams on the connection (RFC 9113 §6.8). A new connection has to be established for new streams.
raise ProtocolError, "Cannot create stream #{id} after GOAWAY!"
end

if @streams.key?(id)
raise ProtocolError, "Cannot create stream with id #{id}, already exists!"
end
Expand Down
5 changes: 5 additions & 0 deletions releases.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,10 @@
# Releases

## Unreleased

- On a graceful `GOAWAY` (error code `0`), keep the connection open until the streams the remote peer accepted have completed, instead of closing it immediately and failing those requests with `EOFError`.
- `Connection#create_stream` refuses to open a locally-initiated stream once a `GOAWAY` has been received, as required by RFC 9113 §6.8.

## v0.26.2

- Ignore the reserved high bit when decoding GOAWAY last stream IDs.
Expand Down
193 changes: 190 additions & 3 deletions test/protocol/http2/connection.rb
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,10 @@
expect(connection).not.to be(:valid_remote_stream_id?, 1)
end

it "does not report any stream_id as being local" do
expect(connection).not.to be(:local_stream_id?, 1)
end

it "rejects a push promise" do
frame = Protocol::HTTP2::PushPromiseFrame.new

Expand Down Expand Up @@ -374,8 +378,12 @@ def before
another_stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)

expect(client.read_frame).to be_a Protocol::HTTP2::GoawayFrame
expect(client.remote_stream_id).to be == 1
expect(client).to be(:closed?)
expect(client.goaway_stream_id).to be == 1
expect(client.remote_stream_id).to be == 0

# The server accepted stream 1 and is still processing it, so the connection is not closed yet:
expect(client).to be(:goaway_received?)
expect(client).not.to be(:closed?)

# The server will ignore this frame as it was sent after the graceful shutdown:
server.read_frame
Expand All @@ -389,6 +397,25 @@ def before
client.read_frame

expect(stream.state).to be == :closed

# There is nothing left to drain, so the connection is closed:
expect(client).to be(:closed?)
end

it "keeps peer-initiated streams open when receiving GOAWAY" do
stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)
server.read_frame

# The client's GOAWAY stream ID applies to streams initiated by the server, not the client's request stream:
client.send_goaway(0)
server.read_frame

expect(server.goaway_stream_id).to be == 0
expect(server.remote_stream_id).to be == 1
expect(server.streams.keys).to be == [1]
expect(server.streams[1].state).to be == :half_closed_remote
expect(server).to be(:goaway_received?)
expect(server).not.to be(:closed?)
end

let(:stream_class) do
Expand Down Expand Up @@ -442,6 +469,165 @@ def closed(error)

# The processed stream (id=1) should still be open:
expect(stream.state).not.to be == :closed

# Unprocessed streams are removed from the connection, so what remains is exactly what we are waiting for:
expect(client.streams.keys).to be == [1]
end

it "uses the lowest stream ID from successive GOAWAY frames" do
stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)

another_stream = client.create_stream do |connection, id|
stream_class.create(connection, id)
end
another_stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)

goaway = Protocol::HTTP2::GoawayFrame.new
goaway.pack(3, 0, "")
server.write_frame(goaway)
client.read_frame

expect(client.goaway_stream_id).to be == 3
expect(another_stream.state).not.to be == :closed

goaway = Protocol::HTTP2::GoawayFrame.new
goaway.pack(1, 0, "")
server.write_frame(goaway)
client.read_frame

expect(client.goaway_stream_id).to be == 1
expect(another_stream.state).to be == :closed
expect(another_stream.error).to be_a(Protocol::HTTP::RefusedError)
expect(client.streams.keys).to be == [1]

goaway = Protocol::HTTP2::GoawayFrame.new
goaway.pack(3, 0, "")
server.write_frame(goaway)
client.read_frame

# A later GOAWAY cannot raise the effective cutoff:
expect(client.goaway_stream_id).to be == 1
end

it "drains the streams which were accepted before a graceful GOAWAY" do
stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)

another_stream = client.create_stream do |connection, id|
stream_class.create(connection, id)
end
another_stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)

# Establish both request streams on the server:
server.read_frame
server.read_frame

# The server accepted both streams, but will not accept any more:
server.send_goaway(0)

client.read_frame

expect(client).to be(:goaway_received?)
expect(client).not.to be(:closed?)
expect(client.streams.keys).to be == [1, 3]

# Both streams still receive their response:
server.streams[1].send_headers(response_headers, Protocol::HTTP2::END_STREAM)
client.read_frame

expect(stream.state).to be == :closed
expect(client).not.to be(:closed?)
expect(client.streams.keys).to be == [3]

server.streams[3].send_headers(response_headers, Protocol::HTTP2::END_STREAM)
client.read_frame

expect(another_stream.state).to be == :closed
expect(another_stream.error).to be_nil

# The last accepted stream completed, so the connection is closed:
expect(client).to be(:closed?)
end

it "refuses to create new streams after a graceful GOAWAY" do
stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)
server.read_frame

server.send_goaway(0)
client.read_frame

expect(client).to be(:goaway_received?)
expect(client).not.to be(:closed?)

expect do
client.create_stream
end.to raise_exception(Protocol::HTTP2::ProtocolError, message: be =~ /Cannot create stream 3 after GOAWAY/)

# The stream the server accepted is unaffected:
expect(client.streams.keys).to be == [1]
end

it "still accepts streams initiated by the remote peer after GOAWAY" do
stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)
server.read_frame

server.send_goaway(0)
client.read_frame

expect(client).to be(:goaway_received?)
expect(client).not.to be(:closed?)

# `last_stream_id` only covers the streams we initiate, so the server can still push on the stream it accepted. Its frames must be processed like any other: dropping them would desynchronise the HPACK context shared by the whole connection.
promised_stream = server.streams[1].send_push_promise(request_headers)
client.read_frame

expect(client.streams.keys).to be == [1, promised_stream.id]

promised_stream.send_headers(response_headers, Protocol::HTTP2::END_STREAM)

expect(client.read_frame).to be_a(Protocol::HTTP2::HeadersFrame)
expect(client.streams[promised_stream.id].state).to be == :half_closed_local
end

it "decides the state of the connection even if a stream callback raises" do
raising_stream_class = Class.new(Protocol::HTTP2::Stream) do
def closed(error)
super

raise "Error in closed callback!" if error
end
end

another_stream = client.create_stream do |connection, id|
raising_stream_class.create(connection, id)
end
another_stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)

# The server did not process anything, so the stream is refused and its callback raises:
server.send_goaway(0)

expect do
client.read_frame
end.to raise_exception(RuntimeError, message: be =~ /Error in closed callback/)

expect(client).to be(:goaway_received?)
expect(client).to be(:closed?)
end

it "closes the connection immediately if there is nothing to drain" do
stream.send_headers(request_headers, Protocol::HTTP2::END_STREAM)
server.read_frame

# The server processes the stream before shutting down:
server.streams[1].send_headers(response_headers, Protocol::HTTP2::END_STREAM)
server.send_goaway(0)

client.read_frame
expect(stream.state).to be == :closed

client.read_frame

expect(client).to be(:goaway_received?)
expect(client).to be(:closed?)
end

it "closes all streams with RefusedError on GOAWAY with last_stream_id=0" do
Expand Down Expand Up @@ -499,7 +685,8 @@ def closed(error)

client.close

expect(client.remote_stream_id).to be == 1
expect(client.goaway_stream_id).to be == 1
expect(client.remote_stream_id).to be == 0
expect(client).to be(:closed?)
end

Expand Down
Loading