diff --git a/exporter/otlp-http/lib/opentelemetry/exporter/otlp/http/trace_exporter.rb b/exporter/otlp-http/lib/opentelemetry/exporter/otlp/http/trace_exporter.rb index f38e0134db..88cddb0748 100644 --- a/exporter/otlp-http/lib/opentelemetry/exporter/otlp/http/trace_exporter.rb +++ b/exporter/otlp-http/lib/opentelemetry/exporter/otlp/http/trace_exporter.rb @@ -25,6 +25,8 @@ class TraceExporter # rubocop:disable Metrics/ClassLength # Default timeouts in seconds. KEEP_ALIVE_TIMEOUT = 30 RETRY_COUNT = 5 + RESPONSE_BODY_LIMIT = 4_194_304 # 4 MB + private_constant(:KEEP_ALIVE_TIMEOUT, :RETRY_COUNT, :RESPONSE_BODY_LIMIT) ERROR_MESSAGE_INVALID_HEADERS = 'headers must be a String with comma-separated URL Encoded UTF-8 k=v pairs or a Hash' @@ -37,7 +39,8 @@ def initialize(endpoint: nil, ssl_verify_mode: fetch_ssl_verify_mode, headers: OpenTelemetry::Common::Utilities.config_opt('OTEL_EXPORTER_OTLP_TRACES_HEADERS', 'OTEL_EXPORTER_OTLP_HEADERS', default: {}), compression: OpenTelemetry::Common::Utilities.config_opt('OTEL_EXPORTER_OTLP_TRACES_COMPRESSION', 'OTEL_EXPORTER_OTLP_COMPRESSION', default: 'gzip'), - timeout: OpenTelemetry::Common::Utilities.config_opt('OTEL_EXPORTER_OTLP_TRACES_TIMEOUT', 'OTEL_EXPORTER_OTLP_TIMEOUT', default: 10)) + timeout: OpenTelemetry::Common::Utilities.config_opt('OTEL_EXPORTER_OTLP_TRACES_TIMEOUT', 'OTEL_EXPORTER_OTLP_TIMEOUT', default: 10), + metrics_reporter: nil) raise ArgumentError, "unsupported compression key #{compression}" unless compression.nil? || %w[gzip none].include?(compression) @uri = prepare_endpoint(endpoint) @@ -48,6 +51,7 @@ def initialize(endpoint: nil, @headers = prepare_headers(headers) @timeout = timeout.to_f @compression = compression + @metrics_reporter = metrics_reporter || OpenTelemetry::SDK::Trace::Export::MetricsReporter @shutdown = false end @@ -145,34 +149,48 @@ def send_bytes(bytes, timeout:) # rubocop:disable Metrics/MethodLength @http.read_timeout = remaining_timeout @http.write_timeout = remaining_timeout @http.start unless @http.started? - response = @http.request(request) - - case response - when Net::HTTPSuccess - response.body # Read and discard body - SUCCESS - when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests - response.body # Read and discard body - redo if backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1, reason: response.code) - FAILURE - when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway - response.body # Read and discard body - redo if backoff?(retry_count: retry_count += 1, reason: response.code) - FAILURE - when Net::HTTPNotFound - log_request_failure(response.code) - FAILURE - when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError - log_status(response.body) - FAILURE - when Net::HTTPRedirection - @http.finish - handle_redirect(response['location']) - redo if backoff?(retry_after: 0, retry_count: retry_count += 1, reason: response.code) - else - @http.finish - FAILURE + result = nil + should_redo = false + + measure_request_duration do + @http.request(request) do |response| + case response + when Net::HTTPSuccess + drain_body(response) + result = SUCCESS + when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests + drain_body(response) + should_redo = backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1, reason: response.code) + result = FAILURE + when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway + drain_body(response) + should_redo = backoff?(retry_count: retry_count += 1, reason: response.code) + result = FAILURE + when Net::HTTPNotFound + drain_body(response) + log_request_failure(response.code) + result = FAILURE + when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError + body, truncated = read_response_body(response) + log_status(body, truncated: truncated) + @metrics_reporter.add_to_counter('otel.otlp_exporter.failure', labels: { 'reason' => response.code }) + result = FAILURE + when Net::HTTPRedirection + drain_body(response) + @http.finish + handle_redirect(response['location']) + should_redo = backoff?(retry_after: 0, retry_count: retry_count += 1, reason: response.code) + else + drain_body(response) + @http.finish + result = FAILURE + end + end end + + redo if should_redo + + result rescue Net::OpenTimeout, Net::ReadTimeout retry if backoff?(retry_count: retry_count += 1, reason: 'timeout') return FAILURE @@ -207,7 +225,13 @@ def handle_redirect(location) # TODO: figure out destination and reinitialize @http and @path end - def log_status(body) + def log_status(body, truncated: false) + if truncated + OpenTelemetry.handle_error(message: "OTLP exporter received an oversized error response body (truncated at #{RESPONSE_BODY_LIMIT} bytes)") + return + end + return if body.nil? || body.empty? + status = Google::Rpc::Status.decode(body) pool = ::Google::Protobuf::DescriptorPool.generated_pool details = status.details.filter_map do |detail| @@ -219,6 +243,59 @@ def log_status(body) OpenTelemetry.handle_error(exception: e, message: 'unexpected error decoding rpc.Status in OTLP::Exporter#log_status') end + # Drains and discards the body without buffering it, preserving keep-alive. + def drain_body(response) + response.read_body { |_| } # rubocop:disable Lint/EmptyBlock + end + + def read_response_body(response) # rubocop:disable Metrics/MethodLength + return ['', false] if response.nil? + + content_length = response['content-length']&.to_i + if content_length && content_length > RESPONSE_BODY_LIMIT + @http.finish # closes socket without reading any of the oversized body + return ['', true] + end + + body = +'' + truncated = false + + response.read_body do |chunk| + remaining = RESPONSE_BODY_LIMIT - body.bytesize + body << chunk.byteslice(0, remaining) + + if chunk.bytesize > remaining + truncated = true + @http.finish # closes socket, nil's the body or else net/http will attempt to read the rest of the response + break + end + end + + body.force_encoding('UTF-8') + body.scrub! if truncated # truncation may have split a multi-byte character + [body, truncated] + rescue IOError + raise unless truncated # we'll handle this when we know net/http is upset trying to read after http.finish + + [body || '', truncated] + rescue StandardError => e + OpenTelemetry.handle_error(exception: e, message: 'error reading response body') + ['', false] + end + + def measure_request_duration + start = Process.clock_gettime(Process::CLOCK_MONOTONIC) + begin + response = yield + ensure + stop = Process.clock_gettime(Process::CLOCK_MONOTONIC) + duration_ms = 1000.0 * (stop - start) + @metrics_reporter.record_value('otel.otlp_exporter.request_duration', + value: duration_ms, + labels: { 'status' => response&.code || 'unknown' }) + end + end + def log_request_failure(response_code) OpenTelemetry.handle_error(message: "OTLP exporter received http.code=#{response_code} for uri='#{@uri}' in OTLP::Exporter#send_bytes") end diff --git a/exporter/otlp-http/test/opentelemetry/exporter/otlp/http/trace_exporter_test.rb b/exporter/otlp-http/test/opentelemetry/exporter/otlp/http/trace_exporter_test.rb index 95dcf66f0f..9f95bd585a 100644 --- a/exporter/otlp-http/test/opentelemetry/exporter/otlp/http/trace_exporter_test.rb +++ b/exporter/otlp-http/test/opentelemetry/exporter/otlp/http/trace_exporter_test.rb @@ -846,4 +846,232 @@ end end end + + describe 'response body reading' do + let(:exporter) { OpenTelemetry::Exporter::OTLP::HTTP::TraceExporter.new } + let(:span_data) { OpenTelemetry::TestHelpers.create_span_data } + + it 'discards body for successful responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 200, body: 'success body') + + result = exporter.export([span_data]) + + _(result).must_equal(success) + end + + it 'discards body for retryable responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/traces') + .to_return(status: 503, body: 'service unavailable', headers: { 'Retry-After' => '0' }) + .then.to_return(status: 200) + + result = exporter.export([span_data]) + + _(result).must_equal(success) + end + + it 'reads and parses error response body smaller than limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: 'error details'))] + status = ::Google::Rpc::Status.encode(::Google::Rpc::Status.new(code: 3, message: 'invalid argument', details: details)) + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: status) + + result = exporter.export([span_data]) + + _(result).must_equal(export_failure) + _(log_stream.string).must_match(/invalid argument/) + _(log_stream.string).wont_match(/truncated/) + ensure + OpenTelemetry.logger = logger + end + + it 'truncates error response body larger than 4 MB limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + # Create a body larger than 4 MB + large_message = 'x' * 5_000_000 # 5 MB + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: large_message))] + large_status = ::Google::Rpc::Status.new(code: 3, message: 'large error', details: details) + large_body = ::Google::Rpc::Status.encode(large_status) + + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: large_body) + + result = exporter.export([span_data]) + + _(result).must_equal(export_failure) + _(log_stream.string).must_match(/oversized error response body/) + ensure + OpenTelemetry.logger = logger + end + + it 'handles malformed error response body gracefully' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: 'not valid protobuf') + + result = exporter.export([span_data]) + + _(result).must_equal(export_failure) + _(log_stream.string).must_match(/unexpected error decoding rpc.Status/) + ensure + OpenTelemetry.logger = logger + end + end + + describe 'response body reading with real sockets' do + # These tests use a real TCPServer instead of WebMock to verify + # behavior against actual Net::HTTP socket I/O. WebMock's patching + # of Net::HTTP changes how read_body works, so it doesn't exercise + # the same code paths as real Net::HTTP. + + def with_fake_server(response_body_size: nil, status: 400, handler: nil) + server = TCPServer.new('127.0.0.1', 0) + port = server.addr[1] + handler ||= ->(srv, stat) { handle_fake_request(srv, 'X' * response_body_size, stat) } + + server_thread = Thread.new { handler.call(server, status) } + + WebMock::HttpLibAdapters::NetHttpAdapter.disable! + yield port + ensure + WebMock::HttpLibAdapters::NetHttpAdapter.enable! + server&.close + server_thread&.join(2) + end + + def handle_fake_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{body.bytesize}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.write body + client.close + rescue StandardError + # client may disconnect early + end + + # Sends only headers announcing an oversized Content-Length, then closes + # without ever writing a body — proves the exporter short-circuits on + # Content-Length rather than attempting to stream-read the body. + def handle_fake_headers_only_request(server, declared_content_length, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{declared_content_length}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + # Writes body using chunked transfer-encoding with no Content-Length + # header, the common shape of a real collector's error response. + def handle_fake_chunked_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Transfer-Encoding: chunked\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + + chunk_size = 65_536 + offset = 0 + while offset < body.bytesize + chunk = body.byteslice(offset, chunk_size) + client.print "#{chunk.bytesize.to_s(16)}\r\n#{chunk}\r\n" + offset += chunk_size + end + client.print "0\r\n\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + def read_content_length(client) + content_length = 0 + while (line = client.gets) && line != "\r\n" + content_length = line.split(': ', 2).last.to_i if line.start_with?('Content-Length') + end + content_length + end + + it 'limits error response body read to 4 MB against a real HTTP server' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + with_fake_server(response_body_size: 5_000_000) do |port| + exporter = OpenTelemetry::Exporter::OTLP::HTTP::TraceExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_span_data]) + + _(result).must_equal(export_failure) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + _(log_stream.string).wont_match(/read_body called twice/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'skips reading the body when Content-Length already exceeds the limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + handler = ->(server, status) { handle_fake_headers_only_request(server, 5_000_000, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::HTTP::TraceExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_span_data]) + + _(result).must_equal(export_failure) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/error reading response body/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'limits chunked (no Content-Length) response bodies to 4 MB' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + large_body = 'X' * 5_000_000 + handler = ->(server, status) { handle_fake_chunked_request(server, large_body, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::HTTP::TraceExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_span_data]) + + _(result).must_equal(export_failure) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + end + ensure + OpenTelemetry.logger = logger + end + end end diff --git a/exporter/otlp-logs/lib/opentelemetry/exporter/otlp/logs/logs_exporter.rb b/exporter/otlp-logs/lib/opentelemetry/exporter/otlp/logs/logs_exporter.rb index f9e9e670df..7a9e909168 100644 --- a/exporter/otlp-logs/lib/opentelemetry/exporter/otlp/logs/logs_exporter.rb +++ b/exporter/otlp-logs/lib/opentelemetry/exporter/otlp/logs/logs_exporter.rb @@ -31,7 +31,8 @@ class LogsExporter # rubocop:disable Metrics/ClassLength # Default timeouts in seconds. KEEP_ALIVE_TIMEOUT = 30 RETRY_COUNT = 5 - private_constant(:KEEP_ALIVE_TIMEOUT, :RETRY_COUNT) + RESPONSE_BODY_LIMIT = 4_194_304 # 4 MB + private_constant(:KEEP_ALIVE_TIMEOUT, :RETRY_COUNT, :RESPONSE_BODY_LIMIT) ERROR_MESSAGE_INVALID_HEADERS = 'headers must be a String with comma-separated URL Encoded UTF-8 k=v pairs or a Hash' private_constant(:ERROR_MESSAGE_INVALID_HEADERS) @@ -167,37 +168,48 @@ def send_bytes(bytes, timeout:) # rubocop:disable Metrics/CyclomaticComplexity, @http.read_timeout = remaining_timeout @http.write_timeout = remaining_timeout @http.start unless @http.started? - response = @http.request(request) - - case response - when Net::HTTPSuccess - response.body # Read and discard body - SUCCESS - when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests - response.body # Read and discard body - handle_http_error(response) - redo if backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1) - FAILURE - when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway - response.body # Read and discard body - handle_http_error(response) - redo if backoff?(retry_count: retry_count += 1) - FAILURE - when Net::HTTPNotFound - handle_http_error(response) - FAILURE - when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError - log_status(response.body) - FAILURE - when Net::HTTPRedirection - @http.finish - handle_redirect(response['location']) - redo if backoff?(retry_after: 0, retry_count: retry_count += 1) - else - @http.finish - handle_http_error(response) - FAILURE + result = nil + should_redo = false + + @http.request(request) do |response| + case response + when Net::HTTPSuccess + drain_body(response) + result = SUCCESS + when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests + drain_body(response) + handle_http_error(response) + should_redo = backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1) + result = FAILURE + when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway + drain_body(response) + handle_http_error(response) + should_redo = backoff?(retry_count: retry_count += 1) + result = FAILURE + when Net::HTTPNotFound + drain_body(response) + handle_http_error(response) + result = FAILURE + when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError + body, truncated = read_response_body(response) + log_status(body, truncated: truncated) + result = FAILURE + when Net::HTTPRedirection + drain_body(response) + @http.finish + handle_redirect(response['location']) + should_redo = backoff?(retry_after: 0, retry_count: retry_count += 1) + else + drain_body(response) + @http.finish + handle_http_error(response) + result = FAILURE + end end + + redo if should_redo + + result rescue Net::OpenTimeout, Net::ReadTimeout => e OpenTelemetry.handle_error(exception: e) retry if backoff?(retry_count: retry_count += 1) @@ -238,7 +250,13 @@ def handle_redirect(location) # TODO: figure out destination and reinitialize @http and @path end - def log_status(body) + def log_status(body, truncated: false) + if truncated + OpenTelemetry.handle_error(message: "OTLP logs exporter received an oversized error response body (truncated at #{RESPONSE_BODY_LIMIT} bytes)") + return + end + return if body.nil? || body.empty? + status = Google::Rpc::Status.decode(body) pool = ::Google::Protobuf::DescriptorPool.generated_pool details = status.details.filter_map do |detail| @@ -250,6 +268,46 @@ def log_status(body) OpenTelemetry.handle_error(exception: e, message: 'unexpected error decoding rpc.Status in OTLP::Exporter#log_status') end + # Drains and discards the body without buffering it, preserving keep-alive. + def drain_body(response) + response.read_body { |_| } # rubocop:disable Lint/EmptyBlock + end + + def read_response_body(response) # rubocop:disable Metrics/CyclomaticComplexity, Metrics/MethodLength, Metrics/PerceivedComplexity + return ['', false] if response.nil? + + content_length = response['content-length']&.to_i + if content_length && content_length > RESPONSE_BODY_LIMIT + @http.finish # closes socket without reading any of the oversized body + return ['', true] + end + + body = +'' + truncated = false + + response.read_body do |chunk| + remaining = RESPONSE_BODY_LIMIT - body.bytesize + body << chunk.byteslice(0, remaining) + + if chunk.bytesize > remaining + truncated = true + @http.finish # closes socket, nil's the body or else net/http will attempt to read the rest of the response + break + end + end + + body.force_encoding('UTF-8') + body.scrub! if truncated # truncation may have split a multi-byte character + [body, truncated] + rescue IOError + raise unless truncated # we'll handle this when we know net/http is upset trying to read after http.finish + + [body || '', truncated] + rescue StandardError => e + OpenTelemetry.handle_error(exception: e, message: 'error reading response body') + ['', false] + end + def backoff?(retry_count:, retry_after: nil) # rubocop:disable Metrics/CyclomaticComplexity return false if retry_count > RETRY_COUNT diff --git a/exporter/otlp-logs/test/opentelemetry/exporter/otlp/logs_exporter_test.rb b/exporter/otlp-logs/test/opentelemetry/exporter/otlp/logs_exporter_test.rb index 4d37bef293..5bba5975ec 100644 --- a/exporter/otlp-logs/test/opentelemetry/exporter/otlp/logs_exporter_test.rb +++ b/exporter/otlp-logs/test/opentelemetry/exporter/otlp/logs_exporter_test.rb @@ -977,4 +977,232 @@ end end end + + describe 'response body reading' do + let(:exporter) { OpenTelemetry::Exporter::OTLP::Logs::LogsExporter.new } + let(:log_record_data) { OpenTelemetry::TestHelpers.create_log_record_data } + + it 'discards body for successful responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/logs').to_return(status: 200, body: 'success body') + + result = exporter.export([log_record_data]) + + _(result).must_equal(SUCCESS) + end + + it 'discards body for retryable responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/logs') + .to_return(status: 503, body: 'service unavailable', headers: { 'Retry-After' => '0' }) + .then.to_return(status: 200) + + result = exporter.export([log_record_data]) + + _(result).must_equal(SUCCESS) + end + + it 'reads and parses error response body smaller than limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: 'error details'))] + status = ::Google::Rpc::Status.encode(::Google::Rpc::Status.new(code: 3, message: 'invalid argument', details: details)) + stub_request(:post, 'http://localhost:4318/v1/logs').to_return(status: 400, body: status) + + result = exporter.export([log_record_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/invalid argument/) + _(log_stream.string).wont_match(/truncated/) + ensure + OpenTelemetry.logger = logger + end + + it 'truncates error response body larger than 4 MB limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + # Create a body larger than 4 MB + large_message = 'x' * 5_000_000 # 5 MB + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: large_message))] + large_status = ::Google::Rpc::Status.new(code: 3, message: 'large error', details: details) + large_body = ::Google::Rpc::Status.encode(large_status) + + stub_request(:post, 'http://localhost:4318/v1/logs').to_return(status: 400, body: large_body) + + result = exporter.export([log_record_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + ensure + OpenTelemetry.logger = logger + end + + it 'handles malformed error response body gracefully' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + stub_request(:post, 'http://localhost:4318/v1/logs').to_return(status: 400, body: 'not valid protobuf') + + result = exporter.export([log_record_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/unexpected error decoding rpc.Status/) + ensure + OpenTelemetry.logger = logger + end + end + + describe 'response body reading with real sockets' do + # These tests use a real TCPServer instead of WebMock to verify + # behavior against actual Net::HTTP socket I/O. WebMock's patching + # of Net::HTTP changes how read_body works, so it doesn't exercise + # the same code paths as real Net::HTTP. + + def with_fake_server(response_body_size: nil, status: 400, handler: nil) + server = TCPServer.new('127.0.0.1', 0) + port = server.addr[1] + handler ||= ->(srv, stat) { handle_fake_request(srv, 'X' * response_body_size, stat) } + + server_thread = Thread.new { handler.call(server, status) } + + WebMock::HttpLibAdapters::NetHttpAdapter.disable! + yield port + ensure + WebMock::HttpLibAdapters::NetHttpAdapter.enable! + server&.close + server_thread&.join(2) + end + + def handle_fake_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{body.bytesize}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.write body + client.close + rescue StandardError + # client may disconnect early + end + + # Sends only headers announcing an oversized Content-Length, then closes + # without ever writing a body — proves the exporter short-circuits on + # Content-Length rather than attempting to stream-read the body. + def handle_fake_headers_only_request(server, declared_content_length, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{declared_content_length}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + # Writes body using chunked transfer-encoding with no Content-Length + # header, the common shape of a real collector's error response. + def handle_fake_chunked_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Transfer-Encoding: chunked\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + + chunk_size = 65_536 + offset = 0 + while offset < body.bytesize + chunk = body.byteslice(offset, chunk_size) + client.print "#{chunk.bytesize.to_s(16)}\r\n#{chunk}\r\n" + offset += chunk_size + end + client.print "0\r\n\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + def read_content_length(client) + content_length = 0 + while (line = client.gets) && line != "\r\n" + content_length = line.split(': ', 2).last.to_i if line.start_with?('Content-Length') + end + content_length + end + + it 'limits error response body read to 4 MB against a real HTTP server' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + with_fake_server(response_body_size: 5_000_000) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Logs::LogsExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/logs" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_log_record_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + _(log_stream.string).wont_match(/read_body called twice/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'skips reading the body when Content-Length already exceeds the limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + handler = ->(server, status) { handle_fake_headers_only_request(server, 5_000_000, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Logs::LogsExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/logs" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_log_record_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/error reading response body/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'limits chunked (no Content-Length) response bodies to 4 MB' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + large_body = 'X' * 5_000_000 + handler = ->(server, status) { handle_fake_chunked_request(server, large_body, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Logs::LogsExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/logs" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_log_record_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + end + ensure + OpenTelemetry.logger = logger + end + end end diff --git a/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/metrics_exporter.rb b/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/metrics_exporter.rb index 85432c9163..888b275752 100644 --- a/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/metrics_exporter.rb +++ b/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/metrics_exporter.rb @@ -26,7 +26,7 @@ module Exporter module OTLP module Metrics # An OpenTelemetry metrics exporter that sends metrics over HTTP as Protobuf encoded OTLP ExportMetricsServiceRequest. - class MetricsExporter < ::OpenTelemetry::SDK::Metrics::Export::MetricReader + class MetricsExporter < ::OpenTelemetry::SDK::Metrics::Export::MetricReader # rubocop:disable Metrics/ClassLength include Util attr_reader :metric_snapshots @@ -93,7 +93,7 @@ def export(metrics, timeout: nil) end end - def send_bytes(bytes, timeout:) + def send_bytes(bytes, timeout:) # rubocop:disable Metrics/MethodLength return FAILURE if bytes.nil? request = Net::HTTP::Post.new(@path) @@ -121,37 +121,49 @@ def send_bytes(bytes, timeout:) @http.read_timeout = remaining_timeout @http.write_timeout = remaining_timeout @http.start unless @http.started? - response = @http.request(request) - case response - when Net::HTTPSuccess - response.body # Read and discard body - SUCCESS - when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests - response.body # Read and discard body - redo if backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1, reason: response.code) - OpenTelemetry.logger.warn('Net::HTTPServiceUnavailable/Net::HTTPTooManyRequests in MetricsExporter#send_bytes') - FAILURE - when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway - response.body # Read and discard body - redo if backoff?(retry_count: retry_count += 1, reason: response.code) - OpenTelemetry.logger.warn('Net::HTTPRequestTimeOut/Net::HTTPGatewayTimeOut/Net::HTTPBadGateway in MetricsExporter#send_bytes') - FAILURE - when Net::HTTPNotFound - OpenTelemetry.handle_error(message: "OTLP metrics_exporter received http.code=404 for uri: '#{@path}'") - FAILURE - when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError - log_status(response.body) - OpenTelemetry.logger.warn('Net::HTTPBadRequest/Net::HTTPClientError/Net::HTTPServerError in MetricsExporter#send_bytes') - FAILURE - when Net::HTTPRedirection - @http.finish - handle_redirect(response['location']) - redo if backoff?(retry_after: 0, retry_count: retry_count += 1, reason: response.code) - else - @http.finish - OpenTelemetry.logger.warn("Unexpected error in OTLP::MetricsExporter#send_bytes - #{response.message}") - FAILURE + result = nil + should_redo = false + + @http.request(request) do |response| + case response + when Net::HTTPSuccess + drain_body(response) + result = SUCCESS + when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests + drain_body(response) + should_redo = backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1, reason: response.code) + OpenTelemetry.logger.warn('Net::HTTPServiceUnavailable/Net::HTTPTooManyRequests in MetricsExporter#send_bytes') + result = FAILURE + when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway + drain_body(response) + should_redo = backoff?(retry_count: retry_count += 1, reason: response.code) + OpenTelemetry.logger.warn('Net::HTTPRequestTimeOut/Net::HTTPGatewayTimeOut/Net::HTTPBadGateway in MetricsExporter#send_bytes') + result = FAILURE + when Net::HTTPNotFound + drain_body(response) + OpenTelemetry.handle_error(message: "OTLP metrics_exporter received http.code=404 for uri: '#{@path}'") + result = FAILURE + when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError + body, truncated = read_response_body(response) + log_status(body, truncated: truncated) + OpenTelemetry.logger.warn('Net::HTTPBadRequest/Net::HTTPClientError/Net::HTTPServerError in MetricsExporter#send_bytes') + result = FAILURE + when Net::HTTPRedirection + drain_body(response) + @http.finish + handle_redirect(response['location']) + should_redo = backoff?(retry_after: 0, retry_count: retry_count += 1, reason: response.code) + else + drain_body(response) + @http.finish + OpenTelemetry.logger.warn("Unexpected error in OTLP::MetricsExporter#send_bytes - #{response.message}") + result = FAILURE + end end + + redo if should_redo + + result rescue Net::OpenTimeout, Net::ReadTimeout retry if backoff?(retry_count: retry_count += 1, reason: 'timeout') OpenTelemetry.logger.warn('Net::OpenTimeout/Net::ReadTimeout in MetricsExporter#send_bytes') diff --git a/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/util.rb b/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/util.rb index 406a1afa10..44d0e04eae 100644 --- a/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/util.rb +++ b/exporter/otlp-metrics/lib/opentelemetry/exporter/otlp/metrics/util.rb @@ -9,11 +9,13 @@ module Exporter module OTLP module Metrics # Util module provide essential functionality for exporter - module Util + module Util # rubocop:disable Metrics/ModuleLength KEEP_ALIVE_TIMEOUT = 30 RETRY_COUNT = 5 + RESPONSE_BODY_LIMIT = 4_194_304 # 4 MB ERROR_MESSAGE_INVALID_HEADERS = 'headers must be a String with comma-separated URL Encoded UTF-8 k=v pairs or a Hash' DEFAULT_USER_AGENT = "OTel-OTLP-MetricsExporter-Ruby/#{OpenTelemetry::Exporter::OTLP::Metrics::VERSION} Ruby/#{RUBY_VERSION} (#{RUBY_PLATFORM}; #{RUBY_ENGINE}/#{RUBY_ENGINE_VERSION})".freeze + private_constant(:RESPONSE_BODY_LIMIT) def http_connection(uri, ssl_verify_mode, certificate_file, client_certificate_file, client_key_file) http = Net::HTTP.new(uri.hostname, uri.port) @@ -110,7 +112,13 @@ def backoff?(retry_count:, reason:, retry_after: nil) true end - def log_status(body) + def log_status(body, truncated: false) + if truncated + OpenTelemetry.handle_error(message: "OTLP metrics_exporter received an oversized error response body (truncated at #{RESPONSE_BODY_LIMIT} bytes)") + return + end + return if body.nil? || body.empty? + status = Google::Rpc::Status.decode(body) pool = ::Google::Protobuf::DescriptorPool.generated_pool details = status.details.filter_map do |detail| @@ -123,6 +131,47 @@ def log_status(body) OpenTelemetry.handle_error(exception: e, message: 'unexpected error decoding rpc.Status in OTLP::MetricsExporter#log_status') end + # Drains and discards the body without buffering it, preserving keep-alive. + def drain_body(response) + response.read_body { |_| } # rubocop:disable Lint/EmptyBlock + end + + def read_response_body(response) + return ['', false] if response.nil? + + content_length = response['content-length']&.to_i + if content_length && content_length > RESPONSE_BODY_LIMIT + @http.finish # closes socket without reading any of the oversized body + return ['', true] + end + + # Stream read with 4 MB limit + body = +'' + truncated = false + + response.read_body do |chunk| + remaining = RESPONSE_BODY_LIMIT - body.bytesize + body << chunk.byteslice(0, remaining) + + if chunk.bytesize > remaining + truncated = true + @http.finish # closes socket, nil's the body or else net/http will attempt to read the rest of the response + break + end + end + + body.force_encoding('UTF-8') + body.scrub! if truncated # truncation may have split a multi-byte character + [body, truncated] + rescue IOError + raise unless truncated # we'll handle this when we know net/http is upset trying to read after http.finish + + [body || '', truncated] + rescue StandardError => e + OpenTelemetry.handle_error(exception: e, message: 'error reading response body') + ['', false] + end + def handle_redirect(location); end end end diff --git a/exporter/otlp-metrics/test/opentelemetry/exporter/otlp/metrics/metrics_exporter_test.rb b/exporter/otlp-metrics/test/opentelemetry/exporter/otlp/metrics/metrics_exporter_test.rb index 1c57517bd1..b2253da71f 100644 --- a/exporter/otlp-metrics/test/opentelemetry/exporter/otlp/metrics/metrics_exporter_test.rb +++ b/exporter/otlp-metrics/test/opentelemetry/exporter/otlp/metrics/metrics_exporter_test.rb @@ -931,4 +931,232 @@ end end end + + describe 'response body reading' do + let(:exporter) { OpenTelemetry::Exporter::OTLP::Metrics::MetricsExporter.new } + let(:metric_data) { create_metrics_data } + + it 'discards body for successful responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/metrics').to_return(status: 200, body: 'success body') + + result = exporter.export([metric_data]) + + _(result).must_equal(METRICS_SUCCESS) + end + + it 'discards body for retryable responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/metrics') + .to_return(status: 503, body: 'service unavailable', headers: { 'Retry-After' => '0' }) + .then.to_return(status: 200) + + result = exporter.export([metric_data]) + + _(result).must_equal(METRICS_SUCCESS) + end + + it 'reads and parses error response body smaller than limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: 'error details'))] + status = ::Google::Rpc::Status.encode(::Google::Rpc::Status.new(code: 3, message: 'invalid argument', details: details)) + stub_request(:post, 'http://localhost:4318/v1/metrics').to_return(status: 400, body: status) + + result = exporter.export([metric_data]) + + _(result).must_equal(METRICS_FAILURE) + _(log_stream.string).must_match(/invalid argument/) + _(log_stream.string).wont_match(/truncated/) + ensure + OpenTelemetry.logger = logger + end + + it 'truncates error response body larger than 4 MB limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + # Create a body larger than 4 MB + large_message = 'x' * 5_000_000 # 5 MB + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: large_message))] + large_status = ::Google::Rpc::Status.new(code: 3, message: 'large error', details: details) + large_body = ::Google::Rpc::Status.encode(large_status) + + stub_request(:post, 'http://localhost:4318/v1/metrics').to_return(status: 400, body: large_body) + + result = exporter.export([metric_data]) + + _(result).must_equal(METRICS_FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + ensure + OpenTelemetry.logger = logger + end + + it 'handles malformed error response body gracefully' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + stub_request(:post, 'http://localhost:4318/v1/metrics').to_return(status: 400, body: 'not valid protobuf') + + result = exporter.export([metric_data]) + + _(result).must_equal(METRICS_FAILURE) + _(log_stream.string).must_match(/unexpected error decoding rpc.Status/) + ensure + OpenTelemetry.logger = logger + end + end + + describe 'response body reading with real sockets' do + # These tests use a real TCPServer instead of WebMock to verify + # behavior against actual Net::HTTP socket I/O. WebMock's patching + # of Net::HTTP changes how read_body works, so it doesn't exercise + # the same code paths as real Net::HTTP. + + def with_fake_server(response_body_size: nil, status: 400, handler: nil) + server = TCPServer.new('127.0.0.1', 0) + port = server.addr[1] + handler ||= ->(srv, stat) { handle_fake_request(srv, 'X' * response_body_size, stat) } + + server_thread = Thread.new { handler.call(server, status) } + + WebMock::HttpLibAdapters::NetHttpAdapter.disable! + yield port + ensure + WebMock::HttpLibAdapters::NetHttpAdapter.enable! + server&.close + server_thread&.join(2) + end + + def handle_fake_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{body.bytesize}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.write body + client.close + rescue StandardError + # client may disconnect early + end + + # Sends only headers announcing an oversized Content-Length, then closes + # without ever writing a body — proves the exporter short-circuits on + # Content-Length rather than attempting to stream-read the body. + def handle_fake_headers_only_request(server, declared_content_length, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{declared_content_length}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + # Writes body using chunked transfer-encoding with no Content-Length + # header, the common shape of a real collector's error response. + def handle_fake_chunked_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Transfer-Encoding: chunked\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + + chunk_size = 65_536 + offset = 0 + while offset < body.bytesize + chunk = body.byteslice(offset, chunk_size) + client.print "#{chunk.bytesize.to_s(16)}\r\n#{chunk}\r\n" + offset += chunk_size + end + client.print "0\r\n\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + def read_content_length(client) + content_length = 0 + while (line = client.gets) && line != "\r\n" + content_length = line.split(': ', 2).last.to_i if line.start_with?('Content-Length') + end + content_length + end + + it 'limits error response body read to 4 MB against a real HTTP server' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + with_fake_server(response_body_size: 5_000_000) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Metrics::MetricsExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/metrics" + ) + result = exporter.export([create_metrics_data]) + + _(result).must_equal(METRICS_FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + _(log_stream.string).wont_match(/read_body called twice/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'skips reading the body when Content-Length already exceeds the limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + handler = ->(server, status) { handle_fake_headers_only_request(server, 5_000_000, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Metrics::MetricsExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/metrics" + ) + result = exporter.export([create_metrics_data]) + + _(result).must_equal(METRICS_FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/error reading response body/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'limits chunked (no Content-Length) response bodies to 4 MB' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + large_body = 'X' * 5_000_000 + handler = ->(server, status) { handle_fake_chunked_request(server, large_body, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Metrics::MetricsExporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/metrics" + ) + result = exporter.export([create_metrics_data]) + + _(result).must_equal(METRICS_FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + end + ensure + OpenTelemetry.logger = logger + end + end end diff --git a/exporter/otlp/lib/opentelemetry/exporter/otlp/exporter.rb b/exporter/otlp/lib/opentelemetry/exporter/otlp/exporter.rb index 68f6ceae18..8fdd3bb7cf 100644 --- a/exporter/otlp/lib/opentelemetry/exporter/otlp/exporter.rb +++ b/exporter/otlp/lib/opentelemetry/exporter/otlp/exporter.rb @@ -28,7 +28,8 @@ class Exporter # rubocop:disable Metrics/ClassLength # Default timeouts in seconds. KEEP_ALIVE_TIMEOUT = 30 RETRY_COUNT = 5 - private_constant(:KEEP_ALIVE_TIMEOUT, :RETRY_COUNT) + RESPONSE_BODY_LIMIT = 4_194_304 # 4 MB + private_constant(:KEEP_ALIVE_TIMEOUT, :RETRY_COUNT, :RESPONSE_BODY_LIMIT) ERROR_MESSAGE_INVALID_HEADERS = 'headers must be a String with comma-separated URL Encoded UTF-8 k=v pairs or a Hash' private_constant(:ERROR_MESSAGE_INVALID_HEADERS) @@ -176,36 +177,49 @@ def send_bytes(bytes, timeout:) # rubocop:disable Metrics/CyclomaticComplexity, @http.read_timeout = remaining_timeout @http.write_timeout = remaining_timeout @http.start unless @http.started? - response = measure_request_duration { @http.request(request) } - - case response - when Net::HTTPSuccess - response.body # Read and discard body - SUCCESS - when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests - response.body # Read and discard body - redo if backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1, reason: response.code) - FAILURE - when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway - response.body # Read and discard body - redo if backoff?(retry_count: retry_count += 1, reason: response.code) - FAILURE - when Net::HTTPNotFound - log_request_failure(response.code) - FAILURE - when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError - log_status(response.body) - @metrics_reporter.add_to_counter('otel.otlp_exporter.failure', labels: { 'reason' => response.code }) - FAILURE - when Net::HTTPRedirection - @http.finish - handle_redirect(response['location']) - redo if backoff?(retry_after: 0, retry_count: retry_count += 1, reason: response.code) - else - @http.finish - log_request_failure(response.code) - FAILURE + result = nil + should_redo = false + + measure_request_duration do + @http.request(request) do |response| + case response + when Net::HTTPSuccess + drain_body(response) + result = SUCCESS + when Net::HTTPServiceUnavailable, Net::HTTPTooManyRequests + drain_body(response) + should_redo = backoff?(retry_after: response['Retry-After'], retry_count: retry_count += 1, reason: response.code) + result = FAILURE + when Net::HTTPRequestTimeOut, Net::HTTPGatewayTimeOut, Net::HTTPBadGateway + drain_body(response) + should_redo = backoff?(retry_count: retry_count += 1, reason: response.code) + result = FAILURE + when Net::HTTPNotFound + drain_body(response) + log_request_failure(response.code) + result = FAILURE + when Net::HTTPBadRequest, Net::HTTPClientError, Net::HTTPServerError + body, truncated = read_response_body(response) + log_status(body, truncated: truncated) + @metrics_reporter.add_to_counter('otel.otlp_exporter.failure', labels: { 'reason' => response.code }) + result = FAILURE + when Net::HTTPRedirection + drain_body(response) + @http.finish + handle_redirect(response['location']) + should_redo = backoff?(retry_after: 0, retry_count: retry_count += 1, reason: response.code) + else + drain_body(response) + @http.finish + log_request_failure(response.code) + result = FAILURE + end + end end + + redo if should_redo + + result rescue Net::OpenTimeout, Net::ReadTimeout retry if backoff?(retry_count: retry_count += 1, reason: 'timeout') return FAILURE @@ -241,7 +255,13 @@ def handle_redirect(location) # TODO: figure out destination and reinitialize @http and @path end - def log_status(body) + def log_status(body, truncated: false) + if truncated + OpenTelemetry.handle_error(message: "OTLP exporter received an oversized error response body (truncated at #{RESPONSE_BODY_LIMIT} bytes) for uri=#{@uri}") + return + end + return if body.nil? || body.empty? + status = Google::Rpc::Status.decode(body) pool = ::Google::Protobuf::DescriptorPool.generated_pool details = status.details.filter_map do |detail| @@ -253,6 +273,51 @@ def log_status(body) OpenTelemetry.handle_error(exception: e, message: 'unexpected error decoding rpc.Status in OTLP::Exporter#log_status') end + # Drains and discards the body without buffering it, preserving keep-alive. + def drain_body(response) + response.read_body { |_| } # rubocop:disable Lint/EmptyBlock + end + + def read_response_body(response) + return ['', false] if response.nil? + + content_length = response['content-length']&.to_i + if content_length && content_length > RESPONSE_BODY_LIMIT + @http.finish + return ['', true] + end + + body, truncated = collect_response_chunks(response) + body.force_encoding('UTF-8') + body.scrub! if truncated # truncation may have split a multi-byte character + [body, truncated] + rescue StandardError => e + OpenTelemetry.handle_error(exception: e, message: 'error reading response body') + ['', false] + end + + def collect_response_chunks(response) + body = +'' + truncated = false + + response.read_body do |chunk| + remaining = RESPONSE_BODY_LIMIT - body.bytesize + body << chunk.byteslice(0, remaining) + + if chunk.bytesize > remaining + truncated = true + @http.finish # closes socket, nil's the body or else net/http will attempt to read the rest of the response + break + end + end + + [body, truncated] + rescue IOError + raise unless truncated # we'll handle this when we know net/http is upset trying to read after http.finish + + [body || '', truncated] + end + def log_request_failure(response_code) OpenTelemetry.handle_error(message: "OTLP exporter received http.code=#{response_code} for uri='#{@uri}' in OTLP::Exporter#send_bytes") @metrics_reporter.add_to_counter('otel.otlp_exporter.failure', labels: { 'reason' => response_code }) diff --git a/exporter/otlp/test/opentelemetry/exporter/otlp/exporter_test.rb b/exporter/otlp/test/opentelemetry/exporter/otlp/exporter_test.rb index e5a1dc5e0a..f06b5d99b0 100644 --- a/exporter/otlp/test/opentelemetry/exporter/otlp/exporter_test.rb +++ b/exporter/otlp/test/opentelemetry/exporter/otlp/exporter_test.rb @@ -1041,4 +1041,323 @@ def create_link(span_context) OpenTelemetry::Trace::Link.new(span_context, { 'link-attribute' => 'link-value' }) end end + + describe 'response body reading with real sockets' do + # These tests use a real TCPServer instead of WebMock to verify + # behavior against actual Net::HTTP socket I/O. WebMock's patching + # of Net::HTTP changes how read_body works — even with + # allow_net_connect!, responses go through WebMock's adapter which + # doesn't exercise the same code paths as real Net::HTTP. + + def with_fake_server(response_body_size: nil, status: 400, handler: nil) + server = TCPServer.new('127.0.0.1', 0) + port = server.addr[1] + handler ||= ->(srv, stat) { handle_fake_request(srv, 'X' * response_body_size, stat) } + + server_thread = Thread.new { handler.call(server, status) } + + # Fully disable WebMock's Net::HTTP adapter so we get real + # socket behavior, not WebMock's patched read_body. + WebMock::HttpLibAdapters::NetHttpAdapter.disable! + yield port + ensure + WebMock::HttpLibAdapters::NetHttpAdapter.enable! + server&.close + server_thread&.join(2) + end + + def handle_fake_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{body.bytesize}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.write body + client.close + rescue StandardError + # client may disconnect early + end + + # Sends only headers announcing an oversized Content-Length, then closes + # without ever writing a body. Proves the exporter short-circuits on + # Content-Length rather than attempting to stream-read the body: if it + # tried, it would hit an EOF instead of the short-circuit's oversized message. + def handle_fake_headers_only_request(server, declared_content_length, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Content-Length: #{declared_content_length}\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + # Writes body using chunked transfer-encoding with no Content-Length + # header, the common shape of a real collector's error response. + def handle_fake_chunked_request(server, body, status) + client = server.accept + content_length = read_content_length(client) + client.read(content_length) if content_length > 0 + + client.print "HTTP/1.1 #{status} Bad Request\r\n" + client.print "Content-Type: application/x-protobuf\r\n" + client.print "Transfer-Encoding: chunked\r\n" + client.print "Connection: close\r\n" + client.print "\r\n" + + chunk_size = 65_536 + offset = 0 + while offset < body.bytesize + chunk = body.byteslice(offset, chunk_size) + client.print "#{chunk.bytesize.to_s(16)}\r\n#{chunk}\r\n" + offset += chunk_size + end + client.print "0\r\n\r\n" + client.close + rescue StandardError + # client may disconnect early + end + + def read_content_length(client) + content_length = 0 + while (line = client.gets) && line != "\r\n" + content_length = line.split(': ', 2).last.to_i if line.start_with?('Content-Length') + end + content_length + end + + it 'limits error response body read to 4 MB against a real HTTP server' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + with_fake_server(response_body_size: 5_000_000) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Exporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + span_data = OpenTelemetry::TestHelpers.create_span_data + result = exporter.export([span_data]) + + _(result).must_equal(FAILURE) + # The chunked reader should cap at 4 MB and log a clear oversized message. + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + # And it should work without hitting "read_body called twice", + # which would mean the body was already fully read into memory. + _(log_stream.string).wont_match(/read_body called twice/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'skips reading the body when Content-Length already exceeds the limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + handler = ->(server, status) { handle_fake_headers_only_request(server, 5_000_000, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Exporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_span_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/error reading response body/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'limits chunked (no Content-Length) response bodies to 4 MB' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + large_body = 'X' * 5_000_000 + handler = ->(server, status) { handle_fake_chunked_request(server, large_body, status) } + with_fake_server(handler: handler) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Exporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + result = exporter.export([OpenTelemetry::TestHelpers.create_span_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + end + ensure + OpenTelemetry.logger = logger + end + + it 'does not buffer beyond the limit internally' do + captured_body_size = nil + internal_body = :not_checked + + spy = Module.new do + define_method(:read_response_body) do |response| + super(response).tap do |result_body, _truncated| + captured_body_size = result_body.bytesize + internal_body = response.body + end + end + end + + limit = OpenTelemetry::Exporter::OTLP::Exporter.const_get(:RESPONSE_BODY_LIMIT) + + with_fake_server(response_body_size: 10_000_000) do |port| + exporter = OpenTelemetry::Exporter::OTLP::Exporter.new( + endpoint: "http://127.0.0.1:#{port}/v1/traces" + ) + exporter.singleton_class.prepend(spy) + + OpenTelemetry.logger = ::Logger.new(File::NULL) + exporter.export([OpenTelemetry::TestHelpers.create_span_data]) + + # read_response_body must cap the returned body at the limit + _(captured_body_size).wont_be_nil + _(captured_body_size).must_be :<=, limit + + # Net::HTTP must NOT have the full 10 MB response buffered internally. + # After @http.finish closes the socket mid-read, response.body is nil + # (block-form read_body doesn't accumulate into @body). + if internal_body.is_a?(String) + _(internal_body.bytesize).must_be :<=, limit, + "Net::HTTP buffered #{internal_body.bytesize} bytes internally, exceeding the #{limit} byte limit" + end + end + end + end + + describe 'response body reading' do + let(:exporter) { OpenTelemetry::Exporter::OTLP::Exporter.new } + let(:span_data) { OpenTelemetry::TestHelpers.create_span_data } + + it 'discards body for successful responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 200, body: 'success body') + + result = exporter.export([span_data]) + + _(result).must_equal(SUCCESS) + end + + it 'discards body for retryable responses without reading into memory' do + stub_request(:post, 'http://localhost:4318/v1/traces') + .to_return(status: 503, body: 'service unavailable', headers: { 'Retry-After' => '0' }) + .then.to_return(status: 200) + + result = exporter.export([span_data]) + + _(result).must_equal(SUCCESS) + end + + it 'reads and parses error response body smaller than limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: 'error details'))] + status = ::Google::Rpc::Status.encode(::Google::Rpc::Status.new(code: 3, message: 'invalid argument', details: details)) + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: status) + + result = exporter.export([span_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/invalid argument/) + _(log_stream.string).wont_match(/truncated/) + ensure + OpenTelemetry.logger = logger + end + + it 'truncates error response body larger than 4 MB limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + # Create a body larger than 4 MB + large_message = 'x' * 5_000_000 # 5 MB + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: large_message))] + large_status = ::Google::Rpc::Status.new(code: 3, message: 'large error', details: details) + large_body = ::Google::Rpc::Status.encode(large_status) + + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: large_body) + + result = exporter.export([span_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + _(log_stream.string).wont_match(/unexpected error decoding/) + ensure + OpenTelemetry.logger = logger + end + + it 'handles error response body at exactly 4 MB limit' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + # Create a body at exactly 4 MB + exact_size_message = 'y' * (4_194_304 - 100) # Account for protobuf overhead + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: exact_size_message))] + exact_status = ::Google::Rpc::Status.new(code: 3, message: 'exact size', details: details) + exact_body = ::Google::Rpc::Status.encode(exact_status) + + # Skip if encoded body is still larger than 4 MB due to protobuf overhead + skip 'Protobuf overhead makes this test impractical' if exact_body.bytesize > 4_194_304 + + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: exact_body) + + result = exporter.export([span_data]) + + _(result).must_equal(FAILURE) + ensure # rubocop:disable Minitest/SkipEnsure + OpenTelemetry.logger = logger + end + + it 'handles malformed error response body gracefully' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: 'not valid protobuf') + + result = exporter.export([span_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/unexpected error decoding rpc.Status/) + ensure + OpenTelemetry.logger = logger + end + + it 'handles truncated protobuf in error response' do + log_stream = StringIO.new + logger = OpenTelemetry.logger + OpenTelemetry.logger = ::Logger.new(log_stream) + + # Create a large protobuf that will be truncated, making it invalid + large_message = 'z' * 5_000_000 + details = [::Google::Protobuf::Any.pack(::Google::Protobuf::StringValue.new(value: large_message))] + large_status = ::Google::Rpc::Status.new(code: 3, message: 'truncation test', details: details) + large_body = ::Google::Rpc::Status.encode(large_status) + + stub_request(:post, 'http://localhost:4318/v1/traces').to_return(status: 400, body: large_body) + + result = exporter.export([span_data]) + + _(result).must_equal(FAILURE) + _(log_stream.string).must_match(/oversized error response body/) + ensure + OpenTelemetry.logger = logger + end + end end