From 5bbb9a9decccd28dfaa2ac73d314778bd0cd4b0d Mon Sep 17 00:00:00 2001 From: Abishekcs Date: Sun, 9 Aug 2026 14:21:31 +0530 Subject: [PATCH] Run Telemetry.every blocks inside a fiber Timer blocks currently execute directly on the reactor thread in the root fiber. Any blocking I/O in a block - e.g. flushing metrics to a collector via HTTP - stalls the entire reactor and freezes the server. Wrap each execution in Fiber.schedule so blocking operations that go through the fiber scheduler (socket I/O, sleep, DNS) yield to the event loop instead of blocking it. --- CHANGELOG.md | 1 + lib/rage/telemetry/telemetry.rb | 13 ++++++++++++ spec/telemetry/telemetry_spec.rb | 36 ++++++++++++++++++++++++++++++++ 3 files changed, 50 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1d3dd6a9..1d4ee96f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ - [OpenAPI] Add support for the `root:` option in Blueprinter response annotations (#343). - Add `FiberScheduler#timeout_after` (#374). +- [Telemetry] Add `Rage::Telemetry.every(interval_ms, &block)`, a generic scheduling primitive for running recurring work on the reactor — e.g. sampling metrics or measuring event loop lag. (#379) ### Fixed diff --git a/lib/rage/telemetry/telemetry.rb b/lib/rage/telemetry/telemetry.rb index 2a21b1a7..45061223 100644 --- a/lib/rage/telemetry/telemetry.rb +++ b/lib/rage/telemetry/telemetry.rb @@ -36,6 +36,19 @@ def self.available_spans __registry.keys 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. + # + # @param interval_ms [Integer] the execution interval in milliseconds + # @example Periodically sample GC stats + # Rage::Telemetry.every(1000) { MyMetrics.record_gc_stats(GC.stat) } + def self.every(interval_ms, &block) + Iodine.run_every(interval_ms) do + Fiber.schedule { block.call } + end + end + # @private def self.__registry @__registry ||= Spans.constants.each_with_object({}) do |const, memo| diff --git a/spec/telemetry/telemetry_spec.rb b/spec/telemetry/telemetry_spec.rb index f78e6cd6..6441b8d5 100644 --- a/spec/telemetry/telemetry_spec.rb +++ b/spec/telemetry/telemetry_spec.rb @@ -56,6 +56,42 @@ 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 + allow(Iodine).to receive(:running?).and_return(false) + expect(Iodine).to receive(:run_every).with(100) + + described_class.every(100) {} + end + end + + context "when the reactor is running" do + it "registers the timer immediately" do + allow(Iodine).to receive(:running?).and_return(true) + expect(Iodine).to receive(:run_every).with(100) + + described_class.every(100) {} + end + + it "executes the block in a fiber via Iodine.run_every" do + allow(Iodine).to receive(:running?).and_return(true) + + received_block = nil + allow(Iodine).to receive(:run_every) { |_ms, &block| received_block = block } + + Fiber.set_scheduler(Rage::FiberScheduler.new) + my_block = -> { :did_run } + described_class.every(100, &my_block) + + fiber = received_block.call + expect(fiber.__get_result).to eq(:did_run) + ensure + Fiber.set_scheduler(nil) + end + end + end + describe "SpanResult" do subject { described_class::SpanResult }