Skip to content

Stream#send_data / send_headers overwrite @state after a write that yielded, resurrecting a stream the peer already reset (async-http then leaks the pool slot and the reactor never exits) #38

Description

@kiyot

Summary

Protocol::HTTP2::Stream#send_data (and #send_headers) read @state before writing the frame and assign the next state after the write returns:

def send_data(*arguments, **options)
  if @state == :open
    frame = write_data(*arguments, **options)   # may yield: write mutex, semaphore, socket back-pressure
    if frame.end_stream?
      @state = :half_closed_local               # unconditional
    end
  elsif @state == :half_closed_remote
    frame = write_data(*arguments, **options)
    if frame.end_stream?
      close!
    end
  ...

write_data → Connection#write_frame can yield to the reactor (io-stream's @writing mutex, the write semaphore in async-http, or EAGAIN on the socket). If the peer closes the stream while that write is parked (RST_STREAM, including NO_ERROR sent after a complete response as allowed by RFC 9113 §8.1, or GOAWAY), the reader task runs close!: @state = :closed, the stream is removed from connection.streams, and closed(error) fires. When the writer resumes, the stale branch assigns @state = :half_closed_local on top of :closed. The result is a "zombie" stream: removed from the connection, closed already fired, but closed? returns false.

Consequence in async-http

Async::HTTP::Protocol::HTTP2::Response#pool= runs when the request task resumes after response.wait, which can be after the reset was processed and after the parked write completed:

if @stream.closed?
  pool.release(@stream.connection)
else
  @stream.pool = pool   # expects `closed` to fire later, but it already did
end

On a zombie stream it takes the second branch, so the connection's usage count in Async::Pool::Controller is never decremented. When the reactor tries to exit, the pool gardener's ensure runs Controller#close → drain → acquire_existing_resource, which waits on @condition for a usage of zero that never comes. The gardener was made non-transient by its own cancel, so Scheduler#run never returns: the thread sits in select with idle CPU forever and Sync { ... } never returns.

Why this only started to hang recently

The stale assignment itself is old (it is the same in 0.26.0). Two changes removed the escape hatch:

With protocol-http2 ≤ 0.27 and async-http ≤ 0.103 the same race raised StreamError: Stream closed! from Response::Stream#wait, so Client#call's ensure released the pool slot: a failed request instead of a hang. That is also what we see in production against Firebase Cloud Messaging: with the older gems the race surfaced as Protocol::HTTP2::StreamError: Stream closed! a few times a month, and after upgrading the same workload started to hang instead (all responses received, Sync never returning).

Reproduction

The script below runs a plain Async::HTTP::Server and Async::HTTP::Client over h2c in one process, no TLS.

  • The server answers as soon as content-length bytes have arrived and then closes the request body, so it sends RST_STREAM(NO_ERROR) after the response (the "early response" case from RFC 9113 §8.1).
  • The client holds io-stream's write mutex for up to 50 ms after each flush, standing in for a socket write that blocks under load.
REQUESTS=100 WRITE_DELAY=0.05 ruby minimal_zombie_stream.rb

prints the number of zombie streams (removed from the connection but state != :closed), typically 90 or more out of 100.

LATE_POOL_ASSIGN=1 REQUESTS=100 WRITE_DELAY=0.05 ruby minimal_zombie_stream.rb

additionally delays the request task's pool= until after the parked write completed (an ordering that occurs by itself under load). Sync then never returns; the only live task is Async::Pool::Controller Gardener inside drain, and the script's watchdog reports the hang after 10 s.

Versions: ruby 4.0.7, async 2.39.0, async-http 0.105.0, protocol-http2 0.29.1, async-pool 0.12.0, io-stream 0.14.0. Reproduced on Linux (epoll) and macOS (kqueue). With protocol-http2 0.26.0 / async-http 0.95.1, the same race (reproduced with our production client under the same conditions) surfaces as StreamError: Stream closed! for the affected requests instead of a hang.

minimal_zombie_stream.rb
# frozen_string_literal: true

# Minimal reproduction: a client stream that is reset by the peer while its END_STREAM write is blocked
# comes back to life as `half_closed_local`, so async-http never releases the connection back to the pool,
# and the pool gardener's `drain` blocks forever when the reactor tries to exit.
#
#   REQUESTS=100 WRITE_DELAY=0.05 ruby minimal_zombie_stream.rb                      # counts zombie streams
#   LATE_POOL_ASSIGN=1 REQUESTS=100 WRITE_DELAY=0.05 ruby minimal_zombie_stream.rb   # hangs on reactor exit
#
# Tested with ruby 4.0.7, async 2.39.0, async-http 0.105.0, protocol-http2 0.29.1, async-pool 0.12.0, io-stream 0.14.0.

require "bundler/setup" if ENV["BUNDLE_GEMFILE"]
require "async"
require "async/barrier"
require "async/http/server"
require "async/http/client"
require "async/http/endpoint"
require "protocol/http/response"

$stdout.sync = true
PORT = Integer(ENV.fetch("REPRO_PORT", 9444))
REQUESTS = Integer(ENV.fetch("REQUESTS", 50))
WRITE_DELAY = Float(ENV.fetch("WRITE_DELAY", 0.02))

puts %w[async async-http protocol-http2 async-pool io-stream].map { "#{it}=#{Gem.loaded_specs[it].version}" }.join(" ")

# Client side only: keep the io-stream write mutex (@writing) for a moment after each flush. This stands in for a
# socket write that blocks (EAGAIN) under load, and gives the peer time to answer before END_STREAM is written.
IO::Stream::Buffered.prepend(Module.new do
  private def drain(buffer)
    result = super
    sleep(rand * WRITE_DELAY) if Thread.current.name == "client"
    result
  end
end)

# Optional: make the request task slow to assign the pool after `response.wait` returns, i.e. the request task is
# scheduled after the peer reset the stream AND after the blocked END_STREAM write completed. This is the ordering
# that leaks the pool slot (it happens by itself under load; here it is forced to make the hang deterministic).
if ENV["LATE_POOL_ASSIGN"] == "1"
  Async::HTTP::Client.prepend(Module.new do
    def make_response(request, connection, attempt)
      response = request.call(connection)
      sleep(WRITE_DELAY * 2)
      response.pool = @pool
      response
    end
  end)
end

# ---- server: answer as soon as the request body (content-length bytes) has arrived, then close the request body.
# The async-http server then sends RST_STREAM(NO_ERROR) once the response is complete (RFC 9113 §8.1). This matches
# what we observe from Firebase Cloud Messaging in production.
app = lambda do |request|
  length = request.body&.length || 0
  received = 0
  while received < length && (chunk = request.body.read)
    received += chunk.bytesize
  end
  request.body&.close
  # A little latency so that HEADERS, DATA, END_STREAM and RST_STREAM tend to land in the client's read buffer together.
  sleep 0.005
  Protocol::HTTP::Response[200, {}, ["ok"]]
end

endpoint = Async::HTTP::Endpoint.parse("http://localhost:#{PORT}", protocol: Async::HTTP::Protocol::HTTP2)
server = Thread.new do
  Thread.current.name = "server"
  Sync { Async::HTTP::Server.new(app, endpoint).run }
end
sleep 0.5

# ---- client
zombies = []
Thread.current.name = "client"
done = false
watchdog = Thread.new do
  sleep 10
  next if done
  puts "HANG: Sync did not return 10s after all #{REQUESTS} responses were received"
  Thread.list.find { it.name == "client" }&.then { |t| puts t.backtrace.first(6) }
  exit!(2)
end

Sync do
  client = Async::HTTP::Client.new(endpoint)
  barrier = Async::Barrier.new
  REQUESTS.times do |i|
    barrier.async do
      body = Protocol::HTTP::Body::Buffered.new(["x" * 1000])
      response = client.post("/#{i}", {}, body)
      response.read
      response.close
    end
  end
  barrier.wait

  puts "all #{REQUESTS} responses received"
  client.pool.resources.each do |connection, usage|
    puts "connection usage=#{usage} active streams in connection=#{connection.streams.size}"
  end
  # Zombie streams: already closed (and removed from the connection) but reporting a different state.
  ObjectSpace.each_object(Async::HTTP::Protocol::HTTP2::Response::Stream) do |stream|
    next if stream.connection.streams.key?(stream.id)
    zombies << stream if stream.state != :closed
  end
  puts "zombie streams (removed from connection, state != closed): #{zombies.size} -> #{zombies.map(&:state).tally}"
  puts "pool slots held by zombies will never be released; the reactor will hang on exit." if zombies.any?
ensure
  barrier&.stop
end
done = true
puts "Sync returned normally"

Suggested fix

Re-read @state after write_data / write_headers return and skip the transition if the stream was closed in the meantime, for example:

frame = write_data(*arguments, **options)
return frame if @state == :closed

in both branches of send_data and in the corresponding branches of send_headers.

Two related hardenings would make the failure mode less severe even if a similar race reappears: async-http's Response#pool= could release immediately whenever the stream is no longer in connection.streams instead of relying on a future closed callback, and Async::Pool::Controller#drain could avoid waiting indefinitely during reactor shutdown.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions