diff --git a/gems/smithy-client/lib/smithy-client.rb b/gems/smithy-client/lib/smithy-client.rb index c9c016c0d..e9e83c779 100644 --- a/gems/smithy-client/lib/smithy-client.rb +++ b/gems/smithy-client/lib/smithy-client.rb @@ -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' diff --git a/gems/smithy-client/lib/smithy-client/http/response.rb b/gems/smithy-client/lib/smithy-client/http/response.rb index e354fcbe0..ef6113244 100644 --- a/gems/smithy-client/lib/smithy-client/http/response.rb +++ b/gems/smithy-client/lib/smithy-client/http/response.rb @@ -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) @@ -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 @@ -153,6 +161,7 @@ def reset @body.truncate(0) @body.rewind @error = nil + @done = nil end private diff --git a/gems/smithy-client/lib/smithy-client/net_http/connection_pool.rb b/gems/smithy-client/lib/smithy-client/net_http/connection_pool.rb index d0ac23f01..eaf06cd00 100644 --- a/gems/smithy-client/lib/smithy-client/net_http/connection_pool.rb +++ b/gems/smithy-client/lib/smithy-client/net_http/connection_pool.rb @@ -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 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 @@ -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 diff --git a/gems/smithy-client/lib/smithy-client/net_http/exchange.rb b/gems/smithy-client/lib/smithy-client/net_http/exchange.rb new file mode 100644 index 000000000..9c58a33ff --- /dev/null +++ b/gems/smithy-client/lib/smithy-client/net_http/exchange.rb @@ -0,0 +1,311 @@ +# frozen_string_literal: true + +require 'net/http' + +module Smithy + module Client + module NetHTTP + # The HTTP/1.1 request/response driving engine, backed by +Net::HTTP+. An + # +Exchange+ runs a single request/response and PUSHES the response into + # the sink supplied at construction (+sink.headers+, then +sink.data+ per + # body chunk, then one terminal). + # + # {Transport} owns it: {Transport#transmit} calls {#drive}, and + # {Transport#transmit_background} calls {#drive_background} and hands the + # +Exchange+ to a {Stream} handle. The verb is validated at construction, so + # an invalid verb raises {ArgumentError} before {#drive} touches the + # network. + # + # {#abort} may be called from another thread (the norm for an event stream + # driven in the background). State transitions are mutex-guarded so an abort + # cannot race completion or deliver to the sink after the abort (see + # {#deliver}), and closing the socket interrupts a blocked read on the + # driving thread. + # @api private + class Exchange + # Internal sentinel raised on the abort path so + # {ConnectionPool#session_for} finishes the checked-out socket instead of + # returning it to the pool. Rescued ahead of the +StandardError+ clause so + # an abort is never reported as a networking failure or surfaced to the + # sink. + # @api private + class InternalAbortSignal < StandardError; end + + # @param [ConnectionPool] pool The connection pool to check a session + # out of. + # @param [Http::Request] request + # @param [#headers, #data, #done, #error] sink The response sink this + # exchange pushes into when driven (see {ResponseSink}). + # @raise [ArgumentError] If the request has an invalid HTTP method + # (validated here, before any network I/O). + def initialize(pool, request, sink) + @pool = pool + @request = request + @sink = sink + @net_request = build_net_request(request) + @session = nil + @bytes_received = 0 + @status = nil + @headers = nil + @done = false + @aborted = false + # Guards the @done/@aborted/@session transitions so a cross-thread + # {#abort} is observed by the driving thread and cannot race normal + # completion. + @mutex = Mutex.new + end + + # Runs the exchange synchronously on the caller's thread, pushing the + # response into the sink to a single terminal. Networking failures are + # surfaced as +sink.error(NetworkingError)+, not raised. + # @return [void] + def drive + run(@net_request) + nil + end + + # Runs {#drive} on a background thread and returns immediately. Net::HTTP + # is blocking, so an OS thread is the concurrency mechanism (a reactor + # transport would use an async task). + # @return [Thread] + def drive_background + Thread.new { drive } + end + + # Aborts the exchange, discarding the underlying session rather than + # returning it to the pool. Idempotent and never raises (it runs on + # teardown paths); safe to call from another thread (see the class-level + # Concurrency note). + # @param [StandardError, nil] _error + # @return [void] + def abort(_error = nil) + session = nil + @mutex.synchronize do + return if @done || @aborted + + @aborted = true + session = @session + @session = nil + end + # Discard through the pool so this finish is serialized against a + # concurrent check-in: the pool removes the session if it was already + # returned, closing the window where abort could finish a just-pooled + # session. Passing the endpoint keeps the pool's removal targeted to + # that endpoint's list rather than scanning the whole pool. + # {ExtendedSession#finish} swallows socket-close errors, so this cannot + # raise on the teardown path. + @pool.finish_session(session, @request.endpoint) + nil + end + + private + + # Runs the full request within the pool's session block, pushing the + # response into the sink. Terminates the sink exactly once. + # @param [Net::HTTPRequest] net_request + def run(net_request) + # The success terminal is emitted HERE, outside #drive_exchange's rescue + # region: +@sink.done+ runs caller :done listeners, and a bug in one + # must propagate to the caller rather than be caught by the rescue and + # turned into a second, error terminal (which would also mask the caller + # bug as a transient NetworkingError and trigger a re-download). + @sink.done if drive_exchange(net_request) + nil + end + + # Drives the request within the pool's session block. Emits the abort and + # networking-failure terminals itself; the success terminal is emitted by + # {#run} (outside this method's rescue region) only when this returns true. + # @param [Net::HTTPRequest] net_request + # @return [Boolean] true on the success path (caller should emit + # +sink.done+); false when aborted or a failure terminal was delivered. + def drive_exchange(net_request) + @pool.session_for(@request.endpoint) do |session| + store_session(session) + # #abort may have raced ahead of checkout with a nil session (nothing + # to finish yet). Raise so #session_for finishes this session instead + # of pooling it, avoiding a leak. The abort is already recorded. + raise InternalAbortSignal if aborted? + + perform_exchange(session, net_request) + # perform_exchange can return NORMALLY even when an abort landed + # mid-body (delivery just stops in #deliver). #release_unless_aborted + # decides the session's fate under one lock: it returns false if + # aborted (raise so #session_for finishes the already-closed socket) + # or relinquishes the session for pooling otherwise. + raise InternalAbortSignal unless release_unless_aborted + end + true + rescue InternalAbortSignal + # Intentional abort path; session_for already discarded the session. + # Nothing to surface - an aborted exchange stays quiet. + mark_done + false + rescue StandardError => e + # A networking failure (the invalid-verb ArgumentError is validated at + # construction, before run). If an abort is concurrently in progress + # the socket close surfaces here as a read error; stay quiet in that + # case since the caller intentionally cancelled. + mark_done + @sink.error(NetworkingError.new(e)) unless aborted? + false + end + + # Issues the request within the pool's session block and pushes the + # response into the sink. Runs on the driving thread. + # @param [Net::HTTPSession] session + # @param [Net::HTTPRequest] net_request + def perform_exchange(session, net_request) + # {Patches} suppresses net-http < 0.7.0's default Content-Type via this + # flag; see there for the version detail and removal TODO. Clear it as + # soon as the request is sent (the response block runs after send), and + # keep the ensure as a backstop. + Thread.current[:net_http_skip_default_content_type] = true + session.request(net_request) do |net_response| + Thread.current[:net_http_skip_default_content_type] = nil + push_response(net_response) + end + ensure + Thread.current[:net_http_skip_default_content_type] = nil + end + + # Pushes one response into the sink. Every sink call goes through + # {#deliver}, so nothing is delivered after an abort. + # @param [Net::HTTPResponse] net_response + def push_response(net_response) + @status = net_response.code.to_i + @headers = extract_headers(net_response) + return unless deliver { @sink.headers(@status, @headers) } + + net_response.read_body do |chunk| + @bytes_received += chunk.bytesize + next if chunk.empty? + + break unless deliver { @sink.data(chunk) } + end + verify_content_length! + end + + # Runs +block+ (a single sink call) if not aborted, atomically with + # respect to {#abort}: the check and the call happen under @mutex, so an + # abort cannot slip between observing "not aborted" and delivering, and no + # headers/data reach the sink after an abort is recorded. Returns whether + # it delivered, so callers stop pushing once aborted. + # + # The sink call runs while @mutex is held, so a sink MUST NOT call back + # into this exchange's #abort (it would deadlock). #abort is driven by the + # consumer thread, not from inside the sink, so this holds. + # @return [Boolean] true if the sink call ran; false if aborted. + def deliver + @mutex.synchronize do + return false if @aborted + + yield + true + end + end + + # @return [Boolean] + def aborted? + @mutex.synchronize { @aborted } + end + + def store_session(session) + @mutex.synchronize { @session = session } + end + + # Atomically decides the session's fate at the end of a normal exchange: + # under one @mutex acquisition, returns false if an abort was recorded + # (leaving @session in place, since abort already took/closed it) so the + # caller raises {InternalAbortSignal}; otherwise marks the exchange done, + # relinquishes @session (so a later abort no-ops), and returns true so + # #session_for re-pools the live connection. + # @return [Boolean] true if the session was released for pooling; false if + # aborted. + def release_unless_aborted + @mutex.synchronize do + return false if @aborted + + @session = nil + @done = true + true + end + end + + # Marks the exchange done and relinquishes ownership of the session so a + # concurrent #abort no-ops (on @done) and, even if it had already read a + # stale @session, finds none to finish. On the error/abort paths that + # call this, #session_for has already finished the socket; nil-ing + # @session here makes "done => nothing left to finish" a hard invariant + # rather than relying on that ordering. + def mark_done + @mutex.synchronize do + @done = true + @session = nil + end + end + + # Constructs a +Net::HTTP+ request object from an {Http::Request}. + # @param [Http::Request] request + # @return [Net::HTTPRequest] + def build_net_request(request) + request_class = net_http_request_class(request) + req = request_class.new(request.endpoint.request_uri, net_headers_for(request)) + # Set the body stream when its size is unknown or greater than 0. + req.body_stream = request.body if !request.body.respond_to?(:size) || request.body.size.positive? + req + end + + # @param [Http::Request] request + # @raise [ArgumentError] If the HTTP method is not a valid verb. + # @return [Class] + def net_http_request_class(request) + ::Net::HTTP.const_get(request.http_method.capitalize) + rescue NameError + raise ArgumentError, "`#{request.http_method}` is not a valid http verb" + end + + # @param [Http::Request] request + # @return [Hash] + def net_headers_for(request) + # When Accept-Encoding is left unset, Net::HTTP injects a default + # value and enables decode_content, transparently decompressing the + # response. Setting it explicitly (to 'identity') opts out of that + # path so the raw body is delivered. + headers = { 'accept-encoding' => 'identity' } + request.headers.each_pair do |key, value| + headers[key] = value + end + headers + end + + # @param [Net::HTTPResponse] response + # @return [Hash] + def extract_headers(response) + response.to_hash.transform_values(&:first) + end + + # Detects short fixed-length bodies that Net::HTTP may otherwise tolerate + # because +ignore_eof+ defaults to +true+. Runs inside the pool session + # block so a truncated connection is finished rather than returned to the + # pool. + # @raise [IOError] When fewer bytes arrived than +Content-Length+. + # @return [void] + def verify_content_length! + return unless should_verify_bytes? + + bytes_expected = @headers['content-length'].to_i + return if bytes_expected == @bytes_received + + raise IOError, "http response body truncated, expected #{bytes_expected} " \ + "bytes, received #{@bytes_received} bytes" + end + + # @return [Boolean] + def should_verify_bytes? + @request.http_method != 'HEAD' && !@headers.nil? && @headers.key?('content-length') + end + end + end + end +end diff --git a/gems/smithy-client/lib/smithy-client/net_http/patches.rb b/gems/smithy-client/lib/smithy-client/net_http/patches.rb index 64ab28720..b0bcca793 100644 --- a/gems/smithy-client/lib/smithy-client/net_http/patches.rb +++ b/gems/smithy-client/lib/smithy-client/net_http/patches.rb @@ -17,7 +17,7 @@ def self.apply! # is set. net-http 0.7.0+ removed this entirely, so the patch # is only applied when the method exists. Unable to remove this # completely due to bundled net-http versions in Ruby 3.2-3.3. - # TODO: remove this patch, the Stream skip-flag that drives it, and its + # TODO: remove this patch, the Exchange skip-flag that drives it, and its # spec once the min supported Ruby ships net-http >= 0.7.0 (i.e. drops # Ruby 3.3/3.4). Keyed on the min Ruby bump so it is not forgotten. # See: https://github.com/ruby/net-http/pull/207 diff --git a/gems/smithy-client/lib/smithy-client/net_http/stream.rb b/gems/smithy-client/lib/smithy-client/net_http/stream.rb index d077c74d0..97c8ae232 100644 --- a/gems/smithy-client/lib/smithy-client/net_http/stream.rb +++ b/gems/smithy-client/lib/smithy-client/net_http/stream.rb @@ -1,339 +1,46 @@ # frozen_string_literal: true -require 'net/http' - module Smithy module Client module NetHTTP - # An HTTP/1.1 stream backed by +Net::HTTP+, implementing the stream - # contract consumed by {SendHandler}: +#response_headers+, +#each_chunk+, - # +#write+, +#close_write+, and +#abort+. - # - # Body delivery is pull-based on the caller's thread: a +Fiber+ keeps - # +Net::HTTP+'s +read_body+ block alive across separate {#response_headers} - # and {#each_chunk} calls, so there is no background reader thread. (On - # JRuby/TruffleRuby a +Fiber+ is thread-backed, so each stream uses one - # background thread; the pull model is unchanged.) + # The HTTP/1.1 {Client::Stream} control handle, returned by + # {Transport#transmit_background} for an event-stream operation. It is a + # thin OUTBOUND + CONTROL wrapper around a backgrounded {Exchange}; it holds + # no driving logic of its own (the {Exchange} owns the request/response + # engine and pushes the response into the sink on a background thread). # - # The fiber is created and must be resumed on the same thread. A different - # thread may call {#abort} to cancel; state transitions are mutex-guarded - # so an abort cannot race normal completion. - # - # HTTP/1.1 cannot write to a request after transmit, so {#write} and - # {#close_write} raise {NotSupportedError}. + # See {Client::Stream} for the contract; the per-method docs below cover + # {#abort} and the {#write}/{#close_write} raises. # @api private class Stream - # Raised when the received response body is shorter than the advertised - # +Content-Length+. This is an HTTP/1.1 wire concern: HTTP/2 detects - # truncation via frame accounting / END_STREAM, not Content-Length. - # @api private - class TruncatedBodyError < IOError - def initialize(bytes_expected, bytes_received) - msg = "http response body truncated, expected #{bytes_expected} " \ - "bytes, received #{bytes_received} bytes" - super(msg) - end - end - - # Internal sentinel used on the abort path: raised inside the driving - # fiber when an {#abort} was recorded before session checkout completed, - # so that {ConnectionPool#session_for} finishes the checked-out socket - # instead of returning it to the pool. Caught explicitly (not via the - # broad StandardError rescue) so it is never mistaken for a real - # networking failure, and never surfaced to the caller. Not a separate - # kind of cancellation from {#abort} - it is how an abort unwinds the - # fiber. Not part of the stream contract. - # @api private - class InternalAbortSignal < StandardError; end - - # @param [ConnectionPool] pool The connection pool to check a session - # out of. - # @param [Http::Request] request - def initialize(pool, request) - @pool = pool - @request = request - @status = nil - @headers = nil - @session = nil - @error = nil - @bytes_received = 0 - @done = false - @aborted = false - # Guards the @done/@aborted/@error/@session transitions so a cross-thread - # {#abort} is observed by the driving thread and cannot race normal - # completion. - @mutex = Mutex.new + # @param [Exchange] exchange The backgrounded request/response engine + # this handle controls. + def initialize(exchange) + @exchange = exchange end - # Sends the request and reads the response status and headers, leaving - # the response body ready to be consumed via {#each_chunk}. Called by - # {Transport#transmit}; not part of the public +Stream+ contract. - # - # @raise [ArgumentError] If the request has an invalid HTTP method. - # @raise [NetworkingError] If a networking error occurs while sending - # the request or reading the response headers. - # @return [self] - # @api private - def send_request - # Build (and validate) the Net::HTTP request before touching the - # network so an invalid verb raises ArgumentError without opening a - # connection or wrapping the error. - net_request = build_net_request(@request) - @fiber = Fiber.new { run(net_request) } - @fiber.resume - raise @error if @error - - self - end - - # @return [Array(Integer, Hash)] The response status code - # and headers. Available immediately after {#send_request} for HTTP/1.1. - def response_headers - [@status, @headers] - end - - # Yields raw response body chunks on the caller's thread until EOF or - # abort. If a fixed-length body ends before the advertised - # +Content-Length+, surfaces a {NetworkingError} after iteration. - # @yieldparam [String] chunk - # @raise [NetworkingError] If a networking error occurs while reading, or - # if the body is shorter than the advertised +Content-Length+. + # Cancels the in-progress exchange, discarding the underlying session + # rather than returning it to the pool. Delegates to the {Exchange}; safe + # to call from another thread, idempotent, and never raises. + # @param [StandardError, nil] error # @return [void] - def each_chunk(&) - return if finished? - - begin - stream_body(&) - rescue StandardError - # The consumer block (or the fiber resume) raised. Tear the - # connection down so the suspended fiber does not leak the - # checked-out socket, then re-raise the original error. - abort - raise - end - - # A cooperative abort stopped iteration early; skip surfacing a - # partial-read result as an error. - return if aborted? - - raise @error if @error - - nil + def abort(error = nil) + @exchange.abort(error) end - # Net::HTTP HTTP/1.1 does not support post-transmit request writes. + # Net::HTTP HTTP/1.1 does not support post-transmit request writes + # (bidirectional streaming). Output-only event streams never call this. # @raise [NotSupportedError] def write(_bytes) raise NotSupportedError, 'HTTP/1.1 does not support writing to a stream after transmit' end - # Net::HTTP HTTP/1.1 does not support post-transmit request writes. + # Net::HTTP HTTP/1.1 does not support post-transmit request writes + # (bidirectional streaming). Output-only event streams never call this. # @raise [NotSupportedError] def close_write raise NotSupportedError, 'HTTP/1.1 does not support writing to a stream after transmit' end - - # Aborts the stream and discards the underlying session rather than - # returning it to the pool. Safe to call from another thread, - # idempotent, and never raises. - # @param [StandardError, nil] error - # @return [void] - def abort(error = nil) - session = nil - @mutex.synchronize do - return if @done || @aborted - - @aborted = true - @error ||= error - session = @session - @session = nil - end - # Discard through the pool so this finish is serialized against a - # concurrent check-in: the pool removes the session if it was already - # returned, closing the window where abort could finish a just-pooled - # session. Errors are swallowed; abort runs on teardown paths and must - # not raise. - @pool.finish_session(session) - nil - end - - private - - # Body of the driving fiber. Runs the full request within the pool's - # session block, suspending (via Fiber.yield) after headers and after - # each body chunk so the caller drives reads. - # @param [Net::HTTPRequest] net_request - def run(net_request) - @pool.session_for(@request.endpoint) do |session| - store_session(session) - # #abort may have raced ahead of checkout with a nil session (nothing - # to finish yet). Raise so #session_for finishes this session instead - # of pooling it, avoiding a leak. The abort is already recorded. - raise InternalAbortSignal if aborted? - - perform_exchange(session, net_request) - # Relinquish the session before #session_for re-pools it: any abort - # after this no-ops, so check-in owns the session and abort cannot - # finish a pooled connection. - release_session - end - nil - rescue InternalAbortSignal - # Intentional abort path; session_for already discarded the session. - mark_done - nil - rescue StandardError => e - # A networking failure (the invalid-verb ArgumentError is validated in - # #send_request, before the fiber exists). Recorded even if an abort is - # concurrent; #each_chunk returns early when aborted, so it stays quiet. - mark_error(NetworkingError.new(e)) - nil - end - - # Sends the request and reads the response within the pool's session - # block, suspending after headers and after each body chunk so the caller - # drives reads. Runs inside the driving fiber. - # @param [Net::HTTPSession] session - # @param [Net::HTTPRequest] net_request - def perform_exchange(session, net_request) - # On net-http < 0.7.0, Net::HTTP applies a default Content-Type when a - # request has a body; {Patches} suppresses that via this flag. Set it - # here (inside the driving fiber) because Thread#[] is fiber-local and - # the request is sent from this fiber. No-op on net-http >= 0.7.0, - # which removed the behavior. - # TODO: remove with {Patches} once min Ruby ships net-http >= 0.7.0. - Thread.current[:net_http_skip_default_content_type] = true - session.request(net_request) do |net_response| - # The request has been sent by the time the response block runs, so - # the skip flag is no longer needed. - Thread.current[:net_http_skip_default_content_type] = nil - @status = net_response.code.to_i - @headers = extract_headers(net_response) - Fiber.yield # headers are ready; hand control back to #send_request - net_response.read_body do |chunk| - @bytes_received += chunk.bytesize - Fiber.yield(chunk) - end - # Verify while still inside the pool session block: a truncated - # (peer-closed) body raises here, so #session_for finishes the - # socket instead of returning it to the pool. - verify_content_length! - end - end - - # Drives the fiber, yielding non-empty body chunks. Stops promptly if - # another thread requested cancellation between chunks (cooperative - # abort). - # @yieldparam [String] chunk - def stream_body - while @fiber.alive? - break if aborted? - - chunk = @fiber.resume - break if chunk.nil? - next if chunk.empty? - - yield(chunk) - end - end - - # @return [Boolean] Whether the stream has completed or been aborted. - def finished? - @mutex.synchronize { @done || @aborted } - end - - # @return [Boolean] - def aborted? - @mutex.synchronize { @aborted } - end - - def store_session(session) - @mutex.synchronize { @session = session } - end - - # Marks the stream done and relinquishes ownership of the session so a - # concurrent #abort will no-op instead of finishing a session that is - # about to be (or has just been) returned to the pool. - def release_session - @mutex.synchronize do - @session = nil - @done = true - end - end - - def mark_done - @mutex.synchronize { @done = true } - end - - def mark_error(error) - @mutex.synchronize do - @error ||= error - @done = true - end - end - - # Constructs a +Net::HTTP+ request object from an {Http::Request}. - # @param [Http::Request] request - # @return [Net::HTTPRequest] - def build_net_request(request) - request_class = net_http_request_class(request) - req = request_class.new(request.endpoint.request_uri, net_headers_for(request)) - # On net-http < 0.7.0, Net::HTTP adds a default Content-Type when a - # body is present; {Patches} suppresses that during send (see - # #send_request). Set the body stream when its size is unknown or - # greater than 0. - req.body_stream = request.body if !request.body.respond_to?(:size) || request.body.size.positive? - req - end - - # @param [Http::Request] request - # @raise [ArgumentError] If the HTTP method is not a valid verb. - # @return [Class] - def net_http_request_class(request) - ::Net::HTTP.const_get(request.http_method.capitalize) - rescue NameError - raise ArgumentError, "`#{request.http_method}` is not a valid http verb" - end - - # @param [Http::Request] request - # @return [Hash] - def net_headers_for(request) - # When Accept-Encoding is left unset, Net::HTTP injects a default - # value and enables decode_content, transparently decompressing the - # response. Setting it explicitly (to 'identity') opts out of that - # path so the raw body is delivered. - headers = { 'accept-encoding' => 'identity' } - request.headers.each_pair do |key, value| - headers[key] = value - end - headers - end - - # @param [Net::HTTPResponse] response - # @return [Hash] - def extract_headers(response) - response.to_hash.transform_values(&:first) - end - - # Detects short fixed-length bodies that Net::HTTP may otherwise tolerate - # because +ignore_eof+ defaults to +true+. Runs inside the pool session - # block so a truncated connection is finished rather than returned to the - # pool. - # @raise [TruncatedBodyError] - # @return [void] - def verify_content_length! - return unless should_verify_bytes? - - bytes_expected = @headers['content-length'].to_i - return if bytes_expected == @bytes_received - - raise TruncatedBodyError.new(bytes_expected, @bytes_received) - end - - # @return [Boolean] - def should_verify_bytes? - @request.http_method != 'HEAD' && !@headers.nil? && @headers.key?('content-length') - end end end end diff --git a/gems/smithy-client/lib/smithy-client/net_http/transport.rb b/gems/smithy-client/lib/smithy-client/net_http/transport.rb index b4569a9b0..d402651c3 100644 --- a/gems/smithy-client/lib/smithy-client/net_http/transport.rb +++ b/gems/smithy-client/lib/smithy-client/net_http/transport.rb @@ -3,6 +3,7 @@ require 'openssl' require_relative 'connection_pool' +require_relative 'exchange' require_relative 'stream' module Smithy @@ -85,19 +86,46 @@ def initialize(options = {}) @pool = ConnectionPool.for(pool_options) end - # Sends the request (headers + body) synchronously and returns a - # {Stream} whose response status and headers are already available. - # Response body chunks are read on the caller's thread via - # {Stream#each_chunk}. + # The bridge queue for event streaming, matching this transport's + # concurrency model. Net::HTTP drives on a background thread, so a + # thread-blocking +SizedQueue+ is correct: the driving thread pushes and + # the consumer thread pops. Supplied to the event stream layer's bridge + # so the transport stays fully push and the bridge is concurrency- + # appropriate (see {Transport}). + # @return [SizedQueue] + def event_queue + SizedQueue.new(64) + end + + # Drives an {Exchange} inline (see {Transport#transmit} for the + # contract): sends and pushes the response into +sink+ to its terminal on + # the caller's thread, then returns nothing. Net::HTTP is synchronous, so + # "inline" is a plain blocking call here. + # @param [Http::Request] request + # @param [#headers, #data, #done, #error] sink The inbound response sink + # (see {ResponseSink}). + # @return [void] + # @raise [ArgumentError] If the request has an invalid HTTP method + # (validated at {Exchange} construction, before any network I/O). + def transmit(request, sink) + Exchange.new(@pool, request, sink).drive + nil + end + + # Drives an {Exchange} on a background thread (Net::HTTP is blocking, so + # the concurrency mechanism is an OS thread) and returns its {Stream} + # handle immediately (see {Transport#transmit_background} for the + # contract). # @param [Http::Request] request + # @param [#headers, #data, #done, #error] sink The inbound response sink + # (see {ResponseSink}). # @return [Stream] - # @raise [ArgumentError] If the request has an invalid HTTP method. - # @raise [NetworkingError] If a networking error occurs while sending - # the request or reading the response headers. - def transmit(request) - stream = Stream.new(@pool, request) - stream.send_request - stream + # @raise [ArgumentError] If the request has an invalid HTTP method + # (validated at {Exchange} construction, before spawning). + def transmit_background(request, sink) + exchange = Exchange.new(@pool, request, sink) + exchange.drive_background + Stream.new(exchange) end private diff --git a/gems/smithy-client/lib/smithy-client/plugins/transport.rb b/gems/smithy-client/lib/smithy-client/plugins/transport.rb index 64c1022b9..048be2164 100644 --- a/gems/smithy-client/lib/smithy-client/plugins/transport.rb +++ b/gems/smithy-client/lib/smithy-client/plugins/transport.rb @@ -16,8 +16,7 @@ module Plugins # These client options cover the shared transport settings exposed on # +Client.new(...)+. Transport-specific options remain adapter-specific and # must be configured on the transport instance itself (see - # {Smithy::Client::NetHTTP::Transport}). If a caller supplies a transport - # via +:transport+, that instance is used as-is. + # {Smithy::Client::NetHTTP::Transport}). # @api private class Transport < Plugin # The common client transport options forwarded to the default @@ -121,8 +120,9 @@ class Transport < Plugin The transport used to send requests. Defaults to an HTTP/1.1 transport based on Net::HTTP ({Smithy::Client::NetHTTP::Transport}), constructed with the resolved common client transport options. Supply a custom object responding to - `#transmit(request)` (returning a stream) to swap the transport, or a - directly-constructed `NetHTTP::Transport` to set Net::HTTP-specific knobs. A + `#transmit(request, sink)` and `#transmit_background(request, sink)` (both push + the response into the sink; see {Smithy::Client::Transport}) to swap the transport, + or a directly-constructed `NetHTTP::Transport` to set Net::HTTP-specific knobs. A caller-supplied transport instance is used as-is. DOCS Client::NetHTTP::Transport.new( @@ -134,10 +134,12 @@ class Transport < Plugin # Validates a customer-supplied +:transport+ before the config is built, # so the default transport is not eagerly constructed just to check it. - # Validates against {Client::Transport::REQUIRED_METHODS} by duck typing, - # so any conforming object is accepted without requiring a particular - # base class or module. The stream a transport returns has its own - # (send-time) contract that cannot be checked here. + # Checks {Client::Transport::REQUIRED_METHODS} by duck typing (no base + # class or module required). This is only a shape check: a single-mode + # transport still answers both methods and raises {Client::NotSupportedError} + # from the one it does not serve (see {Client::Transport}), which + # +respond_to?+ cannot detect - so a mode mismatch, and the returned + # stream's own contract, surface at send time rather than here. # @param [Class] _client_class # @param [Hash] options # @raise [ArgumentError] If a supplied transport does not answer the diff --git a/gems/smithy-client/lib/smithy-client/response_sink.rb b/gems/smithy-client/lib/smithy-client/response_sink.rb new file mode 100644 index 000000000..e647a3d14 --- /dev/null +++ b/gems/smithy-client/lib/smithy-client/response_sink.rb @@ -0,0 +1,48 @@ +# frozen_string_literal: true + +module Smithy + module Client + # Adapts an {Http::Response} to the sink side of the {Transport} contract: + # it forwards the sink calls (+#headers+/+#data+/+#done+/+#error+) onto the + # response's push-based +signal_*+ methods, so the transport never needs to + # know about {Http::Response} directly. See {Transport} for the calling + # order and terminal guarantees. + # @api private + class ResponseSink + # @param [Http::Response] response + def initialize(response) + @response = response + end + + # @param [Integer] status_code + # @param [Hash] headers + # @return [void] + def headers(status_code, headers) + @response.signal_headers(status_code, headers) + end + + # @param [String] chunk + # @return [void] + def data(chunk) + @response.signal_data(chunk) + end + + # Signals successful completion of the response. Terminal: at most one + # terminal ({#done} or {#error}) takes effect per response lifecycle. + # @return [void] + def done + @response.signal_done + end + + # Signals that the exchange failed. Terminal, in place of {#done}: if a + # terminal has already fired, this is a no-op (see {Http::Response}). + # @param [Exception] error The failure cause (a {NetworkingError} for a + # transport networking failure, or an {ArgumentError} for an invalid + # request). + # @return [void] + def error(error) + @response.signal_error(error) + end + end + end +end diff --git a/gems/smithy-client/lib/smithy-client/send_handler.rb b/gems/smithy-client/lib/smithy-client/send_handler.rb index 2d8af1b39..8428426ed 100644 --- a/gems/smithy-client/lib/smithy-client/send_handler.rb +++ b/gems/smithy-client/lib/smithy-client/send_handler.rb @@ -2,23 +2,17 @@ module Smithy module Client - # The +:send+ step handler for {Smithy::Client}. Sends the request using the - # configured transport (+config.transport+) and drives the resulting stream - # into the context's {Http::Response}. + # The +:send+ step handler for {Smithy::Client}. It supplies a + # {ResponseSink} over the context's {Http::Response} and hands it to the + # configured transport (+config.transport+), which pushes the response into + # it - response shaping lives in the {Transport}, not here. # - # This handler is adapter-independent but contract-shaping: it depends on no - # concrete transport, but requires the stream from +#transmit+ to fit a - # staged, pull-based model: - # - # * +#transmit+ returns before the body is consumed, - # * status/headers are a distinct phase via +#response_headers+, - # * the body is pulled in order via +#each_chunk+, and - # * +#abort+ cancels during the exchange. - # - # A push/event-style transport must adapt itself to this lifecycle. - # Protocol-specific concerns (blocking, pooling, truncation detection) live - # in the transport. This handler decides only *when* to block and bridges - # the pulled bytes onto the push-based {Http::Response}. + # The handler only picks WHICH transport method to call, by operation mode: + # {Transport#transmit} for a non-event-stream operation (drives to completion + # synchronously; nothing to store), or {Transport#transmit_background} for an + # event stream (+context[:event_stream]+; returns a {Stream} handle stored on + # the context for the event stream layer). See the methods below for the + # per-mode error/teardown handling. # @api private class SendHandler < Handler # @param [HandlerContext] context @@ -36,38 +30,41 @@ def call(context) # @param [HandlerContext] context # @return [void] def transmit(transport, req, resp, context) - stream = nil - stream = transport.transmit(req) - context[:stream] = stream + # The sink is the inbound destination in all modes (an event stream feeds + # inbound events into it too). + sink = ResponseSink.new(resp) + # TODO: nothing sets context[:event_stream] yet; the event stream layer will. + if context[:event_stream] + # Deliberately outside drive_non_event_stream's rescues: the exchange runs + # on a background thread, so failures surface via sink.error there, and + # teardown belongs to the event stream layer. + context[:stream] = transport.transmit_background(req, sink) + return + end - # Blocking is a handler-stack decision, not a transport concern. A - # duplex event stream would deadlock if blocked here (the server waits - # for input events), so the event stream layer drives it instead. All - # other operations resolve here so retry/error/parse handlers can run. - resolve_response(stream, resp) unless context[:duplex_stream] - rescue ArgumentError => e - # Invalid verb, ArgumentError is a StandardError. Not retryable. - resp.signal_error(e) - rescue StandardError => e - resp.signal_error(e.is_a?(NetworkingError) ? e : NetworkingError.new(e)) - ensure - # Guarantee the connection is released. On the normal path #abort is a - # no-op; if an error escaped before the body was consumed, #abort - # finishes the socket so it is not leaked. Duplex streams are owned and - # closed by the event stream layer. - stream.abort unless stream.nil? || context[:duplex_stream] + drive_non_event_stream(transport, req, resp, sink) end - # Resolves the response by reading headers and draining the body into the - # push-based {Http::Response}. - # @param [#response_headers, #each_chunk] stream - # @param [Http::Response] resp + # Non-event-stream send: transmit drives the response into the sink to its + # terminal synchronously and returns nothing (the transport owns teardown, + # so there is no handle to store or abort). Errors are caught and signaled + # onto the response so the error/retry handlers can run. # @return [void] - def resolve_response(stream, resp) - status, headers = stream.response_headers - resp.signal_headers(status, headers) - stream.each_chunk { |chunk| resp.signal_data(chunk) } - resp.signal_done + def drive_non_event_stream(transport, req, resp, sink) + transport.transmit(req, sink) + rescue ArgumentError, NotSupportedError => e + # Not retryable, and must be rescued ahead of the clause below: + # - ArgumentError (invalid verb) is raised before any network I/O; + # - NotSupportedError (a single-mode transport served only event streams + # and raised from #transmit) is a StandardError, so the clause below + # would otherwise wrap it into a transient NetworkingError and retry it + # with backoff. A mode mismatch is a caller/config error, not a + # networking failure. Signal both as-is. + resp.signal_error(e) + rescue StandardError => e + # Defensive: the transport should surface networking failures via + # sink.error while driving, so reaching here means something escaped. + resp.signal_error(e.is_a?(NetworkingError) ? e : NetworkingError.new(e)) end end end diff --git a/gems/smithy-client/lib/smithy-client/stream.rb b/gems/smithy-client/lib/smithy-client/stream.rb index 43633dd73..cf6e73ad0 100644 --- a/gems/smithy-client/lib/smithy-client/stream.rb +++ b/gems/smithy-client/lib/smithy-client/stream.rb @@ -2,71 +2,60 @@ module Smithy module Client - # The stream contract: the live handle returned by {Transport#transmit} and - # consumed by {SendHandler}. {NetHTTP::Stream} is the built-in HTTP/1.1 - # implementation. - # - # A stream is any object that answers {REQUIRED_METHODS}; it need not include - # this module. The module documents the contract and defines the - # compatibility surface. Conformance can be exercised with the shared - # compliance tests (+spec/support/stream_contract.rb+). - # - # ## Staged, pull-based model - # - # Adapter-independent but contract-shaping: an implementation must present - # its work as these ordered stages, however it fetches data internally. - # - # 1. **Headers.** After {Transport#transmit} returns, +#response_headers+ - # yields the status and headers without further input from the caller. - # 2. **Body.** The caller pulls the body in order via +#each_chunk+ until it - # is complete or the stream is aborted. - # 3. **Terminal.** Iteration ends once: normally, or by surfacing a - # {NetworkingError}. - # - # A live handle also supports +#abort+ for cancellation, and - # +#write+/+#close_write+ for transports that allow post-transmit writes - # (others raise {NotSupportedError}). + # The +Stream+ contract: the control handle returned by + # {Transport#transmit_background} for an EVENT-STREAM operation (it is the + # only path that produces a handle - a non-event-stream operation is driven + # to completion inside {Transport#transmit}, which returns nothing). + # {NetHTTP::Stream} is the built-in HTTP/1.1 implementation. + # + # Inbound response data does not flow back through the +Stream+: the transport + # pushes it into the sink supplied at {Transport#transmit_background} (see + # {ResponseSink}) as events arrive on the transport's own concurrency + # mechanism. The +Stream+ is the OUTBOUND + CONTROL side only - it cancels the + # exchange (+#abort+) and, for bidirectional-capable transports, writes + # request-side bytes after transmit (+#write+/+#close_write+). + # + # A +Stream+ answers {REQUIRED_FOR_OUTPUT_EVENT_STREAM_OPS} for an output-only + # event stream and, to also serve a BIDIRECTIONAL one, the write side in + # {REQUIRED_FOR_BIDI_EVENT_STREAM_OPS}. It need not include this module; + # conformance can be exercised with the shared compliance tests + # (+spec/support/stream_contract.rb+). # # ## Concurrency # - # +#each_chunk+ runs on the caller's thread. +#abort+ MAY be called from - # another thread; implementations must make the abort/completion transition - # race-safe and must not deliver chunks after an abort. - # - # ## Required methods - # - # +#response_headers+ - # * Returns +[Integer, Hash]+ - the response status code and - # headers. Available immediately after the transport returns the stream. + # The exchange runs on the transport's concurrency mechanism, so +#abort+ MAY + # be called from a thread other than the one driving the exchange. + # Implementations must make the abort/completion transition race-safe and must + # not deliver to the sink after an abort. +#abort+ must also be idempotent and + # never raise, since it runs on teardown paths. # - # +#each_chunk { |chunk| ... }+ - # * Yields raw response body chunks (String) in order, on the caller's - # thread, until the body is complete or the stream is aborted. - # * Raises {NetworkingError} on a networking error while reading, or if the - # body ends before its advertised length. - # * Returns +void+. - # - # +#write(bytes)+ - # * Writes +bytes+ (String) to the request after transmit. Only for - # transports that allow post-transmit writes; others raise - # {NotSupportedError}. - # - # +#close_write+ - # * Signals the end of the request body for a writable stream. Only for - # transports that allow post-transmit writes; others raise - # {NotSupportedError}. + # ## Output-only event streams: {REQUIRED_FOR_OUTPUT_EVENT_STREAM_OPS} # # +#abort(error = nil)+ # * Cancels the in-progress exchange and releases the underlying resource # without returning it for reuse. +error+ [{StandardError}, nil] optional - # cause. Must be safe to call from a thread other than the one driving - # +#each_chunk+, idempotent, and must never raise (it runs on teardown - # paths). Returns +void+. + # cause. Cross-thread-safe, idempotent, never raises (see Concurrency). + # Returns +void+. + # + # ## Bidirectional event streams: {REQUIRED_FOR_BIDI_EVENT_STREAM_OPS} + # + # In addition to +#abort+, a bidirectional stream answers the outbound write + # side. Both raise {NotSupportedError} on a transport that supports + # output-only but not bidirectional (e.g. HTTP/1.1): + # + # * +#write(bytes)+ - writes +bytes+ (String) to the request after transmit + # (outbound events). + # * +#close_write+ - signals the end of the request body (outbound). module Stream - # The messages every stream must answer. +#write+ and +#close_write+ are - # part of the surface but may raise {NotSupportedError} for transports that - # do not allow post-transmit writes. - REQUIRED_METHODS = %i[response_headers each_chunk write close_write abort].freeze + # Methods sufficient for an OUTPUT-ONLY event stream (server streams, + # client already sent its single request): cancellation only. + REQUIRED_FOR_OUTPUT_EVENT_STREAM_OPS = %i[abort].freeze + + # Methods sufficient for a BIDIRECTIONAL event stream: cancellation plus the + # outbound write side. The event stream layer checks for these and fails + # fast when a bidirectional operation is invoked against a transport that + # cannot serve it. + REQUIRED_FOR_BIDI_EVENT_STREAM_OPS = %i[abort write close_write].freeze end end end diff --git a/gems/smithy-client/lib/smithy-client/transport.rb b/gems/smithy-client/lib/smithy-client/transport.rb index e5a74c090..deb41cca4 100644 --- a/gems/smithy-client/lib/smithy-client/transport.rb +++ b/gems/smithy-client/lib/smithy-client/transport.rb @@ -2,50 +2,77 @@ module Smithy module Client - # The transport contract: sends a request and returns a {Stream} the handler - # stack consumes. This is the swap point exposed as the +:transport+ client - # option (see {Plugins::Transport}); {NetHTTP::Transport} is the built-in - # HTTP/1.1 implementation. - # - # A transport is any object that answers {REQUIRED_METHODS}; it need not - # include this module. The module documents the contract and defines the - # compatibility surface used to validate a caller-supplied transport. - # Conformance can be exercised with the shared compliance tests + # The transport contract: sends a request and pushes the response into a + # caller-supplied sink. This is the swap point exposed as the +:transport+ + # client option (see {Plugins::Transport}); {NetHTTP::Transport} is the + # built-in HTTP/1.1 implementation. A transport is any object answering + # {REQUIRED_METHODS} (it need not include this module); conformance can be + # exercised with the shared compliance tests # (+spec/support/transport_contract.rb+). # - # ## Contract + # ## Response shaping is the transport's job # - # A transport MUST: + # The transport owns the response lifecycle and PUSHES it into the sink (see + # {ResponseSink}) rather than returning a pull-stream for the caller to walk. + # It calls the sink's methods in order: # - # * expose +#transmit+, accepting an {Http::Request} and returning an object - # satisfying the {Stream} contract; - # * send the request before returning, so the stream's - # {Stream#response_headers} are available without further caller input; - # * return before the body is consumed (body delivery is pull-based, via - # {Stream#each_chunk}); - # * raise {NetworkingError} for networking failures while sending or reading - # headers, and {ArgumentError} for a malformed request (e.g. an invalid - # HTTP method) without opening a connection. + # * +sink.headers(status, headers)+ once, first; + # * +sink.data(chunk)+ zero or more times, in order; + # * exactly one terminal: +sink.done+ (success) or +sink.error(e)+ (failure). + # + # An ABORTED exchange is the one exception: it delivers NO terminal. Once an + # abort is recorded the transport stops calling the sink entirely, so a + # cancelled exchange ends silently and the party that called +abort+ is the + # only one that knows it ended. A caller must not wait on a terminal to + # observe cancellation. + # + # Networking failures are surfaced as +sink.error(NetworkingError)+, not + # raised. An invalid HTTP method is the exception: it raises {ArgumentError} + # synchronously (before opening a connection or spawning background work, and + # without touching the sink), since it is a caller error, not a transport + # failure. + # + # ## Two send methods, one per operation mode + # + # The operation mode - not the wire protocol - decides which method the SDK + # calls. Both take a +request+ ({Http::Request}) and a +sink+ (see above): + # + # * +#transmit(request, sink)+ serves a NON-EVENT-STREAM operation (plain + # request/response and byte streaming). It DRIVES the response into the sink + # to its terminal SYNCHRONOUSLY, then returns nothing. There is no handle to + # abort: the exchange has completed by the time it returns, and the + # transport owns teardown (releasing or discarding the connection itself). + # + # * +#transmit_background(request, sink)+ serves an EVENT-STREAM operation + # (output-only or bidirectional). It starts the exchange CONCURRENTLY - the + # transport owns the concurrency mechanism (a background thread for + # {NetHTTP::Transport}, an async task for an async transport) - and returns + # a {Stream} handle IMMEDIATELY, before the response arrives. The sink is + # fed asynchronously as inbound events arrive; the event stream layer + # consumes them, writes outbound via the handle (for bidirectional streams), + # and owns the handle's lifetime and teardown. + # + # (INPUT-ONLY streams - client-streaming with a unary response - are out of + # scope for this contract and not served by either method.) + # + # ## Serving only one mode + # + # {REQUIRED_METHODS} lists both methods. A transport that serves only one + # mode implements both and raises {NotSupportedError} from the mode it does + # not serve, mirroring {Stream#write}/{Stream#close_write} on a transport + # that cannot do bidirectional streaming. # # A transport MAY own connection management and protocol-specific concerns - # (pooling, blocking strategy, truncation detection); these are not part of - # the contract. - # - # ## Required methods - # - # +#transmit(request)+ - # * +request+ [{Http::Request}] the request to send. - # * Returns an object satisfying the {Stream} contract, with its response - # status and headers already available (for request/response operations); - # the body is read later, on the caller's thread, via {Stream#each_chunk}. - # * Raises {ArgumentError} if the request has an invalid HTTP method (without - # opening a connection). - # * Raises {NetworkingError} if a networking error occurs while sending the - # request or reading the response headers. + # (pooling, concurrency mechanism, truncation detection); these are not part + # of the contract. module Transport - # The messages every transport must answer. This is the compatibility - # surface checked for a caller-supplied +:transport+. - REQUIRED_METHODS = %i[transmit].freeze + # The messages every transport must answer to be a conforming +:transport+ + # (checked by {Plugins::Transport}). Both modes are required; a transport + # may opt out of one by raising {NotSupportedError} from it. + # + # TODO: fail fast in the event stream layer when an event-stream op hits a + # transport whose +#transmit_background+ raises {NotSupportedError}. + REQUIRED_METHODS = %i[transmit transmit_background].freeze end end end diff --git a/gems/smithy-client/spec/smithy-client/http/response_spec.rb b/gems/smithy-client/spec/smithy-client/http/response_spec.rb index f619d2b49..22f3d59f1 100644 --- a/gems/smithy-client/spec/smithy-client/http/response_spec.rb +++ b/gems/smithy-client/spec/smithy-client/http/response_spec.rb @@ -82,6 +82,15 @@ module Http subject.signal_done expect(subject.body.read).to eq('second response body') end + + it 're-arms the done terminal so a re-driven response can emit :done again' do + count = 0 + subject.on_done { count += 1 } + subject.signal_done + subject.reset + subject.signal_done + expect(count).to eq(2) + end end describe '#signal_headers' do @@ -123,6 +132,31 @@ module Http expect(done).to be(true) end + it 'emits :done at most once' do + count = 0 + subject.on_done { count += 1 } + subject.signal_done + subject.signal_done + expect(count).to eq(1) + end + + it 'does not re-emit :done via a subsequent signal_error' do + count = 0 + subject.on_done { count += 1 } + subject.signal_done + subject.signal_error(StandardError.new('late')) + expect(count).to eq(1) + end + + # Documents current behavior, not a guarantee: a late signal_error + # after a successful terminal is dropped. This is the same no-op that + # loses :done listener errors - see the TODO on Http::Response#signal_error. + it 'does not set an error after a successful terminal has fired' do + subject.signal_done + subject.signal_error(StandardError.new('late')) + expect(subject.error).to be_nil + end + it 'rewinds the body' do body = StringIO.new('data') subject.body = body diff --git a/gems/smithy-client/spec/smithy-client/net_http/connection_pool_spec.rb b/gems/smithy-client/spec/smithy-client/net_http/connection_pool_spec.rb index b7d8eca50..90b088de1 100644 --- a/gems/smithy-client/spec/smithy-client/net_http/connection_pool_spec.rb +++ b/gems/smithy-client/spec/smithy-client/net_http/connection_pool_spec.rb @@ -55,6 +55,32 @@ module NetHTTP end expect(sessions).to eq([session, session]) end + + it 'finishes the session and does not pool it when the block raises a StandardError' do + session = double('Net::HTTPSession').as_null_object + pool = ConnectionPool.for({}) + allow(pool).to receive(:start_session).and_return(session) + expect(session).to receive(:finish) + expect do + pool.session_for(URI.parse(endpoint)) { raise 'boom' } + end.to raise_error('boom') + expect(pool.size).to eq(0) + end + + it 'finishes the session on a non-StandardError unwind so the socket is not leaked' do + # The teardown must be +ensure+-based, not +rescue StandardError+: + # Interrupt/SystemExit/Timeout::Error/Thread#kill are not + # StandardError, and previously escaped without finishing the + # checked-out session, leaking the socket. + session = double('Net::HTTPSession').as_null_object + pool = ConnectionPool.for({}) + allow(pool).to receive(:start_session).and_return(session) + expect(session).to receive(:finish) + expect do + pool.session_for(URI.parse(endpoint)) { raise Interrupt } + end.to raise_error(Interrupt) + expect(pool.size).to eq(0) + end end describe '#finish_session' do @@ -76,6 +102,26 @@ module NetHTTP pool.finish_session(session) expect(pool.size).to eq(0) end + + it 'removes the session from the pool using the given endpoint' do + session = double('Net::HTTPSession').as_null_object + pool = ConnectionPool.for({}) + allow(pool).to receive(:start_session).and_return(session) + pool.session_for(URI.parse(endpoint), &:request) + expect(pool.size).to eq(1) + # Passing the endpoint scopes the removal to that endpoint's list. + pool.finish_session(session, URI.parse(endpoint)) + expect(pool.size).to eq(0) + end + + it 'still finishes when the session is not pooled (e.g. aborted in flight)' do + session = double('Net::HTTPSession') + pool = ConnectionPool.for({}) + # Not in the pool at all; abort discards an in-flight session. + expect(session).to receive(:finish) + expect { pool.finish_session(session, URI.parse(endpoint)) }.not_to raise_error + expect(pool.size).to eq(0) + end end describe '#size' do diff --git a/gems/smithy-client/spec/smithy-client/net_http/exchange_spec.rb b/gems/smithy-client/spec/smithy-client/net_http/exchange_spec.rb new file mode 100644 index 000000000..487c20e69 --- /dev/null +++ b/gems/smithy-client/spec/smithy-client/net_http/exchange_spec.rb @@ -0,0 +1,195 @@ +# frozen_string_literal: true + +require_relative '../../spec_helper' + +module Smithy + module Client + module NetHTTP + describe Exchange do + # Records the pushed response lifecycle for assertions. + let(:sink) { RecordingSink.new } + let(:pool) { ConnectionPool.for({}) } + let(:endpoint) { 'https://example.com' } + let(:http_method) { 'GET' } + let(:body) { nil } + let(:request) do + Http::Request.new(endpoint: endpoint, http_method: http_method, body: body) + end + + subject { described_class.new(pool, request, sink) } + + describe '#initialize' do + it 'raises ArgumentError for an invalid http verb (without networking)' do + request.http_method = 'bogus' + expect { described_class.new(pool, request, sink) } + .to raise_error(ArgumentError, /not a valid http verb/) + end + end + + describe '#drive' do + it 'pushes headers then data then done for a successful response' do + stub_request(:get, endpoint) + .to_return(status: 201, headers: { 'X-Foo' => 'bar' }, body: 'hello-world') + subject.drive + expect(sink.status).to eq(201) + expect(sink.headers_hash['x-foo']).to eq('bar') + expect(sink.body).to eq('hello-world') + expect(sink.terminal).to eq(:done) + end + + it 'pushes headers and done with no data for an empty body' do + stub_request(:get, endpoint).to_return(status: 200, body: '') + subject.drive + expect(sink.chunks).to be_empty + expect(sink.terminal).to eq(:done) + end + + it 'sends the request body' do + @body = StringIO.new('request-body') + stub = stub_request(:post, endpoint).with(body: 'request-body') + post = Http::Request.new(endpoint: endpoint, http_method: 'POST', body: @body) + described_class.new(pool, post, RecordingSink.new).drive + expect(stub).to have_been_requested + end + + it 'surfaces networking errors as sink.error(NetworkingError) (not raised)' do + stub_request(:get, endpoint).to_raise(EOFError) + expect { subject.drive }.not_to raise_error + expect(sink.terminal).to eq(:error) + expect(sink.error_value).to be_a(Smithy::Client::NetworkingError) + end + + it 'does not convert a raising :done listener into a second (error) terminal' do + # @sink.done runs caller :done listeners. A bug in one must propagate + # to the caller, NOT be caught by the networking-failure rescue and + # turned into a second sink.error terminal (which would also make the + # caller bug look like a transient NetworkingError and re-download). + stub_request(:get, endpoint).to_return(status: 200, body: 'ok') + raising_sink = RecordingSink.new + boom = RuntimeError.new('listener blew up') + raising_sink.define_singleton_method(:done) do + super() + raise boom + end + expect { described_class.new(pool, request, raising_sink).drive } + .to raise_error(boom) + # Terminal recorded is :done (from super), never overwritten by :error. + expect(raising_sink.terminal).to eq(:done) + end + + it 'surfaces a short Content-Length body as a NetworkingError terminal' do + stub_request(:get, endpoint) + .to_return(status: 200, headers: { 'Content-Length' => '100' }, body: 'short') + subject.drive + expect(sink.terminal).to eq(:error) + expect(sink.error_value).to be_a(Smithy::Client::NetworkingError) + end + + it 'wraps the underlying truncation error as the NetworkingError cause' do + # Truncation is delivered as a NetworkingError (so it retries like any + # networking failure), but callers can distinguish it via the cause's + # message. Pin that so it stays a supported detection path. + stub_request(:get, endpoint) + .to_return(status: 200, headers: { 'Content-Length' => '100' }, body: 'short') + subject.drive + expect(sink.error_value.original_error).to be_a(IOError) + expect(sink.error_value.original_error.message) + .to eq('http response body truncated, expected 100 bytes, received 5 bytes') + end + + it 'does not verify bytes for HEAD requests' do + request.http_method = 'HEAD' + stub_request(:head, endpoint).to_return(headers: { 'Content-Length' => '100' }) + subject.drive + expect(sink.terminal).to eq(:done) + end + end + + describe '#drive_background' do + it 'drives the exchange on a background thread' do + stub_request(:get, endpoint).to_return(status: 200, body: 'ok') + thread = subject.drive_background + expect(thread).to be_a(Thread) + thread.join(5) + expect(sink.status).to eq(200) + expect(sink.body).to eq('ok') + expect(sink.terminal).to eq(:done) + end + end + + describe '#abort' do + it 'does not raise after the exchange has completed' do + stub_request(:get, endpoint).to_return(body: 'data') + subject.drive + expect { subject.abort }.not_to raise_error + end + + it 'is a no-op once the exchange has completed and the session is pooled' do + # After drive completes the session has been returned to the pool and + # the exchange relinquished ownership; a late cross-thread abort must + # not reach through and finish the pooled session. + stub_request(:get, endpoint).to_return(body: 'data') + subject.drive + expect(pool).not_to receive(:finish_session) + expect { subject.abort }.not_to raise_error + end + + it 'discards the session through the pool when aborting before completion' do + expect(pool).to receive(:finish_session) + subject.abort + end + + it 'is idempotent' do + expect(pool).to receive(:finish_session).once + subject.abort + expect { subject.abort }.not_to raise_error + end + end + + describe 'no sink delivery after abort' do + # Records every sink call so we can assert what was (not) delivered. + let(:recording_sink) do + Class.new do + def initialize = @calls = [] + + attr_reader :calls + + def headers(status, _headers) = @calls << [:headers, status] + def data(chunk) = @calls << [:data, chunk] + def done = @calls << [:done] + def error(err) = @calls << [:error, err] + end + end + + it 'delivers nothing when the exchange is aborted before it is driven' do + # abort() before drive() records @aborted; #deliver then suppresses + # both headers and body, and #run emits no success terminal. This is + # the deterministic core of the abort/no-deliver contract: #deliver + # checks @aborted and calls the sink atomically under @mutex, so once + # @aborted is observed nothing further is delivered. + stub_request(:get, endpoint).to_return(status: 200, body: 'body-bytes') + sink = recording_sink.new + exchange = described_class.new(pool, request, sink) + + exchange.abort # record the abort first + exchange.drive # drive; every deliver{} sees @aborted + + expect(sink.calls).to be_empty # no headers, data, done, or error + end + + it 'does not re-pool the session when delivery stops on abort' do + # When #deliver suppresses delivery (rather than the socket close + # raising), perform_exchange can return normally; #run must still + # discard the session instead of letting #session_for pool a socket + # that abort already finished. + stub_request(:get, endpoint).to_return(status: 200, body: 'body-bytes') + exchange = described_class.new(pool, request, recording_sink.new) + exchange.abort + exchange.drive + expect(pool.size).to eq(0) # nothing pooled + end + end + end + end + end +end diff --git a/gems/smithy-client/spec/smithy-client/net_http/stream_spec.rb b/gems/smithy-client/spec/smithy-client/net_http/stream_spec.rb index a3117afd6..72049ae03 100644 --- a/gems/smithy-client/spec/smithy-client/net_http/stream_spec.rb +++ b/gems/smithy-client/spec/smithy-client/net_http/stream_spec.rb @@ -6,140 +6,70 @@ module Smithy module Client module NetHTTP describe Stream do + let(:sink) { RecordingSink.new } let(:pool) { ConnectionPool.for({}) } let(:endpoint) { 'https://example.com' } - let(:http_method) { 'GET' } - let(:body) { nil } let(:request) do - Http::Request.new(endpoint: endpoint, http_method: http_method, body: body) + Http::Request.new(endpoint: endpoint, http_method: 'GET', body: nil) end + let(:exchange) { Exchange.new(pool, request, sink) } - subject { described_class.new(pool, request) } + subject { described_class.new(exchange) } it_behaves_like 'a stream' do - def build_stream(body:, status: 200, headers: {}) + def build_handle(sink, body:, status: 200, headers: {}) stub_request(:get, 'https://example.com') .to_return(status: status, headers: headers, body: body) - Smithy::Client::NetHTTP::Stream.new( - ConnectionPool.for({}), - Http::Request.new(endpoint: 'https://example.com', http_method: 'GET', body: nil) - ).send_request - end - end - - describe '#send_request' do - it 'returns self and reads status and headers' do - stub_request(:get, endpoint) - .to_return(status: 201, headers: { 'X-Foo' => 'bar' }, body: '') - expect(subject.send_request).to be(subject) - status, headers = subject.response_headers - expect(status).to eq(201) - expect(headers['x-foo']).to eq('bar') - end - - it 'sends the request body' do - @body = StringIO.new('request-body') - stub = stub_request(:post, endpoint).with(body: 'request-body') - post = Http::Request.new(endpoint: endpoint, http_method: 'POST', body: @body) - described_class.new(pool, post).send_request - expect(stub).to have_been_requested - end - - it 'raises ArgumentError for an invalid http verb (without networking)' do - request.http_method = 'bogus' - expect { subject.send_request }.to raise_error(ArgumentError, /not a valid http verb/) + pool = ConnectionPool.for({}) + request = Http::Request.new(endpoint: 'https://example.com', http_method: 'GET', body: nil) + @exchange = Exchange.new(pool, request, sink) + Smithy::Client::NetHTTP::Stream.new(@exchange) end - it 'wraps networking errors in a NetworkingError' do - stub_request(:get, endpoint).to_raise(EOFError) - expect { subject.send_request }.to raise_error(Smithy::Client::NetworkingError) + # Synchronously runs the exchange behind the handle to completion so + # the abort contract can assert an observable "no terminal delivered" + # effect (see stream_contract.rb). + def drive_handle(_handle) + @exchange.drive end end - describe '#each_chunk' do - it 'yields the response body chunks' do - stub_request(:get, endpoint).to_return(body: 'hello-world') - subject.send_request - chunks = [] - subject.each_chunk { |c| chunks << c } - expect(chunks.join).to eq('hello-world') + describe 'Stream contract conformance' do + it 'answers the output-only event-stream method tier' do + expect(subject).to respond_to(*Client::Stream::REQUIRED_FOR_OUTPUT_EVENT_STREAM_OPS) end - it 'yields nothing for an empty body' do - stub_request(:get, endpoint).to_return(body: '') - subject.send_request - chunks = [] - subject.each_chunk { |c| chunks << c } - expect(chunks).to be_empty + it 'answers the bidirectional event-stream method tier' do + expect(subject).to respond_to(*Client::Stream::REQUIRED_FOR_BIDI_EVENT_STREAM_OPS) end + end - it 'is safe to call after the body has been consumed' do - stub_request(:get, endpoint).to_return(body: 'data') - subject.send_request - subject.each_chunk { |_c| } # drain - expect { |b| subject.each_chunk(&b) }.not_to yield_control + describe '#abort' do + it 'delegates to the exchange' do + expect(exchange).to receive(:abort).with(nil) + subject.abort end - it 'raises a NetworkingError when the body is shorter than Content-Length' do - stub_request(:get, endpoint) - .to_return(status: 200, headers: { 'Content-Length' => '100' }, body: 'short') - subject.send_request - chunks = [] - expect { subject.each_chunk { |c| chunks << c } } - .to raise_error(Smithy::Client::NetworkingError) + it 'forwards an error argument to the exchange' do + error = StandardError.new('boom') + expect(exchange).to receive(:abort).with(error) + subject.abort(error) end - it 're-raises the consumer block error and aborts the stream' do - stub_request(:get, endpoint).to_return(body: 'hello-world') - subject.send_request - expect { subject.each_chunk { |_c| raise 'consumer boom' } } # rubocop:disable Lint/UnreachableLoop - .to raise_error(RuntimeError, 'consumer boom') - # The stream was torn down; a subsequent drain yields nothing. - expect { |b| subject.each_chunk(&b) }.not_to yield_control + it 'does not raise' do + expect { subject.abort }.not_to raise_error end end describe '#write / #close_write' do - before do - stub_request(:get, endpoint).to_return(body: '') - subject.send_request - end - - it '#write raises NotSupportedError' do + it '#write raises NotSupportedError (HTTP/1.1 is not bidirectional)' do expect { subject.write('data') }.to raise_error(Smithy::Client::NotSupportedError) end - it '#close_write raises NotSupportedError' do + it '#close_write raises NotSupportedError (HTTP/1.1 is not bidirectional)' do expect { subject.close_write }.to raise_error(Smithy::Client::NotSupportedError) end end - - describe '#abort' do - it 'does not raise and prevents further reads' do - stub_request(:get, endpoint).to_return(body: 'data') - subject.send_request - expect { subject.abort }.not_to raise_error - expect { |b| subject.each_chunk(&b) }.not_to yield_control - end - - it 'is a no-op once the body has completed and the session is pooled' do - # After a full read the session has been returned to the pool and the - # stream has relinquished ownership; a late cross-thread abort must - # not reach through and finish the pooled session. - stub_request(:get, endpoint).to_return(body: 'data') - subject.send_request - subject.each_chunk { |_c| } # drain to completion, session re-pooled - expect(pool).not_to receive(:finish_session) - expect { subject.abort }.not_to raise_error - end - - it 'discards the session through the pool when aborting mid-stream' do - stub_request(:get, endpoint).to_return(body: 'data') - subject.send_request - expect(pool).to receive(:finish_session) - subject.abort - end - end end end end diff --git a/gems/smithy-client/spec/smithy-client/net_http/transport_spec.rb b/gems/smithy-client/spec/smithy-client/net_http/transport_spec.rb index c368d607c..8159cd998 100644 --- a/gems/smithy-client/spec/smithy-client/net_http/transport_spec.rb +++ b/gems/smithy-client/spec/smithy-client/net_http/transport_spec.rb @@ -9,7 +9,6 @@ module Client module NetHTTP describe Transport do let(:endpoint) { 'https://example.com' } - let(:request) { Http::Request.new(endpoint: endpoint, http_method: 'GET') } subject { described_class.new } @@ -17,27 +16,20 @@ module NetHTTP let(:transport) { described_class.new } end - describe '#transmit' do - it 'returns a Stream with the response headers available' do - stub_request(:get, endpoint).to_return(status: 200, headers: { 'X-A' => 'b' }, body: 'ok') - stream = subject.transmit(request) - expect(stream).to be_a(Stream) - status, headers = stream.response_headers - expect(status).to eq(200) - expect(headers['x-a']).to eq('b') - chunks = [] - stream.each_chunk { |c| chunks << c } - expect(chunks.join).to eq('ok') - end - - it 'propagates ArgumentError for an invalid verb' do - request.http_method = 'nope' - expect { subject.transmit(request) }.to raise_error(ArgumentError) - end + # NOTE: transmit / transmit_background behavior (push into sink, return + # value, invalid-verb ArgumentError, networking-failure terminal, handle + # returned for the background path) is covered by the shared 'a transport' + # compliance examples above. The specs below cover only what is specific + # to this transport: event_queue, default wiring, and option mapping. - it 'raises NetworkingError on a networking failure' do - stub_request(:get, endpoint).to_raise(SocketError) - expect { subject.transmit(request) }.to raise_error(Smithy::Client::NetworkingError) + describe '#event_queue' do + it 'returns a fresh SizedQueue with the bridge capacity' do + q1 = subject.event_queue + q2 = subject.event_queue + expect(q1).to be_a(SizedQueue) + expect(q1.max).to eq(64) + # Fresh per call - the event stream layer gets its own bridge queue. + expect(q1).not_to be(q2) end end diff --git a/gems/smithy-client/spec/smithy-client/plugins/transport_spec.rb b/gems/smithy-client/spec/smithy-client/plugins/transport_spec.rb index ad0062051..eb9848838 100644 --- a/gems/smithy-client/spec/smithy-client/plugins/transport_spec.rb +++ b/gems/smithy-client/spec/smithy-client/plugins/transport_spec.rb @@ -31,7 +31,7 @@ module Plugins end it 'uses a caller-supplied transport as-is' do - custom = double('transport', transmit: nil) + custom = double('transport', transmit: nil, transmit_background: nil) expect(client_class.new(transport: custom).config.transport).to be(custom) end @@ -43,7 +43,16 @@ module Plugins it 'raises when a caller-supplied transport does not satisfy the contract' do expect { client_class.new(transport: Object.new) } - .to raise_error(ArgumentError, /does not implement the transport contract \(transmit\)/) + .to raise_error(ArgumentError, /does not implement the transport contract/) + end + + it 'raises when a transport answers only transmit (both modes required)' do + # Both transmit and transmit_background are required; a transport that + # serves only one mode must still answer both (raising NotSupportedError + # from the unsupported one), so a transmit-only object is not conforming. + partial = double('transport', transmit: nil) + expect { client_class.new(transport: partial) } + .to raise_error(ArgumentError, /transmit_background/) end it 'forwards the transport-agnostic options to the default transport' do diff --git a/gems/smithy-client/spec/smithy-client/send_handler_spec.rb b/gems/smithy-client/spec/smithy-client/send_handler_spec.rb index ef41c0744..5f0357671 100644 --- a/gems/smithy-client/spec/smithy-client/send_handler_spec.rb +++ b/gems/smithy-client/spec/smithy-client/send_handler_spec.rb @@ -41,10 +41,11 @@ def endpoint expect(make_request.context).to be(context) end - it 'stores the stream on the context' do + it 'does not store a stream on the context for a non-event-stream operation' do stub_request(:any, endpoint) make_request - expect(context[:stream]).to be_a(NetHTTP::Stream) + # transmit drives inline and returns nothing; there is no handle. + expect(context[:stream]).to be_nil end describe 'request' do @@ -96,6 +97,18 @@ def endpoint expect(make_request.error).to be_a(NetworkingError) end + it 'signals NotSupportedError as-is (not wrapped into a retryable NetworkingError)' do + # A single-mode transport raises NotSupportedError from #transmit. + # NotSupportedError < StandardError, so it must be rescued ahead of + # the networking-failure clause; otherwise it would be wrapped into a + # transient NetworkingError and retried with backoff for an operation + # that can never succeed. + error = NotSupportedError.new('this transport does not serve request/response') + allow(context.config.transport).to receive(:transmit).and_raise(error) + expect(make_request.error).to be(error) + expect(make_request.error).not_to be_a(NetworkingError) + end + it 'raises when content length and body length mismatch' do stub_request(:any, endpoint).to_return(body: 'foo', headers: { 'Content-Length' => 1 }) expect(make_request.error).to be_a(NetworkingError) @@ -108,16 +121,17 @@ def endpoint end end - describe 'duplex (bidirectional) streams' do - it 'stores the stream but does not resolve the response' do + describe 'event streams' do + it 'stores the stream (via transmit_background) but does not resolve the response inline' do stub_request(:any, endpoint).to_return(status: 200, body: 'data') - context[:duplex_stream] = true + context[:event_stream] = true make_request - # Stream is available for the event stream layer to drive... + # The handle is available for the event stream layer to pump. expect(context[:stream]).to be_a(NetHTTP::Stream) - # ...but the handler did not block for / populate the response. - expect(context.http_response.status_code).to eq(0) - expect(context.http_response.body.read).to eq('') + expect(context[:stream]).to respond_to(:abort) + # The handler did not drive synchronously; clean up the background + # exchange. + context[:stream].abort end end end diff --git a/gems/smithy-client/spec/spec_helper.rb b/gems/smithy-client/spec/spec_helper.rb index bfb85192a..aa21dae9e 100644 --- a/gems/smithy-client/spec/spec_helper.rb +++ b/gems/smithy-client/spec/spec_helper.rb @@ -15,6 +15,7 @@ require 'smithy' require_relative 'support/client_helper' +require_relative 'support/recording_sink' require_relative 'support/transport_contract' require_relative 'support/stream_contract' diff --git a/gems/smithy-client/spec/support/recording_sink.rb b/gems/smithy-client/spec/support/recording_sink.rb new file mode 100644 index 000000000..3897d628e --- /dev/null +++ b/gems/smithy-client/spec/support/recording_sink.rb @@ -0,0 +1,69 @@ +# frozen_string_literal: true + +require 'timeout' + +# A test sink that records the response lifecycle a transport pushes into it. +# Satisfies the sink side of the {Smithy::Client::Transport} contract +# (headers / data / done / error) and exposes the recorded values for +# assertions. Recorded-value readers use distinct names (received_*) so they do +# not collide with the sink methods themselves. +class RecordingSink + def initialize + @status = nil + @received_headers = nil + @chunks = [] + @terminal = nil + @received_error = nil + end + + attr_reader :status, :received_headers, :chunks, :terminal, :received_error + + # --- sink contract --- + + def headers(status_code, headers) + @status = status_code + @received_headers = headers + end + + def data(chunk) + @chunks << chunk + end + + def done + @terminal = :done + end + + def error(error) + @terminal = :error + @received_error = error + end + + # --- convenience accessors for assertions --- + + # @return [Hash, nil] the received response headers. + def headers_hash + @received_headers + end + + # @return [StandardError, nil] the error passed to the error terminal. + def error_value + @received_error + end + + # @return [String] the concatenated body chunks. + def body + @chunks.join + end + + # Blocks until a terminal (#done or #error) has been recorded, for use with a + # background exchange (transmit_background / drive_background) whose driving + # thread pushes into this sink asynchronously. Raises Timeout::Error if no + # terminal arrives within +timeout+ seconds so a hung exchange fails the + # example instead of blocking the suite. + # @param [Numeric] timeout Seconds to wait. + # @return [Symbol] the terminal (+:done+ or +:error+). + def wait_for_terminal(timeout: 5) + Timeout.timeout(timeout) { sleep 0.01 until @terminal } + @terminal + end +end diff --git a/gems/smithy-client/spec/support/stream_contract.rb b/gems/smithy-client/spec/support/stream_contract.rb index 0ed54684b..d07d41bef 100644 --- a/gems/smithy-client/spec/support/stream_contract.rb +++ b/gems/smithy-client/spec/support/stream_contract.rb @@ -1,55 +1,67 @@ # frozen_string_literal: true -# Shared compliance tests for the {Smithy::Client::Stream} contract. Any stream -# implementation should exercise these to validate conformance. +# Shared compliance tests for the {Smithy::Client::Stream} control-handle +# contract - the handle returned by {Smithy::Client::Transport#transmit_background} +# for an event-stream operation. Any handle implementation should exercise these +# to validate conformance. # # Usage: # # it_behaves_like 'a stream' do -# # return a fresh, already-transmitted stream for the given body -# def build_stream(body:, status: 200, headers: {}) +# # return a fresh control handle wired to the given sink +# def build_handle(sink, body:, status: 200, headers: {}) # ... # end # end # -# Requires the including group to define a +build_stream+ helper that returns a -# stream whose request has already been transmitted (response headers ready). +# Requires the including group to define a +build_handle(sink, ...)+ helper that +# returns a control handle wired to push into +sink+, and a +# +drive_handle(handle)+ helper that synchronously runs the handle's exchange to +# completion (so the +#abort+ examples can assert an OBSERVABLE effect: an +# aborted stream delivers no terminal into the sink, rather than merely not +# raising). +# +# A +Stream+ is an OUTBOUND + CONTROL handle only. Inbound response data flows +# into the sink, not back through the handle, so this contract covers +#abort+ +# (output-only tier) and +#write+/+#close_write+ (bidirectional tier). RSpec.shared_examples 'a stream' do - it 'exposes response status and headers' do - stream = build_stream(body: '', status: 201, headers: { 'X-Foo' => 'bar' }) - status, headers = stream.response_headers - expect(status).to eq(201) - expect(headers['x-foo']).to eq('bar') - end + let(:sink) { RecordingSink.new } - it 'yields body chunks in order' do - stream = build_stream(body: 'hello-world') - chunks = [] - stream.each_chunk { |c| chunks << c } - expect(chunks.join).to eq('hello-world') - end - - it 'yields nothing for an empty body' do - expect { |b| build_stream(body: '').each_chunk(&b) }.not_to yield_control - end + describe '#write / #close_write (non-bidirectional / HTTP/1.1)' do + it '#write raises NotSupportedError' do + handle = build_handle(sink, body: '') + expect { handle.write('data') }.to raise_error(Smithy::Client::NotSupportedError) + end - it 'is safe to iterate again after the body is consumed' do - stream = build_stream(body: 'data') - stream.each_chunk { |_c| nil } - expect { |b| stream.each_chunk(&b) }.not_to yield_control + it '#close_write raises NotSupportedError' do + handle = build_handle(sink, body: '') + expect { handle.close_write }.to raise_error(Smithy::Client::NotSupportedError) + end end describe '#abort' do - it 'does not raise and stops further reads' do - stream = build_stream(body: 'data') - expect { stream.abort }.not_to raise_error - expect { |b| stream.each_chunk(&b) }.not_to yield_control + it 'does not raise' do + handle = build_handle(sink, body: 'data') + expect { handle.abort }.not_to raise_error end it 'is idempotent' do - stream = build_stream(body: 'data') - stream.abort - expect { stream.abort }.not_to raise_error + handle = build_handle(sink, body: 'data') + handle.abort + expect { handle.abort }.not_to raise_error + end + + it 'delivers no terminal into the sink once aborted' do + # Observable effect (not just "doesn't raise"): a stream aborted before it + # is driven delivers NOTHING to the sink - no headers, data, done, or + # error. Gutting the abort would let the driven exchange push a terminal + # here and fail this example. + handle = build_handle(sink, body: 'data') + handle.abort + drive_handle(handle) + expect(sink.terminal).to be_nil + expect(sink.status).to be_nil + expect(sink.chunks).to be_empty end end end diff --git a/gems/smithy-client/spec/support/transport_contract.rb b/gems/smithy-client/spec/support/transport_contract.rb index 40e20941d..c88a462a0 100644 --- a/gems/smithy-client/spec/support/transport_contract.rb +++ b/gems/smithy-client/spec/support/transport_contract.rb @@ -14,6 +14,7 @@ # * +transport+ - the transport instance under test # * +endpoint+ - a String endpoint the transport can reach (stubbed via WebMock) RSpec.shared_examples 'a transport' do + let(:sink) { RecordingSink.new } let(:request) do Smithy::Client::Http::Request.new( endpoint: endpoint, http_method: 'GET', body: nil @@ -24,37 +25,55 @@ expect(transport).to respond_to(:transmit) end - it 'returns a stream conforming to the stream contract' do - stub_request(:get, endpoint).to_return(status: 200, body: 'ok') - stream = transport.transmit(request) - expect(stream).to respond_to(:response_headers, :each_chunk, :abort) - end - - it 'makes response headers available without consuming the body' do + it 'pushes the response lifecycle into the sink and returns nothing' do stub_request(:get, endpoint) - .to_return(status: 201, headers: { 'X-Foo' => 'bar' }, body: 'body') - stream = transport.transmit(request) - status, headers = stream.response_headers - expect(status).to eq(201) - expect(headers['x-foo']).to eq('bar') - end - - it 'delivers the body in order via #each_chunk' do - stub_request(:get, endpoint).to_return(status: 200, body: 'hello-world') - stream = transport.transmit(request) - chunks = [] - stream.each_chunk { |c| chunks << c } - expect(chunks.join).to eq('hello-world') + .to_return(status: 201, headers: { 'X-Foo' => 'bar' }, body: 'hello-world') + result = transport.transmit(request, sink) + # transmit drives to completion synchronously and returns nothing. + expect(result).to be_nil + expect(sink.status).to eq(201) + expect(sink.headers_hash['x-foo']).to eq('bar') + expect(sink.body).to eq('hello-world') + expect(sink.terminal).to eq(:done) end - it 'raises ArgumentError for an invalid http method' do + it 'raises ArgumentError for an invalid http method (at transmit, no network)' do request.http_method = 'bogus' - expect { transport.transmit(request) }.to raise_error(ArgumentError) + expect { transport.transmit(request, sink) }.to raise_error(ArgumentError) end - it 'raises NetworkingError on a networking failure' do + it 'surfaces a networking failure as a sink.error terminal (not raised)' do stub_request(:get, endpoint).to_raise(EOFError) - expect { transport.transmit(request) } - .to raise_error(Smithy::Client::NetworkingError) + expect { transport.transmit(request, sink) }.not_to raise_error + expect(sink.terminal).to eq(:error) + expect(sink.error_value).to be_a(Smithy::Client::NetworkingError) + end + + describe '#transmit_background (event streams)' do + it 'returns a Stream handle immediately' do + stub_request(:get, endpoint).to_return(status: 200, body: 'ok') + handle = transport.transmit_background(request, sink) + expect(handle).to respond_to(:abort) + # Let the background exchange complete so the thread is joined cleanly. + sink.wait_for_terminal + handle.abort + end + + it 'drives the response into the sink on the background exchange' do + stub_request(:get, endpoint) + .to_return(status: 200, headers: { 'X-Foo' => 'bar' }, body: 'hello-world') + handle = transport.transmit_background(request, sink) + # The exchange runs concurrently; wait for it to terminate. + sink.wait_for_terminal + expect(sink.status).to eq(200) + expect(sink.body).to eq('hello-world') + expect(sink.terminal).to eq(:done) + handle.abort + end + + it 'raises ArgumentError for an invalid http method (before spawning)' do + request.http_method = 'bogus' + expect { transport.transmit_background(request, sink) }.to raise_error(ArgumentError) + end end end