Skip to content
Closed
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
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,13 @@

### Fixed

- Retry a request **once** when the pooled keep-alive socket is already dead
(`unexpected eof`, `Connection reset by peer`, `TCPSocket:(closed)`).
Detected from the error message, not Faraday class (so cert failures and
connection refused are not retried). Applies to POST as well as GET; DNS
failures and real read timeouts are not retried. No backoff; the opt-in
`retry_config:` policy is unchanged.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

README.md still says retries are opt in, disabled clients make exactly one attempt, and writes are never retried. This change makes one retry automatic and retries writes. Update the public retry section to describe the new default.


- Default `idle_timeout` for `net_http_persistent` is now `25` seconds (was `55`).
GCP SSL-proxy load balancers close idle TLS around ~30s; keeping pooled
connections for 55s (AWS ALB idle minus 5s) caused
Expand Down
7 changes: 6 additions & 1 deletion lib/getstream_ruby/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -216,7 +216,7 @@ def request(method, path, data = {}, request_timeout: nil)
return make_multipart_request(method, path, query_params, data) if multipart_request?(data)

body_json = data.to_json
attempt = 0
attempt = stale_retries = 0

begin
started = monotonic_now
Expand All @@ -238,6 +238,11 @@ def request(method, path, data = {}, request_timeout: nil)
handle_response(response)
rescue Faraday::Error => e
error = TransportError.new("Request failed: #{e.message}", error_type: ErrorMapping.classify_faraday_error(e))
if stale_retries.zero? && ErrorMapping.stale_keep_alive?(e)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This retries every HTTP method after errors that can arrive while reading the response. The server may already have applied a POST, PATCH, or DELETE, so the retry can repeat a write. I reproduced a POST whose complete body reached the server twice. Restrict automatic retries to safe methods, or require a server enforced idempotency key for writes.

log_retry_attempt(method, path, error, stale_retries, started)
stale_retries += 1

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

stale_retries does not increment attempt, so this retry sits outside retry_config.max_attempts. With max_attempts: 3, I reproduced four requests. Count this retry in the same total attempt budget.

retry
end
if retry_eligible?(method, error, attempt)
wait_before_retry(method, path, error, attempt, started)
attempt += 1
Expand Down
36 changes: 35 additions & 1 deletion lib/getstream_ruby/error_mapping.rb
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,16 @@ module ErrorMapping

module_function

STALE_KEEP_ALIVE_PATTERN = /
connection\ reset\ by\ peer
|unexpected\ eof
|tcpsocket:\(closed\)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

TCPSocket:(closed) is also the normal message after a real Net::ReadTimeout because Net::HTTP closes the socket before Faraday formats the exception. A local real timeout sent two requests through this branch. Remove this match or detect a stale pooled socket without using the closed socket suffix.

|broken\ pipe
|connection\ is\ closed
|end\ of\ file\ reached
|tls_retry_write_records

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Four accepted messages have no tests: broken pipe, connection is closed, end of file reached, and tls_retry_write_records. The wrapped exception message branch is also only tested for DNS. Add table cases for every accepted message, including one stale message from wrapped_exception.

/ix.freeze

# Raises the appropriate `ApiError` / `RateLimitError` for a non-2xx
# `Faraday::Response`.
def raise_api_error(response)
Expand Down Expand Up @@ -111,7 +121,7 @@ def classify_faraday_error(error)
end

def classify_connection_failure(error)
wrapped = error.respond_to?(:wrapped_exception) ? error.wrapped_exception : nil
wrapped = wrapped_exception(error)
case wrapped
when SocketError
'dns_failure'
Expand All @@ -120,6 +130,30 @@ def classify_connection_failure(error)
end
end

# True when the failure looks like a reused keep-alive socket that the peer
# already closed. Match the error text, not Faraday class: SSLError and
# ConnectionFailed also cover cert failures and connection refused.
# DNS failures and real read timeouts are not stale-pool errors.
def stale_keep_alive?(error)
return false if error.nil?
return false if classify_faraday_error(error) == 'dns_failure'

stale_keep_alive_message?(error)
end

def wrapped_exception(error)
return nil unless error.respond_to?(:wrapped_exception)

error.wrapped_exception
end

def stale_keep_alive_message?(error)
texts = [error.message]
wrapped = wrapped_exception(error)
texts << wrapped.message if wrapped.respond_to?(:message)
texts.compact.any? { |text| text.match?(STALE_KEEP_ALIVE_PATTERN) }
end

def build_task_error(task_id, error_payload)
hash = if error_payload.respond_to?(:to_h)
error_payload.to_h
Expand Down
35 changes: 35 additions & 0 deletions spec/errors_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -316,6 +316,41 @@

end

describe 'ErrorMapping.stale_keep_alive?' do

it 'treats SSL EOF and connection reset as stale keep-alive' do

ssl = Faraday::SSLError.new('SSL_read: unexpected eof while reading')
reset = Faraday::ConnectionFailed.new('Connection reset by peer')
closed = Faraday::TimeoutError.new('Net::ReadTimeout with #<TCPSocket:(closed)>')
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(ssl)).to be(true)
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(reset)).to be(true)
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(closed)).to be(true)

end

it 'does not treat a real timeout or DNS failure as stale keep-alive' do

timeout = Faraday::TimeoutError.new('Net::ReadTimeout')
dns = Faraday::ConnectionFailed.new(SocketError.new('getaddrinfo failed'))
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(timeout)).to be(false)
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(dns)).to be(false)

end

it 'does not treat connection refused or cert failures as stale keep-alive' do

refused = Faraday::ConnectionFailed.new('Connection refused')
cert = Faraday::SSLError.new('certificate verify failed')
ssl_read = Faraday::SSLError.new('SSL_read: wrong version number')
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(refused)).to be(false)
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(cert)).to be(false)
expect(GetStreamRuby::ErrorMapping.stale_keep_alive?(ssl_read)).to be(false)

end

end

describe 'transport-layer failures' do

it 'wraps Faraday::ConnectionFailed as TransportError with cause preserved' do
Expand Down
177 changes: 172 additions & 5 deletions spec/retry_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,176 @@ def enabled(max_attempts: 3, max_backoff: 30.0)

end

describe 'stale keep-alive retry' do

it 'retries POST once on Connection reset by peer without retry_config' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::ConnectionFailed, 'Connection reset by peer' if calls == 1

[200, {}, '{"ok":true}']

end

end
client = build_client(stubs)
client.make_request(:post, '/x')
expect(calls).to eq(2)
expect(client).not_to have_received(:sleep)

end

it 'retries POST once on SSL_read unexpected eof' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::SSLError, 'SSL_read: unexpected eof while reading' if calls == 1

[200, {}, '{"ok":true}']

end

end
client = build_client(stubs)
client.make_request(:post, '/x')
expect(calls).to eq(2)

end

it 'retries POST once on ReadTimeout with a closed TCPSocket' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::TimeoutError, 'Net::ReadTimeout with #<TCPSocket:(closed)>' if calls == 1

[200, {}, '{"ok":true}']

end

end
client = build_client(stubs)
client.make_request(:post, '/x')
expect(calls).to eq(2)

end

it 'does not retry POST on a real read timeout' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::TimeoutError, 'Net::ReadTimeout'

end

end
client = build_client(stubs)
expect { client.make_request(:post, '/x') }.to raise_error(GetStreamRuby::TransportError)
expect(calls).to eq(1)

end

it 'does not retry POST on connection refused' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::ConnectionFailed, 'Connection refused'

end

end
client = build_client(stubs)
expect { client.make_request(:post, '/x') }.to raise_error(GetStreamRuby::TransportError)
expect(calls).to eq(1)

end

it 'does not retry POST on an SSL certificate failure' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::SSLError, 'certificate verify failed'

end

end
client = build_client(stubs)
expect { client.make_request(:post, '/x') }.to raise_error(GetStreamRuby::TransportError)
expect(calls).to eq(1)

end

it 'does not retry a DNS failure' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
# Faraday needs the SocketError as wrapped_exception for DNS classification.
raise Faraday::ConnectionFailed.new( # rubocop:disable Style/RaiseArgs
SocketError.new('getaddrinfo: nodename nor servname provided'),
)

end

end
client = build_client(stubs)
expect { client.make_request(:post, '/x') }.to raise_error(GetStreamRuby::TransportError) do |err|

expect(err.error_type).to eq('dns_failure')

end
expect(calls).to eq(1)

end

it 'retries a stale connection only once' do

calls = 0
stubs = Faraday::Adapter::Test::Stubs.new do |stub|

stub.post(%r{/x}) do

calls += 1
raise Faraday::ConnectionFailed, 'Connection reset by peer'

end

end
client = build_client(stubs)
expect { client.make_request(:post, '/x') }.to raise_error(GetStreamRuby::TransportError)
expect(calls).to eq(2)

end

end

it 'never retries an unrecoverable 429' do

calls = 0
Expand Down Expand Up @@ -202,11 +372,8 @@ def enabled(max_attempts: 3, max_backoff: 30.0)
end
client = build_client(stubs, retry_config: enabled)
file = GetStream::Generated::Models::FileUploadRequest.new(file: __FILE__)
expect do

client.send(:make_multipart_request, :post, '/upload', {}, file)

end.to raise_error(GetStreamRuby::RateLimitError)
expect { client.send(:make_multipart_request, :post, '/upload', {}, file) }
.to raise_error(GetStreamRuby::RateLimitError)
expect(calls).to eq(1)

end
Expand Down
Loading