diff --git a/CHANGELOG.md b/CHANGELOG.md index 846df76..6b529bc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,23 @@ # Changelog +## Unreleased + +- `timeout` now covers every phase of the default transport, not only connect + and read. A stalled TLS handshake or a stalled upload used to fall back to + `Net::HTTP`'s own 60-second write timeout, well past both `timeout` and + `total_timeout`. +- `Net::WriteTimeout` is classified as a timeout, so `retry_timeouts` governs + it like the other two instead of it surfacing as an unretried + `TransportError`. +- `total_timeout` is enforced as a deadline. Each attempt is given the smaller + of `timeout` and the remaining budget, an attempt that would start with no + budget left raises `TimeoutError`, and a delay landing exactly on the + deadline now stops the retry loop rather than allowing one more attempt. +- The transport contract accepts an optional `timeout:` keyword carrying that + per-attempt budget. Transports that do not declare it are called unchanged. +- `timeout:` is validated when the client is built: nil, or a finite positive + number. Anything else raises `ConfigurationError`. + ## 0.1.0 - 2026-09-18 Provider-neutral release. One `Client`, two providers behind it. diff --git a/README.md b/README.md index 79638d6..92f564d 100644 --- a/README.md +++ b/README.md @@ -76,7 +76,7 @@ RubyDecisionModel::Client.new( api_key: nil, # overrides the provider's env var model: nil, # nil means the provider default; see aliases below base_url: nil, # overrides the provider base URL - timeout: 5, # open and read timeout in seconds + timeout: 5, # per-attempt timeout in seconds, nil for none retry: { max_retries: 2 }, # RetryPolicy or a Hash of overrides transport: nil # see Transport ) @@ -147,9 +147,11 @@ Retry behaviour follows the official Typesafe SDKs and lives in | `retry_timeouts` | `true` | Retry open and read timeouts | | `total_timeout` | `30.0` | Budget in seconds across attempts and delays; `nil` disables | -When the next delay would push past `total_timeout`, the client stops and -raises the last error instead of sleeping. The budget governs whether another -attempt starts; an attempt already in flight still runs to its own `timeout`. +`total_timeout` is a deadline, not just a gate between attempts. Each attempt +is given the smaller of `timeout` and what is left of the budget, so a single +slow attempt cannot outlive the whole call. When the next delay would reach or +pass the budget the client stops and raises the last error instead of sleeping, +and an attempt that would start with nothing left raises `TimeoutError`. Invalid settings (a negative duration, a non-integer `max_retries`, a jitter outside 0..1, a NaN budget) raise `ConfigurationError` when the client is built. @@ -161,17 +163,24 @@ RubyDecisionModel::Client.new(retry: RubyDecisionModel::RetryPolicy.new(max_retr ## Transport -The client uses `Net::HTTP` by default. Inject `transport:` with any callable -that accepts `url:`, `headers:`, `body:` and returns -`[status, body_string, headers_hash]`. A two-element `[status, body_string]` -return is still accepted and treated as having no headers, which means no -`Retry-After` support and a nil `request_id`. +The client uses `Net::HTTP` by default, with `timeout` applied to all four of +its phases: connect, TLS handshake, write, and read. + +Inject `transport:` with any callable that accepts `url:`, `headers:`, `body:` +and returns `[status, body_string, headers_hash]`. A two-element +`[status, body_string]` return is still accepted and treated as having no +headers, which means no `Retry-After` support and a nil `request_id`. + +A transport that also declares a `timeout:` keyword (or `**`) is handed the +number of seconds this attempt may take, already clamped to what is left of +`total_timeout`. Transports that do not declare it are called exactly as +before. ## Errors | Error | Meaning | | --- | --- | -| `ConfigurationError` | No provider could be resolved, missing api_key, unknown provider, or bad `retry:` value | +| `ConfigurationError` | No provider could be resolved, missing api_key, unknown provider, or bad `timeout:` or `retry:` value | | `RequestError` | Questions hash was empty | | `TransportError` (`TimeoutError`) | Network or timeout failure after retries, carries `#cause_error` | | `ApiError` | Non-2xx response, carries `#status`, `#body`, and `#headers` | diff --git a/lib/ruby_decision_model/client.rb b/lib/ruby_decision_model/client.rb index 3457726..701119f 100644 --- a/lib/ruby_decision_model/client.rb +++ b/lib/ruby_decision_model/client.rb @@ -41,7 +41,7 @@ def initialize(provider: nil, api_key: nil, model: nil, base_url: nil, timeout: end @model = @provider.resolve_model(model) - @timeout = timeout + @timeout = validate_timeout(timeout) @transport = transport || default_transport @sleeper = sleeper @retry_policy = RetryPolicy.from(binding.local_variable_get(:retry)) @@ -89,12 +89,18 @@ def resolve_provider(provider, api_key:, base_url:) end def default_transport - lambda do |url:, headers:, body:| + lambda do |url:, headers:, body:, timeout: @timeout| uri = URI.parse(url) http = Net::HTTP.new(uri.host, uri.port) http.use_ssl = uri.scheme == "https" - http.open_timeout = @timeout - http.read_timeout = @timeout + # All four, not just open and read: a stalled TLS handshake or a + # stalled upload otherwise falls back to Net::HTTP's own defaults + # (60s for the write), which blows through both `timeout` and the + # retry policy's total budget. + http.open_timeout = timeout + http.ssl_timeout = timeout + http.read_timeout = timeout + http.write_timeout = timeout request = Net::HTTP::Post.new(uri.request_uri) headers.each { |k, v| request[k] = v } @@ -105,6 +111,34 @@ def default_transport end end + def validate_timeout(timeout) + return timeout if timeout.nil? + return timeout if timeout.is_a?(Numeric) && timeout.finite? && timeout.positive? + + raise ConfigurationError, "timeout must be nil or a finite positive number, got #{timeout.inspect}" + end + + # The transport contract grew a `timeout:` keyword so the client can hand + # each attempt what is left of the retry budget. Transports written + # against the old three-keyword contract still work: they are called the + # way they always were. + def transport_accepts_timeout? + return @transport_accepts_timeout unless @transport_accepts_timeout.nil? + + parameters = @transport.respond_to?(:parameters) ? @transport.parameters : @transport.method(:call).parameters + @transport_accepts_timeout = parameters.any? do |kind, name| + kind == :keyrest || (%i[key keyreq].include?(kind) && name == :timeout) + end + end + + def call_transport(url:, headers:, body:, timeout:) + if transport_accepts_timeout? + @transport.call(url: url, headers: headers, body: body, timeout: timeout) + else + @transport.call(url: url, headers: headers, body: body) + end + end + def perform_with_retry(url:, headers:, body:) policy = @retry_policy started_at = @clock.call @@ -113,7 +147,8 @@ def perform_with_retry(url:, headers:, body:) loop do begin status, response_body, response_headers = normalize_transport_result( - @transport.call(url: url, headers: headers, body: body) + call_transport(url: url, headers: headers, body: body, + timeout: attempt_timeout(policy, started_at)) ) rescue Error raise @@ -146,7 +181,32 @@ def perform_with_retry(url:, headers:, body:) def budget_exceeded?(policy, started_at, delay) return false if policy.total_timeout.nil? - (@clock.call - started_at) + delay > policy.total_timeout + # >=, not >: a delay that lands exactly on the deadline has used the + # whole budget, and the attempt after it would start with nothing left. + (@clock.call - started_at) + delay >= policy.total_timeout + end + + # What is left of the budget, or nil when there is no budget. Handed to + # the transport so a single attempt cannot outlive the whole call: with + # `timeout: 5` and 2s of budget left, the attempt gets 2s. + def remaining_budget(policy, started_at) + return nil if policy.total_timeout.nil? + + policy.total_timeout - (@clock.call - started_at) + end + + def attempt_timeout(policy, started_at) + remaining = remaining_budget(policy, started_at) + return @timeout if remaining.nil? + + if remaining <= 0 + raise TimeoutError.new( + "request budget of #{policy.total_timeout}s was exhausted before the attempt started", + cause_error: nil + ) + end + + @timeout.nil? ? remaining : [@timeout, remaining].min end def normalize_transport_result(result) diff --git a/lib/ruby_decision_model/retry_policy.rb b/lib/ruby_decision_model/retry_policy.rb index 00774c5..75bda58 100644 --- a/lib/ruby_decision_model/retry_policy.rb +++ b/lib/ruby_decision_model/retry_policy.rb @@ -21,7 +21,10 @@ class RetryPolicy < Data.define( ) DEFAULT_STATUSES = ([408, 429] + (500..599).to_a).freeze - TIMEOUT_EXCEPTIONS = [Net::OpenTimeout, Net::ReadTimeout].freeze + # Net::WriteTimeout fires when the request body stalls on the way out. + # It is as much a timeout as the other two, and leaving it off this list + # made it a plain TransportError that was never retried. + TIMEOUT_EXCEPTIONS = [Net::OpenTimeout, Net::ReadTimeout, Net::WriteTimeout].freeze CONNECTION_EXCEPTIONS = [ Errno::ECONNRESET, Errno::ECONNREFUSED, diff --git a/test/deadline_test.rb b/test/deadline_test.rb new file mode 100644 index 0000000..29e349e --- /dev/null +++ b/test/deadline_test.rb @@ -0,0 +1,223 @@ +# frozen_string_literal: true + +require "test_helper" + +# A transport that declares the `timeout:` keyword and records what it was +# handed, so the budget the client computed is observable. +class TimeoutRecordingTransport + attr_reader :timeouts + + def initialize(responses, clock: nil, cost: 0.0) + @responses = responses + @clock = clock + @cost = cost + @timeouts = [] + end + + def call(url:, headers:, body:, timeout: nil) + @timeouts << timeout + @clock&.call(@cost) + response = @responses.length > 1 ? @responses.shift : @responses.first + raise response if response.is_a?(Exception) + + response + end +end + +# The old three-keyword contract, which must keep working untouched. +class LegacyTransport + attr_reader :calls + + def initialize(response) + @response = response + @calls = 0 + end + + def call(url:, headers:, body:) + @calls += 1 + @response + end +end + +class DeadlineTest < Minitest::Test + def questions + { "urgent" => RubyDecisionModel::Questions.noul("Is this urgent?") } + end + + def success_body + JSON.generate( + "answers" => { "urgent" => { "type" => "noul", "noul" => 0.5 } }, + "usage" => { "input_tokens" => 1, "output_tokens" => 1 } + ) + end + + def build_client(transport, **options) + RubyDecisionModel::Client.new( + **{ api_key: "test-key", transport: transport, sleeper: no_sleep, random: -> { 0.0 } }.merge(options) + ) + end + + # --- timeout validation --- + + def test_timeout_must_be_a_finite_positive_number + [0, -1, "5", Float::INFINITY, Float::NAN, []].each do |bad| + assert_raises(RubyDecisionModel::ConfigurationError, "expected #{bad.inspect} to be rejected") do + build_client(FakeTransport.new([]), timeout: bad) + end + end + end + + def test_nil_timeout_is_allowed_and_means_no_per_attempt_limit + client = build_client(FakeTransport.new([]), timeout: nil) + assert_nil client.timeout + end + + # --- the per-attempt timeout is bounded by what is left of the budget --- + + def test_attempt_timeout_is_the_smaller_of_timeout_and_remaining_budget + now = 0.0 + transport = TimeoutRecordingTransport.new([[200, success_body]], clock: ->(cost) { now += cost }) + client = build_client(transport, timeout: 5, clock: -> { now }, retry: { total_timeout: 2.0 }) + + client.ask(state: {}, questions: questions) + + assert_equal [2.0], transport.timeouts + end + + def test_attempt_timeout_is_the_configured_timeout_when_the_budget_is_larger + client_transport = TimeoutRecordingTransport.new([[200, success_body]]) + client = build_client(client_transport, timeout: 5, retry: { total_timeout: 30.0 }) + + client.ask(state: {}, questions: questions) + + assert_equal [5], client_transport.timeouts + end + + def test_a_later_attempt_gets_only_what_the_budget_has_left + now = 0.0 + # Each attempt burns 3s of the 10s budget; the sleeper burns its delay. + transport = TimeoutRecordingTransport.new([[503, "{}"], [200, success_body]], + clock: ->(_c) { now += 3.0 }, cost: 3.0) + client = build_client(transport, timeout: 30, clock: -> { now }, + sleeper: ->(seconds) { now += seconds }, + retry: { total_timeout: 10.0, backoff_initial: 1.0, backoff_jitter: 0.0 }) + + client.ask(state: {}, questions: questions) + + # First attempt: 10s left. Second: 10 - 3 (attempt) - 1 (backoff) = 6s. + assert_equal [10.0, 6.0], transport.timeouts + end + + def test_no_budget_means_the_transport_gets_the_plain_timeout + transport = TimeoutRecordingTransport.new([[200, success_body]]) + client = build_client(transport, timeout: 4, retry: { total_timeout: nil }) + + client.ask(state: {}, questions: questions) + + assert_equal [4], transport.timeouts + end + + def test_an_exhausted_budget_raises_instead_of_starting_an_attempt + now = 0.0 + transport = TimeoutRecordingTransport.new([[200, success_body]]) + client = build_client(transport, clock: -> { now }, retry: { total_timeout: 0.0 }) + + error = assert_raises(RubyDecisionModel::TimeoutError) { client.ask(state: {}, questions: questions) } + + assert_match(/budget/, error.message) + assert_empty transport.timeouts, "no request should have gone out" + end + + # --- backwards compatibility of the transport contract --- + + def test_a_transport_without_a_timeout_keyword_is_still_called_the_old_way + transport = LegacyTransport.new([200, success_body]) + client = build_client(transport, timeout: 5) + + response = client.ask(state: {}, questions: questions) + + assert_equal 1, transport.calls + assert_in_delta 0.5, response["urgent"].noul + end + + def test_a_lambda_transport_with_a_splat_receives_the_timeout + seen = nil + transport = lambda do |url:, headers:, body:, **rest| + seen = rest[:timeout] + [200, success_body] + end + client = build_client(transport, timeout: 7, retry: { total_timeout: nil }) + + client.ask(state: {}, questions: questions) + + assert_equal 7, seen + end + + # --- the budget boundary --- + + def test_a_delay_landing_exactly_on_the_deadline_stops_the_retry + now = 0.0 + sleeper = lambda { |seconds| now += seconds } + transport = FakeTransport.new([[503, "{}"]]) + client = build_client(transport, clock: -> { now }, sleeper: sleeper, + retry: { max_retries: 5, total_timeout: 1.0, backoff_initial: 1.0, + backoff_jitter: 0.0 }) + + # The first backoff is exactly 1.0s, which uses the entire budget, so the + # 503 is returned rather than slept on. + assert_raises(RubyDecisionModel::ApiError) { client.ask(state: {}, questions: questions) } + assert_equal 1, transport.calls.length + assert_in_delta 0.0, now + end + + # --- write timeouts --- + + def test_a_write_timeout_is_treated_as_a_timeout_and_retried + transport = FakeTransport.new([Net::WriteTimeout.new, [200, success_body]]) + client = build_client(transport) + + response = client.ask(state: {}, questions: questions) + + assert_in_delta 0.5, response["urgent"].noul + assert_equal 2, transport.calls.length + end + + def test_an_unretried_write_timeout_surfaces_as_a_timeout_error + transport = FakeTransport.new([Net::WriteTimeout.new]) + client = build_client(transport, retry: { max_retries: 0 }) + + error = assert_raises(RubyDecisionModel::TimeoutError) { client.ask(state: {}, questions: questions) } + assert_instance_of Net::WriteTimeout, error.cause_error + end + + def test_write_timeouts_can_be_switched_off_like_the_other_timeouts + transport = FakeTransport.new([Net::WriteTimeout.new]) + client = build_client(transport, retry: { retry_timeouts: false }) + + assert_raises(RubyDecisionModel::TimeoutError) { client.ask(state: {}, questions: questions) } + assert_equal 1, transport.calls.length + end + + # --- the default transport covers every phase, not just open and read --- + + def test_the_default_transport_sets_all_four_net_http_timeouts + client = RubyDecisionModel::Client.new(api_key: "test-key", timeout: 3) + http = nil + fake_http = Class.new do + attr_accessor :use_ssl, :open_timeout, :ssl_timeout, :read_timeout, :write_timeout + + def request(_request) + Struct.new(:code, :body).new("200", "{}").tap { |r| def r.each_header = {}.each } + end + end + + Net::HTTP.stub(:new, ->(_host, _port) { http = fake_http.new }) do + client.send(:default_transport).call(url: "https://example.com/x", headers: {}, body: "{}", timeout: 3) + end + + assert_equal 3, http.open_timeout + assert_equal 3, http.ssl_timeout + assert_equal 3, http.read_timeout + assert_equal 3, http.write_timeout + end +end