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
28 changes: 17 additions & 11 deletions lib/ld-eventsource/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ def initialize(uri,
retry_enabled: true,
http_client_options: nil)
@uri = URI(uri)
@stopped = Concurrent::AtomicBoolean.new(false)
@stop_event = Concurrent::Event.new
@retry_enabled = retry_enabled

@headers = headers.clone
Expand Down Expand Up @@ -278,7 +278,7 @@ def query_params(&action)
# has no effect if called a second time.
#
def close
if @stopped.make_true
if @stop_event.try?
reset_http
end
end
Expand All @@ -289,7 +289,7 @@ def close
# @return [Boolean] true if the client has been shut down
#
def closed?
@stopped.value
@stop_event.set?
end

private
Expand All @@ -314,7 +314,7 @@ def default_logger
end

def run_stream
until @stopped.value
until @stop_event.set?
close_connection
begin
resp = connect
Expand All @@ -323,12 +323,16 @@ def run_stream
end
# There's a potential race if close was called in the middle of the previous line, i.e. after we
# connected but before @cxn was set. Checking the variable again is a bit clunky but avoids that.
return if @stopped.value
if @stop_event.set?
# close did not see this connection, so close it here.
reset_http
return
end
read_stream(resp) unless resp.nil?
rescue => e
# When we deliberately close the connection, it will usually trigger an exception. The exact type
# of exception depends on the specific Ruby runtime. But @stopped will always be set in this case.
if @stopped.value
# of exception depends on the specific Ruby runtime. But the client will always be closed in this case.
if @stop_event.set?
@logger.info { "Stream connection closed" }
else
log_and_dispatch_error(e, "Unexpected error from event source")
Expand All @@ -347,12 +351,14 @@ def run_stream
# Try to establish a streaming connection. Returns the StreamingHTTPConnection object if successful.
def connect
loop do
return if @stopped.value
return if @stop_event.set?
interval = @first_attempt ? 0 : @backoff.next_interval
@first_attempt = false
if interval > 0
@logger.info { "Will retry connection after #{'%.3f' % interval} seconds" }
sleep(interval)
# This wait ends early when close sets the event.
@stop_event.wait(interval)
return if @stop_event.set?
end
cxn = nil
begin
Expand Down Expand Up @@ -395,7 +401,7 @@ def read_stream(cxn)

chunks = Enumerator.new do |gen|
loop do
if @stopped.value
if @stop_event.set?
break
else
begin
Expand All @@ -417,7 +423,7 @@ def read_stream(cxn)
event_parser = Impl::EventParser.new(Impl::BufferedLineReader.lines_from(chunks), @last_id)

event_parser.items.each do |item|
return if @stopped.value
return if @stop_event.set?
case item
when StreamEvent
dispatch_event(item)
Expand Down
73 changes: 73 additions & 0 deletions spec/client_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -454,6 +454,79 @@ def send_stream_content(res, content, keep_open:)
end
end

describe "close during a reconnect wait" do
let(:long_reconnect_time) { 5 }

def start_failing_client(server, attempts)
server.setup_response("/") do |req,res|
attempts << Time.now
res.status = 500
res.body = "sorry"
res.keep_alive = false
end

error_sink = Queue.new
threads_before = Thread.list
client = subject.new(server.base_uri, reconnect_time: long_reconnect_time) do |c|
c.on_error { |error| error_sink << error }
end
worker = (Thread.list - threads_before).detect { |t| t.name == 'LD/SSEClient' }
error_sink.pop # the client now waits before it tries again
[client, worker]
end

it "stops the worker thread promptly" do
with_server do |server|
client, worker = start_failing_client(server, Queue.new)
with_client(client) do |c|
c.close
expect(worker.join(1)).not_to be_nil
end
end
end

it "does not send another request" do
with_server do |server|
attempts = Queue.new
client, worker = start_failing_client(server, attempts)
with_client(client) do |c|
c.close
worker.join(1)
expect(attempts.size).to eq 1
end
end
end
end

it "closes a connection that opens while the client is closed" do
server = TCPServer.new("127.0.0.1", 0)
server_result = Queue.new
server_thread = Thread.new do
socket = server.accept
while (line = socket.gets) && line != "\r\n"; end
socket.write("HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\n\r\n")
readable = IO.select([socket], nil, nil, 2)
eof = readable && socket.read_nonblock(1, exception: false).nil?
server_result << (eof ? :closed : :open)
socket.close
end

# query_params runs after the last stop check and before the request, so close cannot close this connection.
client = subject.new("http://127.0.0.1:#{server.addr[1]}") do |c|
c.query_params do
c.close
{}
end
end

with_client(client) do
expect(server_result.pop).to eq :closed
end
ensure
server_thread&.join(1)
server&.close
end

describe "HTTP method parameter" do
it "defaults to GET method" do
with_server do |server|
Expand Down
Loading