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
29 changes: 23 additions & 6 deletions lib/datastar/dispatcher.rb
Original file line number Diff line number Diff line change
Expand Up @@ -386,14 +386,33 @@ def wrap_socket(socket)
@compressor.wrap_socket(socket)
end

# Errors raised when writing to a socket the client has closed.
DISCONNECT_ERRORS = [IOError, Errno::EPIPE, Errno::ECONNRESET].freeze

# Whether an error means the client went away.
#
# Servers may wrap the underlying socket error in their own class
# (e.g. protocol-http1 >= 0.41 raises +Protocol::HTTP::RemoteError+
# from a +rescue Errno::EPIPE+), so the +cause+ chain is checked too.
#
# @param error [Exception]
# @return [Boolean]
def client_disconnect?(error)
while error
return true if DISCONNECT_ERRORS.any? { |klass| error.is_a?(klass) }

error = error.cause
end
false
end

# Handle errors caught during streaming
# @param error [Exception] the error that occurred
# @param socket [IO] the socket to pass to error handlers
def handle_streaming_error(error, socket)
case error
when IOError, Errno::EPIPE, Errno::ECONNRESET
if client_disconnect?(error)
@on_client_disconnect.each { |callable| callable.call(socket) }
when Exception
else
@on_error.each { |callable| callable.call(error) }
end
end
Expand All @@ -407,10 +426,8 @@ def handling_sync_errors(generator, socket, &)
yield

@on_server_disconnect.each { |callable| callable.call(generator) }
rescue IOError, Errno::EPIPE, Errno::ECONNRESET => e
@on_client_disconnect.each { |callable| callable.call(socket) }
rescue Exception => e
@on_error.each { |callable| callable.call(e) }
handle_streaming_error(e, socket)
end

# Parse signals from the request
Expand Down
49 changes: 49 additions & 0 deletions spec/dispatcher_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,18 @@ def wait_for_close(&)
end
end

# A server-side wrapper around the socket error, as raised by e.g.
# protocol-http1 >= 0.41 (Protocol::HTTP::RemoteError from a rescue Errno::EPIPE).
class WrappedDisconnectError < StandardError; end

class WrappingTestSocket < TestSocket
def <<(line)
super
rescue Errno::EPIPE
raise WrappedDisconnectError, 'Remote connection closed during flush!'
end
end

RSpec.describe Datastar::Dispatcher do
include DispatcherExamples

Expand Down Expand Up @@ -635,6 +647,43 @@ def self.render_in(view_context) = %(<div id="foo">\n<span>#{view_context}</span
expect(events).to eq([true, false])
end

specify '#on_client_disconnect when the server wraps the socket error' do
events = []
errors = []
dispatcher
.on_client_disconnect { |conn| events << :disconnect }
.on_error { |err| errors << err }

dispatcher.stream do |sse|
sse.patch_signals(foo: 'bar')
end
socket = WrappingTestSocket.new(open: false)

dispatcher.response.body.call(socket)
socket.wait_for_close
expect(events).to eq([:disconnect])
expect(errors).to eq([])
end

specify 'concurrent streams: #on_client_disconnect when the server wraps the socket error' do
dispatcher = Datastar.new(request:, response:, heartbeat: 0.001)
events = []
errors = []
dispatcher
.on_client_disconnect { |conn| events << :disconnect }
.on_error { |err| errors << err }

dispatcher.stream do |sse|
sleep 10
end
socket = WrappingTestSocket.new(open: false)

dispatcher.response.body.call(socket)
socket.wait_for_close
expect(events).to eq([:disconnect])
expect(errors).to eq([])
end

specify '#check_connection triggers #on_client_disconnect' do
events = []
dispatcher
Expand Down
Loading