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
2 changes: 2 additions & 0 deletions gems/smithy-client/lib/smithy-client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -50,9 +50,11 @@
require_relative 'smithy-client/http/headers'
require_relative 'smithy-client/http/response'
require_relative 'smithy-client/http/request'
require_relative 'smithy-client/response_sink'
require_relative 'smithy-client/stream'
require_relative 'smithy-client/transport'
require_relative 'smithy-client/net_http/connection_pool'
require_relative 'smithy-client/net_http/exchange'
require_relative 'smithy-client/net_http/stream'
require_relative 'smithy-client/net_http/transport'
require_relative 'smithy-client/send_handler'
Expand Down
9 changes: 9 additions & 0 deletions gems/smithy-client/lib/smithy-client/http/response.rb
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,8 @@ def signal_done(options = {})
signal_data(options[:body])
signal_done
elsif options.empty?
return if @done

@body.rewind if @body.respond_to?(:rewind)
@done = true
emit(:done)
Expand All @@ -103,7 +105,13 @@ def signal_done(options = {})
end

# @param [StandardError] error
# TODO: errors raised from a :done listener are lost and the operation
# reports success (because @done is already set, this no-ops and never
# records @error). Fix when refactoring the listener API.
# https://github.com/smithy-lang/smithy-ruby/pull/363#discussion_r4149310679
def signal_error(error)
return if @done

@error = error
signal_done
end
Expand Down Expand Up @@ -153,6 +161,7 @@ def reset
@body.truncate(0)
@body.rewind
@error = nil
@done = nil
end

private
Expand Down
43 changes: 32 additions & 11 deletions gems/smithy-client/lib/smithy-client/net_http/connection_pool.rb
Original file line number Diff line number Diff line change
Expand Up @@ -134,33 +134,49 @@ def session_for(endpoint)
session = @pool[endpoint].shift if @pool.key?(endpoint)
end

pooled = false
begin
session ||= start_session(endpoint)
yield(session)
rescue StandardError
session&.finish
raise
else
@pool_mutex.synchronize do
@pool[endpoint] = [] unless @pool.key?(endpoint)
@pool[endpoint] << session
end
pooled = true
ensure
# +pooled+ is set true only after the session has been returned to
# the pool, so any exit before that point - a normal return, a raised
# error, or a non-StandardError unwind (+Interrupt+, +SystemExit+,
# +Timeout::Error+, +Thread#kill+) - leaves it false and finishes the
# session here. This must be +ensure+, not +rescue StandardError+, to
# cover the non-StandardError unwinds. Once pooled, this is a no-op.
session&.finish unless pooled
end
Comment thread
jterapin marked this conversation as resolved.
nil
end

# Finishes (closes) a session and guarantees it is not left in the pool.
# Serialized against {#session_for} check-in under +@pool_mutex+ so a
# Removal-from-pool and +finish+ happen together under +@pool_mutex+ so a
# cross-thread abort and a normal check-in cannot both own the same
# session: if already returned to the pool it is removed here before
# finishing. Never raises (see {ExtendedSession#finish}).
# session: if the session was already returned to the pool it is removed
# here before finishing. +finish+ (a blocking socket close) is held under
# the lock so this atomicity holds. Never raises (see
# {ExtendedSession#finish}).
# @param [Net::HTTPSession, nil] session
# @param [URI::HTTP, URI::HTTPS, nil] endpoint The endpoint the session
# was checked out for. When given, only that endpoint's list is
# searched instead of the whole pool.
# @return [nil]
def finish_session(session)
def finish_session(session, endpoint = nil)
return if session.nil?

@pool_mutex.synchronize do
@pool.each_value { |sessions| sessions.delete(session) }
if endpoint
key = remove_path_and_query(endpoint)
@pool[key]&.delete(session)
else
@pool.each_value { |sessions| sessions.delete(session) }
end
session.finish
end
nil
Expand Down Expand Up @@ -324,10 +340,15 @@ def request(...)
@last_used = Process.clock_gettime(Process::CLOCK_MONOTONIC, :millisecond)
end

# Attempts to close/finish the session without raising an error.
# Attempts to close/finish the session without raising. Both the
# +session_for+ +ensure+ teardown and +finish_session+ (the abort path,
# which contractually must not raise) rely on this: a socket close can
# surface any number of errors (a not-yet-started session, or an
# +Errno+ / +OpenSSL+ error from the underlying close), and none should
# escape teardown, so every +StandardError+ is swallowed.
def finish
@http.finish
rescue IOError
rescue StandardError
nil
end
end
Expand Down
Loading
Loading