From c3e2d6381173486c7e170f285bf92b2c39faa94c Mon Sep 17 00:00:00 2001 From: Roman Samoilov <2270393+rsamoilov@users.noreply.github.com> Date: Tue, 18 Aug 2026 16:25:56 +0100 Subject: [PATCH 1/2] Use one fiber per WS connection --- lib/rage/cable/cable.rb | 25 +++++++++++++++++++------ 1 file changed, 19 insertions(+), 6 deletions(-) diff --git a/lib/rage/cable/cable.rb b/lib/rage/cable/cable.rb index 4fab4d6f..2a75ae87 100644 --- a/lib/rage/cable/cable.rb +++ b/lib/rage/cable/cable.rb @@ -118,13 +118,26 @@ def on_shutdown(connection) private - def schedule_fiber(connection) - Fiber.schedule do - @log_processor.init_request_logger(connection.env) - yield - rescue => e - log_error(e) + def schedule_fiber(connection, &block) + connection.env["rage.worker"] ||= begin + worker = Fiber.schedule do + env = Fiber.yield + Rage.__log_processor.init_request_logger(env) + + loop do + work = Fiber.yield + work.call + rescue => e + log_error(e) + end + end + + worker.resume(connection.env) + + worker end + + connection.env["rage.worker"].resume(block) end def log_error(e) From 6072e2469ceafb50fe15c59dca8525a5b428868b Mon Sep 17 00:00:00 2001 From: Roman Samoilov <2270393+rsamoilov@users.noreply.github.com> Date: Tue, 18 Aug 2026 16:26:46 +0100 Subject: [PATCH 2/2] Add `Rage::Signal` --- lib/rage-rb.rb | 1 + lib/rage/signal.rb | 88 ++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 89 insertions(+) create mode 100644 lib/rage/signal.rb diff --git a/lib/rage-rb.rb b/lib/rage-rb.rb index f31240e0..4a177277 100644 --- a/lib/rage-rb.rb +++ b/lib/rage-rb.rb @@ -208,6 +208,7 @@ module ActiveRecord autoload :Events, "rage/events/events" autoload :PubSub, "rage/pubsub/pubsub" autoload :Daemon, "rage/daemon" + autoload :Signal, "rage/signal" end module RageController diff --git a/lib/rage/signal.rb b/lib/rage/signal.rb new file mode 100644 index 00000000..ac4ea053 --- /dev/null +++ b/lib/rage/signal.rb @@ -0,0 +1,88 @@ +# frozen_string_literal: true + +# @private +# This is an in-progress PoC component +class Rage::Signal + class << self + def __signals + @__signals ||= Hash.new do |h, k| + h[k] = Hash.new { |h_in, k_in| h_in[k_in] = Set.new } + end + end + + def __subscriptions + @__subscriptions ||= Hash.new { |h, k| h[k] = {} } + end + + def on(signal_id, subscription_id, &block) + if __subscriptions.has_key?(signal_id) && __subscriptions[signal_id].has_key?(subscription_id) + raise ArgumentError, "subscription already exists" + end + + signal_worker = build_signal_worker + + __signals[signal_id][signal_worker] << block + __subscriptions[signal_id][subscription_id] = block + subscribe_to_signal(signal_id) + + true + end + + def emit(signal_id, message = signal_id) + serialized_signal_id = Rage::Internal.stream_name_for(signal_id) + Iodine.publish("signal:#{serialized_signal_id}", message) if Iodine.running? + # TODO: redis adapter + # TODO: filter sender out? + + true + end + + def off(signal_id, subscription_id) + if !__subscriptions.has_key?(signal_id) && !__subscriptions[signal_id].has_key?(subscription_id) + raise ArgumentError, "subscription doesn't exist" + end + + callback = __subscriptions[signal_id].delete(subscription_id) + + signal_worker = Fiber[:signal_worker] + raise unless signal_worker # TODO: iterate through all workers instead? + + worker_callbacks = __signals[signal_id][signal_worker] + worker_callbacks.delete(callback) + __signals[signal_id].delete(signal_worker) if worker_callbacks.empty? + + true + end + + private + + def build_signal_worker + # new fiber inherits parent's fiber storage with `live_components`; + # the fiber is cached per connection in case `on` is called for multiple components in a loop + Fiber[:signal_worker] ||= Fiber.schedule do + loop do + callback, message = Fiber.yield + callback.call(message) + rescue => e + Rage::Errors.report(e) + Rage.logger.error("#{e.class} (#{e.message}):\n#{e.backtrace.join("\n")}") + end + end + end + + def subscribe_to_signal(signal_id) + serialized_signal_id = Rage::Internal.stream_name_for(signal_id) + return if Iodine.subscribed?("signal:#{serialized_signal_id}") + + Iodine.subscribe("signal:#{serialized_signal_id}") do |_, msg| + __signals[signal_id].each do |worker, callbacks| + callbacks.each { |callback| worker.resume([callback, msg]) } + end + end + + Iodine.on_state(:start_shutdown) do + Iodine.unsubscribe("signal:#{serialized_signal_id}") + end + end + end +end