From 18758dab99a6fb86759053915067bfaef2c178b8 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 19:08:19 +1200 Subject: [PATCH 1/8] Probe TLS readability through the SSL socket --- lib/io/stream/buffered.rb | 29 +++++++++++++- lib/io/stream/readable.rb | 2 +- releases.md | 4 ++ test/io/stream/tls_readable.rb | 71 ++++++++++++++++++++++++++++++++++ 4 files changed, 104 insertions(+), 2 deletions(-) create mode 100644 test/io/stream/tls_readable.rb diff --git a/lib/io/stream/buffered.rb b/lib/io/stream/buffered.rb index 645dadd..bdb6c23 100644 --- a/lib/io/stream/buffered.rb +++ b/lib/io/stream/buffered.rb @@ -91,7 +91,34 @@ def close_write # Check if the stream is readable. # @returns [Boolean] True if the stream is readable. def readable? - super && @io.readable? + return false unless super + return true unless @read_buffer.empty? + + # Probe through the wrapped IO rather than its underlying descriptor. This is + # important for layered transports such as TLS, where encrypted data on the + # socket may decode to an EOF (close_notify). Preserve any byte consumed by + # the probe in the stream's read buffer. + result = @io.read_nonblock(1, @read_buffer, exception: false) + + case result + when :wait_readable, :wait_writable + return true + when nil + @finished = true + return false + else + return true + end + rescue OpenSSL::SSL::SSLError => error + if error.message =~ /unexpected eof while reading/ + @finished = true + return false + end + + raise + rescue Errno::ECONNRESET, Errno::EBADF, IOError + @finished = true + return false end protected diff --git a/lib/io/stream/readable.rb b/lib/io/stream/readable.rb index a0792ac..df69113 100644 --- a/lib/io/stream/readable.rb +++ b/lib/io/stream/readable.rb @@ -36,7 +36,7 @@ module Readable getbyte: :readable, readline: :readable, readlines: :readable, - readable?: true, + readable?: :readable, fill_read_buffer: :readable, eof?: :readable, finished?: :readable, diff --git a/releases.md b/releases.md index 29f7849..d62a36e 100644 --- a/releases.md +++ b/releases.md @@ -1,5 +1,9 @@ # Releases +## Unreleased + + - Probe readability through layered transports so a TLS `close_notify` is detected before reusing a connection. + ## v0.13.1 - Set minimum Ruby verison to 3.3.6 to avoid hanging `close` issue in older Ruby versions. diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb new file mode 100644 index 0000000..b8bc07e --- /dev/null +++ b/test/io/stream/tls_readable.rb @@ -0,0 +1,71 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "io/stream/buffered" + +require "sus/fixtures/async/reactor_context" +require "sus/fixtures/openssl/verified_certificate_context" +require "sus/fixtures/openssl/valid_certificate_context" + +describe IO::Stream::Buffered do + include Sus::Fixtures::Async::ReactorContext + include Sus::Fixtures::OpenSSL::VerifiedCertificateContext + include Sus::Fixtures::OpenSSL::ValidCertificateContext + + before do + listener = TCPServer.new("localhost", 0) + port = listener.local_address.ip_port + + @sockets = [ + TCPSocket.new("localhost", port), + listener.accept, + ] + listener.close + + client = OpenSSL::SSL::SSLSocket.new(@sockets[0], client_context) + server = OpenSSL::SSL::SSLSocket.new(@sockets[1], server_context) + + client.sync_close = true + server.sync_close = true + + [ + Async {server.accept}, + Async {client.connect}, + ].each(&:wait) + + @client = IO::Stream::Buffered.wrap(client) + @server = IO::Stream::Buffered.wrap(server) + end + + after do + @client&.close + @server&.close + @sockets.each{|socket| socket.close unless socket.closed?} + end + + attr :client + attr :server + + it "detects a TLS close notification" do + closing = reactor.async do + server.close + end + + @sockets[0].wait_readable(1) + + expect(client).not.to be(:readable?) + closing.wait + end + + it "preserves data consumed by the readability probe" do + server.write("Hello") + server.flush + + @sockets[0].wait_readable(1) + + expect(client).to be(:readable?) + expect(client.read(5)).to be == "Hello" + end +end From 6574764aafc2281fec27f64c7897ac1bc28daf7c Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 19:14:00 +1200 Subject: [PATCH 2/8] Avoid trailing conditionals --- lib/io/stream/buffered.rb | 9 +++++++-- test/io/stream/tls_readable.rb | 6 +++++- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/lib/io/stream/buffered.rb b/lib/io/stream/buffered.rb index bdb6c23..d1658b6 100644 --- a/lib/io/stream/buffered.rb +++ b/lib/io/stream/buffered.rb @@ -91,8 +91,13 @@ def close_write # Check if the stream is readable. # @returns [Boolean] True if the stream is readable. def readable? - return false unless super - return true unless @read_buffer.empty? + unless super + return false + end + + unless @read_buffer.empty? + return true + end # Probe through the wrapped IO rather than its underlying descriptor. This is # important for layered transports such as TLS, where encrypted data on the diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb index b8bc07e..8bad2fc 100644 --- a/test/io/stream/tls_readable.rb +++ b/test/io/stream/tls_readable.rb @@ -42,7 +42,11 @@ after do @client&.close @server&.close - @sockets.each{|socket| socket.close unless socket.closed?} + @sockets.each do |socket| + unless socket.closed? + socket.close + end + end end attr :client From 237a4d079c035ba4eb13582ed3d52f404248482d Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 19:16:00 +1200 Subject: [PATCH 3/8] Treat SSL probe errors as non-viable --- lib/io/stream/buffered.rb | 9 +-------- test/io/stream/tls_readable.rb | 9 +++++++++ 2 files changed, 10 insertions(+), 8 deletions(-) diff --git a/lib/io/stream/buffered.rb b/lib/io/stream/buffered.rb index d1658b6..55f4c7b 100644 --- a/lib/io/stream/buffered.rb +++ b/lib/io/stream/buffered.rb @@ -114,14 +114,7 @@ def readable? else return true end - rescue OpenSSL::SSL::SSLError => error - if error.message =~ /unexpected eof while reading/ - @finished = true - return false - end - - raise - rescue Errno::ECONNRESET, Errno::EBADF, IOError + rescue OpenSSL::SSL::SSLError, Errno::ECONNRESET, Errno::EBADF, IOError @finished = true return false end diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb index 8bad2fc..a20a613 100644 --- a/test/io/stream/tls_readable.rb +++ b/test/io/stream/tls_readable.rb @@ -63,6 +63,15 @@ closing.wait end + it "detects an abrupt TLS connection close" do + @sockets.last.close + @server = nil + + @sockets.first.wait_readable(1) + + expect(client).not.to be(:readable?) + end + it "preserves data consumed by the readability probe" do server.write("Hello") server.flush From 7032fe967d2bfc0b5d60a75627b11030b7b1cac1 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 19:19:58 +1200 Subject: [PATCH 4/8] Avoid ambiguous block syntax --- test/io/stream/tls_readable.rb | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb index a20a613..fc5018b 100644 --- a/test/io/stream/tls_readable.rb +++ b/test/io/stream/tls_readable.rb @@ -30,10 +30,15 @@ client.sync_close = true server.sync_close = true - [ - Async {server.accept}, - Async {client.connect}, - ].each(&:wait) + accept = Async do + server.accept + end + + connect = Async do + client.connect + end + + [accept, connect].each(&:wait) @client = IO::Stream::Buffered.wrap(client) @server = IO::Stream::Buffered.wrap(server) From c161c3febc44eaeeb615186d0bf4e237cff2dbeb Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 19:22:50 +1200 Subject: [PATCH 5/8] Separate consuming transport probe --- lib/io/stream/buffered.rb | 12 +++++++++++- lib/io/stream/readable.rb | 3 ++- releases.md | 2 +- test/io/stream/tls_readable.rb | 12 +++++++++--- 4 files changed, 23 insertions(+), 6 deletions(-) diff --git a/lib/io/stream/buffered.rb b/lib/io/stream/buffered.rb index 55f4c7b..16e3a1f 100644 --- a/lib/io/stream/buffered.rb +++ b/lib/io/stream/buffered.rb @@ -6,6 +6,7 @@ require_relative "generic" require_relative "connection_reset_error" +# Provides buffered IO streams with consistent read, write, and transport semantics. module IO::Stream # A buffered stream implementation that wraps an underlying IO object to provide efficient buffered reading and writing. class Buffered < Generic @@ -91,7 +92,16 @@ def close_write # Check if the stream is readable. # @returns [Boolean] True if the stream is readable. def readable? - unless super + super && @io.readable? + end + + # Probe whether the stream can be read without blocking. + # + # This operation may consume one byte from the wrapped IO. Any byte consumed is preserved in the read buffer. It must not be called concurrently with another read operation. + # + # @returns [Boolean] True if the stream is readable. + def probe_readable? + unless readable? return false end diff --git a/lib/io/stream/readable.rb b/lib/io/stream/readable.rb index df69113..8a77053 100644 --- a/lib/io/stream/readable.rb +++ b/lib/io/stream/readable.rb @@ -36,7 +36,8 @@ module Readable getbyte: :readable, readline: :readable, readlines: :readable, - readable?: :readable, + readable?: true, + probe_readable?: :readable, fill_read_buffer: :readable, eof?: :readable, finished?: :readable, diff --git a/releases.md b/releases.md index d62a36e..4b10d89 100644 --- a/releases.md +++ b/releases.md @@ -2,7 +2,7 @@ ## Unreleased - - Probe readability through layered transports so a TLS `close_notify` is detected before reusing a connection. + - Add `IO::Stream::Buffered#probe_readable?` to probe through layered transports and detect a TLS `close_notify`. ## v0.13.1 diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb index fc5018b..c17a3cf 100644 --- a/test/io/stream/tls_readable.rb +++ b/test/io/stream/tls_readable.rb @@ -64,7 +64,8 @@ @sockets[0].wait_readable(1) - expect(client).not.to be(:readable?) + expect(client).not.to be(:probe_readable?) + expect(client).not.to be(:probe_readable?) closing.wait end @@ -74,7 +75,11 @@ @sockets.first.wait_readable(1) - expect(client).not.to be(:readable?) + expect(client).not.to be(:probe_readable?) + end + + it "reports an open TLS connection as probe readable" do + expect(client).to be(:probe_readable?) end it "preserves data consumed by the readability probe" do @@ -83,7 +88,8 @@ @sockets[0].wait_readable(1) - expect(client).to be(:readable?) + expect(client).to be(:probe_readable?) + expect(client).to be(:probe_readable?) expect(client.read(5)).to be == "Hello" end end From 53436a72a8609eacca3c514995d7f6f0f569bf1f Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 20:03:00 +1200 Subject: [PATCH 6/8] Replace TLS probe with partial peek --- lib/io/stream/buffered.rb | 39 +++++----------------------------- lib/io/stream/readable.rb | 39 ++++++++++++++++++++++++++++++++-- releases.md | 2 +- test/io/stream/tls_readable.rb | 18 ++++++++++------ 4 files changed, 54 insertions(+), 44 deletions(-) diff --git a/lib/io/stream/buffered.rb b/lib/io/stream/buffered.rb index 16e3a1f..9da02e4 100644 --- a/lib/io/stream/buffered.rb +++ b/lib/io/stream/buffered.rb @@ -95,40 +95,6 @@ def readable? super && @io.readable? end - # Probe whether the stream can be read without blocking. - # - # This operation may consume one byte from the wrapped IO. Any byte consumed is preserved in the read buffer. It must not be called concurrently with another read operation. - # - # @returns [Boolean] True if the stream is readable. - def probe_readable? - unless readable? - return false - end - - unless @read_buffer.empty? - return true - end - - # Probe through the wrapped IO rather than its underlying descriptor. This is - # important for layered transports such as TLS, where encrypted data on the - # socket may decode to an EOF (close_notify). Preserve any byte consumed by - # the probe in the stream's read buffer. - result = @io.read_nonblock(1, @read_buffer, exception: false) - - case result - when :wait_readable, :wait_writable - return true - when nil - @finished = true - return false - else - return true - end - rescue OpenSSL::SSL::SSLError, Errno::ECONNRESET, Errno::EBADF, IOError - @finished = true - return false - end - protected def sysclose @@ -139,6 +105,11 @@ def syswrite(buffer) return @io.write(buffer) end + # Attempts to read data from the underlying stream without blocking. + def sysread_nonblock(size, buffer) + return @io.read_nonblock(size, buffer, exception: false) + end + # Reads data from the underlying stream as efficiently as possible. def sysread(size, buffer) # Come on Ruby, why couldn't this just return `nil`? EOF is not exceptional. Every file has one. diff --git a/lib/io/stream/readable.rb b/lib/io/stream/readable.rb index 8a77053..cdc08b7 100644 --- a/lib/io/stream/readable.rb +++ b/lib/io/stream/readable.rb @@ -23,7 +23,7 @@ module IO::Stream # A module providing readable stream functionality. # - # You must implement the `sysread` method to read data from the underlying IO. + # You must implement the `sysread` and `sysread_nonblock` methods to read data from the underlying IO. module Readable ASYNC_SAFE = { read: :readable, @@ -31,13 +31,13 @@ module Readable read_exactly: :readable, read_until: :readable, peek: :readable, + peek_partial: :readable, gets: :readable, getc: :readable, getbyte: :readable, readline: :readable, readlines: :readable, readable?: true, - probe_readable?: :readable, fill_read_buffer: :readable, eof?: :readable, finished?: :readable, @@ -266,6 +266,41 @@ def peek(size = nil) return @read_buffer end + # Peek at data without consuming it, making at most one non-blocking read attempt. + # + # Any data read from the underlying stream is preserved in the read buffer. If + # the read would block or the stream is at EOF, this method returns `nil`. + # + # After this method returns `nil`, {readable?} indicates whether the read would + # block or EOF was observed. + # + # @parameter size [Integer] The maximum number of bytes to peek at. + # @returns [String | Nil] The immediately available data, or nil if no data can be read without blocking. + def peek_partial(size = @minimum_read_size) + if size == 0 + return String.new(encoding: Encoding::BINARY) + end + + if @read_buffer.empty? + if @finished + return nil + end + + read_size = [size, @maximum_read_size].min + + result = sysread_nonblock(read_size, @read_buffer) + case result + when :wait_readable, :wait_writable + return nil + when nil + @finished = true + return nil + end + end + + return @read_buffer.byteslice(0, [size, @read_buffer.bytesize].min) + end + # Read a line from the stream, similar to IO#gets. # @parameter separator [String] The line separator to search for. # @parameter limit [Integer | Nil] The maximum number of bytes to read. diff --git a/releases.md b/releases.md index 4b10d89..1abb706 100644 --- a/releases.md +++ b/releases.md @@ -2,7 +2,7 @@ ## Unreleased - - Add `IO::Stream::Buffered#probe_readable?` to probe through layered transports and detect a TLS `close_notify`. + - Add `IO::Stream::Readable#peek_partial` to peek through layered transports without blocking or consuming application data. ## v0.13.1 diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb index c17a3cf..ad5eddd 100644 --- a/test/io/stream/tls_readable.rb +++ b/test/io/stream/tls_readable.rb @@ -64,8 +64,8 @@ @sockets[0].wait_readable(1) - expect(client).not.to be(:probe_readable?) - expect(client).not.to be(:probe_readable?) + expect(client.peek_partial(1)).to be_nil + expect(client.peek_partial(1)).to be_nil closing.wait end @@ -75,11 +75,15 @@ @sockets.first.wait_readable(1) - expect(client).not.to be(:probe_readable?) + expect do + client.peek_partial(1) + end.to raise_exception(OpenSSL::SSL::SSLError) end - it "reports an open TLS connection as probe readable" do - expect(client).to be(:probe_readable?) + it "reports when reading an open TLS connection would block" do + expect(client.peek_partial(0)).to be == "" + expect(client.peek_partial(1)).to be_nil + expect(client).to be(:readable?) end it "preserves data consumed by the readability probe" do @@ -88,8 +92,8 @@ @sockets[0].wait_readable(1) - expect(client).to be(:probe_readable?) - expect(client).to be(:probe_readable?) + expect(client.peek_partial(1)).to be == "H" + expect(client.peek_partial(1)).to be == "H" expect(client.read(5)).to be == "Hello" end end From e8cb5af0bdb1cf21043d5c600b01c671e625007e Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 20:21:09 +1200 Subject: [PATCH 7/8] Make nonblocking reads optional --- lib/io/stream/readable.rb | 8 +++++++- test/io/stream/generic.rb | 7 +++++++ 2 files changed, 14 insertions(+), 1 deletion(-) diff --git a/lib/io/stream/readable.rb b/lib/io/stream/readable.rb index cdc08b7..cc01125 100644 --- a/lib/io/stream/readable.rb +++ b/lib/io/stream/readable.rb @@ -23,7 +23,7 @@ module IO::Stream # A module providing readable stream functionality. # - # You must implement the `sysread` and `sysread_nonblock` methods to read data from the underlying IO. + # You must implement the `sysread` method to read data from the underlying IO. You may implement `sysread_nonblock` to support non-blocking partial peeks. module Readable ASYNC_SAFE = { read: :readable, @@ -396,6 +396,12 @@ def close_read private + # Attempts to read data from the underlying stream without blocking. + # Implementations may override this method when non-blocking reads are supported. + def sysread_nonblock(size, buffer) + return :wait_readable + end + # Fills the buffer from the underlying stream. def fill_read_buffer(size = @minimum_read_size) # Limit the read size to avoid exceeding SSIZE_MAX and to manage memory usage. diff --git a/test/io/stream/generic.rb b/test/io/stream/generic.rb index 1aa7c72..ce8ea0a 100644 --- a/test/io/stream/generic.rb +++ b/test/io/stream/generic.rb @@ -20,6 +20,13 @@ end end + with "#peek_partial" do + it "should default to no immediately available data" do + expect(stream.peek_partial(1)).to be_nil + expect(stream).to be(:readable?) + end + end + with "#flush" do it "should raise NotImplementedError" do expect{stream.write("hello"); stream.flush}.to raise_exception(NotImplementedError) From a5bcb4b14a9d097b6c17aac9d89cda205fa57a78 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Sat, 25 Jul 2026 20:25:23 +1200 Subject: [PATCH 8/8] Probe stream when peeking zero bytes --- lib/io/stream/readable.rb | 9 +++++---- test/io/stream/tls_readable.rb | 7 ++++--- 2 files changed, 9 insertions(+), 7 deletions(-) diff --git a/lib/io/stream/readable.rb b/lib/io/stream/readable.rb index cc01125..70b0f65 100644 --- a/lib/io/stream/readable.rb +++ b/lib/io/stream/readable.rb @@ -273,20 +273,21 @@ def peek(size = nil) # # After this method returns `nil`, {readable?} indicates whether the read would # block or EOF was observed. + # A size of zero still performs the read attempt, returning an empty string if + # data was available and preserving that data in the read buffer. # # @parameter size [Integer] The maximum number of bytes to peek at. # @returns [String | Nil] The immediately available data, or nil if no data can be read without blocking. def peek_partial(size = @minimum_read_size) - if size == 0 - return String.new(encoding: Encoding::BINARY) - end - if @read_buffer.empty? if @finished return nil end read_size = [size, @maximum_read_size].min + if read_size == 0 + read_size = @minimum_read_size + end result = sysread_nonblock(read_size, @read_buffer) case result diff --git a/test/io/stream/tls_readable.rb b/test/io/stream/tls_readable.rb index ad5eddd..657736c 100644 --- a/test/io/stream/tls_readable.rb +++ b/test/io/stream/tls_readable.rb @@ -64,8 +64,9 @@ @sockets[0].wait_readable(1) + expect(client.peek_partial(0)).to be_nil expect(client.peek_partial(1)).to be_nil - expect(client.peek_partial(1)).to be_nil + expect(client).not.to be(:readable?) closing.wait end @@ -81,7 +82,7 @@ end it "reports when reading an open TLS connection would block" do - expect(client.peek_partial(0)).to be == "" + expect(client.peek_partial(0)).to be_nil expect(client.peek_partial(1)).to be_nil expect(client).to be(:readable?) end @@ -92,7 +93,7 @@ @sockets[0].wait_readable(1) - expect(client.peek_partial(1)).to be == "H" + expect(client.peek_partial(0)).to be == "" expect(client.peek_partial(1)).to be == "H" expect(client.read(5)).to be == "Hello" end