From 429ffa79d91d4f4bcb5d4b786e54898837723912 Mon Sep 17 00:00:00 2001 From: Abishekcs Date: Thu, 3 Sep 2026 20:06:09 +0530 Subject: [PATCH 1/4] Add Rage::Telemetry.queued_connections Exposes the number of connections currently waiting in the kernel's accept queue for the server's listening sockets how many clients have connected but haven't been picked up by the app yet. Backed by Iodine.queued_connections, which walks Iodine's own listeners and reads each one's queue depth from the kernel. Meant to be sampled on a timer via Rage::Telemetry.every, the same way GC stats are, so it can be reported alongside other reactor health metrics. Adds two specs: one confirms the reported count matches the number of connections actually parked in the queue, the other confirms the method works correctly when called from inside a .every block. --- lib/rage/telemetry/telemetry.rb | 7 +++++ spec/telemetry/telemetry_spec.rb | 53 ++++++++++++++++++++++++++++++++ 2 files changed, 60 insertions(+) diff --git a/lib/rage/telemetry/telemetry.rb b/lib/rage/telemetry/telemetry.rb index 45061223..360688cd 100644 --- a/lib/rage/telemetry/telemetry.rb +++ b/lib/rage/telemetry/telemetry.rb @@ -36,6 +36,13 @@ def self.available_spans __registry.keys end + # Returns the number of established connections currently waiting in the + # kernel's accept queue for the server's listening socket(s). + # @return [Integer] the accept-queue depth + def self.queued_connections + Iodine.queued_connections + end + # Registers a block to be executed repeatedly at a fixed interval while the server is running. # The block is run inside a fiber, so blocking I/O inside it (e.g. flushing metrics to a # collector) will not block the server. diff --git a/spec/telemetry/telemetry_spec.rb b/spec/telemetry/telemetry_spec.rb index 6441b8d5..fafd8e6f 100644 --- a/spec/telemetry/telemetry_spec.rb +++ b/spec/telemetry/telemetry_spec.rb @@ -56,6 +56,59 @@ end end + describe ".queued_connections" do + before do + Fiber.set_scheduler(Rage::FiberScheduler.new) + end + + after do + Fiber.set_scheduler(nil) + end + + it "reports the accept-queue depth of the server's listeners" do + require "socket" + + s = TCPServer.new("127.0.0.1", 0) + port = s.addr[1] + s.close + result = nil + + Iodine.workers = 1 + Iodine.on_state(:on_start) do + Iodine.listen(port: port, handler: -> { [200, {}, ["ok"]] }) + socks = 5.times.map { TCPSocket.new("127.0.0.1", port) } + Iodine.run { result = described_class.queued_connections } + socks.each(&:close) + Iodine.run { Iodine.stop } + end + Iodine.start + + expect(result).to eq(5) + end + + it "is callable from within the reactor via .every" do + require "socket" + + s = TCPServer.new("127.0.0.1", 0) + port = s.addr[1] + s.close + result = nil + + Iodine.workers = 1 + Iodine.on_state(:on_start) do + Iodine.listen(port: port, handler: -> { [200, {}, ["ok"]] }) + Rage::Telemetry.every(5) do + Iodine.run { result = described_class.queued_connections } + Iodine.stop + end + Iodine.run_after(200) { Iodine.stop } + end + Iodine.start + + expect(result).to be_a(Integer) + end + end + describe ".every" do context "when the reactor is not running" do it "registers the timer, which Iodine defers until the server starts" do From a0ba1544f8733d86ec32da7680b5e77c7994edb1 Mon Sep 17 00:00:00 2001 From: Abishekcs Date: Mon, 7 Sep 2026 20:42:47 +0530 Subject: [PATCH 2/4] Move `queued_connections` to a new module `Rage::Telemetry::Capacity` - Rename Iodine.queued_connections to Iodine::Perf.queued_connections as `queued_connections` has been moved into the `Perf` module in iodine codebase. --- lib/rage/telemetry/capacity.rb | 47 +++++++++++++++++++++++++++ lib/rage/telemetry/telemetry.rb | 8 +---- spec/telemetry/capacity_spec.rb | 56 ++++++++++++++++++++++++++++++++ spec/telemetry/telemetry_spec.rb | 53 ------------------------------ 4 files changed, 104 insertions(+), 60 deletions(-) create mode 100644 lib/rage/telemetry/capacity.rb create mode 100644 spec/telemetry/capacity_spec.rb diff --git a/lib/rage/telemetry/capacity.rb b/lib/rage/telemetry/capacity.rb new file mode 100644 index 00000000..9e49cd6c --- /dev/null +++ b/lib/rage/telemetry/capacity.rb @@ -0,0 +1,47 @@ +# frozen_string_literal: true + +module Rage::Telemetry + ## + # The `Rage::Telemetry::Capacity` module provides read-only access to + # metrics describing the server's current load and resource utilization. + # Example: how much work is queued up, and how close the server is to its limits. + # + # Unlike spans, capacity metrics aren't tied to a specific operation or + # event, they reflect the server state at the moment they are read, and + # are meant to be sampled periodically rather than triggered by handlers. + # + # Combine with {Rage::Telemetry.every} to report a metric on a fixed + # interval: + # + # Rage::Telemetry.every(1000) do + # MyMetrics.gauge("server.queued_connections", Rage::Telemetry::Capacity.queued_connections) + # end + # + # # Available Metrics + # + # | ---------- Method -------|--------Description-------- | + # | `.queued_connections` | The number of established connections currently waiting in the kernel's accept queue, across the server listening sockets | + # + # @see Rage::Telemetry.every + # + module Capacity + class << self + # Returns the number of established connections currently waiting in the + # kernel accept queue for the server listening sockets. This is the + # count of clients that have already completed the TCP handshake but + # haven't yet been picked up by the application via `accept()`. + # + # A value that stays close to the configured backlog limit is a sign the + # server isn't accepting connections fast enough to keep up with incoming + # traffic. + # + # @return [Integer] the accept-queue depth + def queued_connections + accept_queue_depth = Iodine.queued_connections + raise NotImplementedError if accept_queue_depth.nil? + + accept_queue_depth + end + end + end +end diff --git a/lib/rage/telemetry/telemetry.rb b/lib/rage/telemetry/telemetry.rb index 360688cd..728112c5 100644 --- a/lib/rage/telemetry/telemetry.rb +++ b/lib/rage/telemetry/telemetry.rb @@ -36,13 +36,6 @@ def self.available_spans __registry.keys end - # Returns the number of established connections currently waiting in the - # kernel's accept queue for the server's listening socket(s). - # @return [Integer] the accept-queue depth - def self.queued_connections - Iodine.queued_connections - end - # Registers a block to be executed repeatedly at a fixed interval while the server is running. # The block is run inside a fiber, so blocking I/O inside it (e.g. flushing metrics to a # collector) will not block the server. @@ -141,4 +134,5 @@ def success? require_relative "tracer" require_relative "handler" +require_relative "capacity" Dir["#{__dir__}/spans/*.rb"].each { |span| require_relative span } diff --git a/spec/telemetry/capacity_spec.rb b/spec/telemetry/capacity_spec.rb new file mode 100644 index 00000000..497ac8c2 --- /dev/null +++ b/spec/telemetry/capacity_spec.rb @@ -0,0 +1,56 @@ +# frozen_string_literal: true + +RSpec.describe Rage::Telemetry::Capacity do + describe ".queued_connections" do + before do + Fiber.set_scheduler(Rage::FiberScheduler.new) + end + + after do + Fiber.set_scheduler(nil) + end + + it "reports the accept-queue depth of the server's listeners" do + require "socket" + + s = TCPServer.new("127.0.0.1", 0) + port = s.addr[1] + s.close + result = nil + + Iodine.workers = 1 + Iodine.on_state(:on_start) do + Iodine.listen(port: port, handler: -> { [200, {}, ["ok"]] }) + socks = 5.times.map { TCPSocket.new("127.0.0.1", port) } + Iodine.run { result = described_class.queued_connections } + socks.each(&:close) + Iodine.run { Iodine.stop } + end + Iodine.start + + expect(result).to eq(5) + end + + it "is callable from within the reactor via .every" do + require "socket" + + s = TCPServer.new("127.0.0.1", 0) + port = s.addr[1] + s.close + result = nil + + Iodine.workers = 1 + Iodine.on_state(:on_start) do + Iodine.listen(port: port, handler: -> { [200, {}, ["ok"]] }) + Rage::Telemetry.every(5) do + Iodine.run { result = described_class.queued_connections } + Iodine.stop + end + Iodine.run_after(200) { Iodine.stop } + end + Iodine.start + + expect(result).to be_a(Integer) + end + end +end diff --git a/spec/telemetry/telemetry_spec.rb b/spec/telemetry/telemetry_spec.rb index fafd8e6f..6441b8d5 100644 --- a/spec/telemetry/telemetry_spec.rb +++ b/spec/telemetry/telemetry_spec.rb @@ -56,59 +56,6 @@ end end - describe ".queued_connections" do - before do - Fiber.set_scheduler(Rage::FiberScheduler.new) - end - - after do - Fiber.set_scheduler(nil) - end - - it "reports the accept-queue depth of the server's listeners" do - require "socket" - - s = TCPServer.new("127.0.0.1", 0) - port = s.addr[1] - s.close - result = nil - - Iodine.workers = 1 - Iodine.on_state(:on_start) do - Iodine.listen(port: port, handler: -> { [200, {}, ["ok"]] }) - socks = 5.times.map { TCPSocket.new("127.0.0.1", port) } - Iodine.run { result = described_class.queued_connections } - socks.each(&:close) - Iodine.run { Iodine.stop } - end - Iodine.start - - expect(result).to eq(5) - end - - it "is callable from within the reactor via .every" do - require "socket" - - s = TCPServer.new("127.0.0.1", 0) - port = s.addr[1] - s.close - result = nil - - Iodine.workers = 1 - Iodine.on_state(:on_start) do - Iodine.listen(port: port, handler: -> { [200, {}, ["ok"]] }) - Rage::Telemetry.every(5) do - Iodine.run { result = described_class.queued_connections } - Iodine.stop - end - Iodine.run_after(200) { Iodine.stop } - end - Iodine.start - - expect(result).to be_a(Integer) - end - end - describe ".every" do context "when the reactor is not running" do it "registers the timer, which Iodine defers until the server starts" do From 691633ce4d5c1e0d746517bec629c13439f13b8e Mon Sep 17 00:00:00 2001 From: Abishekcs Date: Mon, 7 Sep 2026 20:56:52 +0530 Subject: [PATCH 3/4] Update Changelog and update change `rage-iodine` to v5.6 in spec dependency --- CHANGELOG.md | 1 + rage.gemspec | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 78378c5c..3857312e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ - [OpenAPI] Add `openapi:validate` Rake task for OpenAPI tags validation (#163). - [Logger] Add `config.log_redact_keys=` for redacting structured log context. +- [Telemetry] Add `Rage::Telemetry::Capacity` to provide read-only access to metrics describing the server's current load and resource utilization. ### Fixed diff --git a/rage.gemspec b/rage.gemspec index 6c4da2be..8d3d676b 100644 --- a/rage.gemspec +++ b/rage.gemspec @@ -29,7 +29,7 @@ Gem::Specification.new do |spec| spec.add_dependency "thor", "~> 1.0" spec.add_dependency "rack", "< 4" - spec.add_dependency "rage-iodine", "~> 5.5" + spec.add_dependency "rage-iodine", "~> 5.6" spec.add_dependency "zeitwerk", "~> 2.6" spec.add_dependency "rack-test", "~> 2.1" spec.add_dependency "rake", ">= 12.0" From e1bde73620a1b8d152efca41479ee7cc209f9832 Mon Sep 17 00:00:00 2001 From: Abishekcs Date: Mon, 7 Sep 2026 21:06:27 +0530 Subject: [PATCH 4/4] Fix queued_connections API to use Iodine::Perf module --- lib/rage/telemetry/capacity.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/rage/telemetry/capacity.rb b/lib/rage/telemetry/capacity.rb index 9e49cd6c..9b1cb12e 100644 --- a/lib/rage/telemetry/capacity.rb +++ b/lib/rage/telemetry/capacity.rb @@ -37,7 +37,7 @@ class << self # # @return [Integer] the accept-queue depth def queued_connections - accept_queue_depth = Iodine.queued_connections + accept_queue_depth = Iodine::Perf.queued_connections raise NotImplementedError if accept_queue_depth.nil? accept_queue_depth