diff --git a/CHANGES.txt b/CHANGES.txt index 64abd4a8..f7b5d755 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,8 @@ CHANGES +8.12.0 (unreleased) +- Added SSE::EventSource::Client#last_activity_at / #seconds_since_last_activity (and SSEHandler delegates) exposing the monotonic time of the last bytes read from the streaming connection, keepalives included, so applications can monitor stream liveness. + 8.11.3 (Sep, 3, 2026) - Updated concurrent-ruby version to higher than 1.0.4, this removed thread_safe, which was merged into concurrent-ruby. diff --git a/lib/splitclient-rb/sse/event_source/client.rb b/lib/splitclient-rb/sse/event_source/client.rb index f47cb0ea..58256e59 100644 --- a/lib/splitclient-rb/sse/event_source/client.rb +++ b/lib/splitclient-rb/sse/event_source/client.rb @@ -32,6 +32,7 @@ def initialize(config, @status_queue = status_queue @read_timeout = read_timeout @connected = Concurrent::AtomicBoolean.new(false) + @last_activity_at = Concurrent::AtomicReference.new(nil) # Set while close is tearing a connection down on purpose, so the reader # thread does not report the resulting socket error as a retryable failure. @shutdown = Concurrent::AtomicBoolean.new(false) @@ -54,6 +55,7 @@ def close(status = nil) @shutdown.make_true push_status(status) @connected.make_false + @last_activity_at.set(nil) socket = @socket @socket = nil @@ -91,8 +93,24 @@ def connected? @connected.value end + # Monotonic timestamp (Process::CLOCK_MONOTONIC seconds) of the last bytes read + # from the stream, including :keepalive comments. nil when not connected. + def last_activity_at + @last_activity_at.get + end + + # Seconds since the last bytes were read from the stream, nil when not connected. + def seconds_since_last_activity + at = last_activity_at + at.nil? ? nil : monotonic_now - at + end + private + def monotonic_now + Process.clock_gettime(Process::CLOCK_MONOTONIC) + end + def connect_thread(latch) # Never orphan the previous reader thread: its handle lives in a single slot # and would otherwise become unreachable, leaking the thread forever. @@ -124,6 +142,7 @@ def connect_stream(latch, generation = nil) if IO.select([socket], nil, nil, @read_timeout) begin partial_data = socket.readpartial(10_000) + @last_activity_at.set(monotonic_now) if first_event first_event = false diff --git a/lib/splitclient-rb/sse/sse_handler.rb b/lib/splitclient-rb/sse/sse_handler.rb index 24eec281..f2f7f10a 100644 --- a/lib/splitclient-rb/sse/sse_handler.rb +++ b/lib/splitclient-rb/sse/sse_handler.rb @@ -30,6 +30,14 @@ def connected? @sse_client&.connected? || false end + def last_activity_at + @sse_client&.last_activity_at + end + + def seconds_since_last_activity + @sse_client&.seconds_since_last_activity + end + def start_workers @splits_worker.start @segments_worker.start diff --git a/spec/sse/event_source/client_spec.rb b/spec/sse/event_source/client_spec.rb index 39827dc4..8045b096 100644 --- a/spec/sse/event_source/client_spec.rb +++ b/spec/sse/event_source/client_spec.rb @@ -50,6 +50,29 @@ let(:event_error) { "d4\r\nevent: error\ndata: {\"message\":\"Token expired\",\"code\":40142,\"statusCode\":401,\"href\":\"https://help.ably.io/error/40142\"}" } context 'tests' do + it 'tracks stream activity until the client closes' do + sse_client = subject.new(config, api_token, telemetry_runtime_producer, event_parser, notification_manager_keeper, notification_processor, push_status_queue) + + expect(sse_client.last_activity_at).to be_nil + expect(sse_client.seconds_since_last_activity).to be_nil + + mock_server do |server| + server.setup_response('/') do |_, res| + send_stream_content(res, "c\r\n:keepalive\n\n\r\n") + end + + expect(sse_client.start(server.base_uri)).to eq(true) + expect(sse_client.last_activity_at).to be_a(Float) + expect(sse_client.seconds_since_last_activity).to be >= 0 + expect(sse_client.seconds_since_last_activity).to be < 5 + + sse_client.close + + expect(sse_client.last_activity_at).to be_nil + expect(sse_client.seconds_since_last_activity).to be_nil + end + end + it 'receive split update event' do stub_request(:get, 'https://sdk.split.io/api/splitChanges?s=1.3&since=-1&rbSince=-1') .with(headers: { 'Authorization' => 'Bearer client-spec-key' }) diff --git a/spec/sse/sse_handler_spec.rb b/spec/sse/sse_handler_spec.rb index 582afb07..ce1d9f82 100644 --- a/spec/sse/sse_handler_spec.rb +++ b/spec/sse/sse_handler_spec.rb @@ -79,6 +79,17 @@ end end + it 'delegates stream activity accessors' do + sse_handler = subject.new(config, splits_worker, segments_worker, sse_client) + timestamp = Process.clock_gettime(Process::CLOCK_MONOTONIC) + + allow(sse_client).to receive(:last_activity_at).and_return(timestamp) + allow(sse_client).to receive(:seconds_since_last_activity).and_return(1.5) + + expect(sse_handler.last_activity_at).to eq(timestamp) + expect(sse_handler.seconds_since_last_activity).to eq(1.5) + end + private def send_content(res, content)