diff --git a/ruby/CHANGELOG.md b/ruby/CHANGELOG.md index ff57759aa..6c82e2c6e 100644 --- a/ruby/CHANGELOG.md +++ b/ruby/CHANGELOG.md @@ -9,8 +9,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- Receive Postgres and SQLite notifications for insertions, queue controls, cancellation, and leadership handoffs, with reconnect recovery and a `poll_only` option. [PR #1496](https://github.com/riverqueue/river/pull/1496). +- Add `Client#request_resign` to request that the maintenance leader relinquish leadership after the caller's transaction commits. [PR #1496](https://github.com/riverqueue/river/pull/1496). - Allow `River::Workers#add` to register a kind and a work block, using the client's default retry and timeout policies. [PR #1486](https://github.com/riverqueue/river/pull/1486). +### Fixed + +- Wake workers when scheduled jobs become available or jobs are manually retried, including workers in other clients. [PR #1496](https://github.com/riverqueue/river/pull/1496). +- Honor cancellations that arrive between claiming a job and starting its worker. [PR #1496](https://github.com/riverqueue/river/pull/1496). +- Preserve worker wakeups that arrive between a fetch and the producer going to sleep. [PR #1496](https://github.com/riverqueue/river/pull/1496). + ## [0.13.0] - 2026-10-08 ### Added diff --git a/ruby/docs/README.md b/ruby/docs/README.md index 1c9c7dece..33a53bc68 100644 --- a/ruby/docs/README.md +++ b/ruby/docs/README.md @@ -712,10 +712,32 @@ A custom service implements `run(client, driver, now)` and runs only while this Service errors are logged without interrupting other services, stuck-job rescue, or cleanup. A failed service is retried on the next scheduled maintenance pass. +Call `client.request_resign` to ask the current maintenance leader to resign. +The leader releases its lease and briefly delays reelection so another client +can take over. The request joins the caller's transaction and is delivered only +on commit. + On SQLite, the leader also removes notification outbox entries older than five -minutes in bounded batches. Cancellation requests write Go-compatible control -notifications in the same transaction as the job update; Ruby workers continue -to observe cancellation through polling. +minutes in bounded batches. + +### Notifications + +Started clients receive job insertion, queue control, cancellation, and +leadership notifications. Postgres uses one dedicated LISTEN connection per +started client, separate from the application's connection pool. SQLite polls +the notification outbox with an independent cursor for each client. Inserts, +scheduled jobs becoming available, and manual retries wake workers promptly. + +The receiver starts before workers and stays active during graceful shutdown. +Startup fails if the initial subscription fails; an established receiver retries +after connection errors and refreshes persisted state on reconnect. Normal job +and queue polling also remains active. + +Set `poll_only: true` in `River::Config` to disable receiving notifications, for +example when connecting through a proxy that cannot support LISTEN. Job, queue, +and cancellation changes are then observed at the polling interval, and +`request_resign` is not received. Backends without native notification support, +such as Yugabyte with LISTEN/NOTIFY disabled, use polling automatically. ### [Renaming job kinds](https://riverqueue.com/docs/renaming-jobs) @@ -746,7 +768,7 @@ The Ruby configuration file must return an unstarted client. The following clien methods are for applications managing their own runtime lifecycle: `client.stop` stops fetching and waits for active jobs to finish. `client.stop_and_cancel` interrupts active worker threads and returns their jobs to `available` without consuming the interrupted attempt. -While draining, the client continues polling for cancellation requests from +While draining, the client continues receiving and polling for cancellation requests from other clients, including for workers with no timeout. Transient polling errors are retried until the active attempts finish. diff --git a/ruby/docs/conformance.md b/ruby/docs/conformance.md index 9b1e8aec3..c83029282 100644 --- a/ruby/docs/conformance.md +++ b/ruby/docs/conformance.md @@ -23,7 +23,7 @@ The fixture checks cover: - All snooze-counter fixtures through real worker attempts and both SQL drivers. - Cron occurrences, named zones, daylight-saving transitions, interval schedules, and invalid expressions, with the Ruby API differences below checked explicitly. -- Insertion, cancellation, queue pause/resume/metadata, and leadership resignation +- Insertion, cancellation, queue pause/resume/metadata, and leadership resignation/request notification encoding through both SQL drivers, using in-memory SQLite and the bundled canonical migrations. @@ -39,13 +39,13 @@ protocol is not part of the Ruby API. ## Coverage gaps and API differences -Ruby's runtime still polls job and queue state rather than consuming Postgres -LISTEN events or the SQLite outbox. The fixture target checks notification -emission, not dispatch. Ruby does not emit or handle `request_resign`; other -clients' requests therefore cannot force a Ruby leader to resign. Completing -that coverage would require a notification listener and runtime dispatch path. -The ordinary driver suites verify Postgres delivery, commit ordering, and -rollback for queue controls and resignations as well as insertion/cancellation. +Ruby consumes Postgres LISTEN events and the SQLite notification outbox. The +ordinary driver suites exercise insert wakeups, cancellation while draining, +queue pause/resume, and leadership handoff on `request_resign` through both +backends and both SQL drivers. They also check commit ordering, rollback, +independent listeners, startup failures, and recovery after listener failures. +The fixture target checks wire formats; live mixed-language worker execution +remains outside its scope. Shared snooze-counter fixtures cover non-negative integers and absent counters. Recovery from other JSON values is implementation-specific. Ruby's shared driver diff --git a/ruby/docs/yugabyte.md b/ruby/docs/yugabyte.md index 8582f4d67..c23f2e563 100644 --- a/ruby/docs/yugabyte.md +++ b/ruby/docs/yugabyte.md @@ -44,8 +44,9 @@ markers and inserted jobs become visible to other connections only when the application transaction commits; rollback leaves neither behind. When native notifications are enabled, insertion and cancellation publish the -same commit-bound notifications as Go. Ruby still consumes changes by polling; -it does not yet implement a notification listener. +same commit-bound notifications as Go. Ruby listens for inserts, queue controls, +and leadership messages on a dedicated connection. Set `poll_only: true` in +`River::Config` to disable the receiver explicitly. Native notifications require YugabyteDB 2025.2.3 or later and `ysql_yb_enable_listen_notify=true` on both Masters and TServers. Follow diff --git a/ruby/driver/riverqueue-activerecord/lib/driver.rb b/ruby/driver/riverqueue-activerecord/lib/driver.rb index 2f45e0d9e..44acae777 100644 --- a/ruby/driver/riverqueue-activerecord/lib/driver.rb +++ b/ruby/driver/riverqueue-activerecord/lib/driver.rb @@ -157,6 +157,10 @@ def transaction(&) time.getutc.round(3).strftime("%Y-%m-%d %H:%M:%S.%3N") end + private def notification_query(sql) + @connection_class.connection_pool.with_connection { runtime_query_rows(sql) } + end + private def postgres_insert_params_to_hash(insert_params, nonce) metadata = insert_params.metadata || {} metadata = metadata.merge(UNIQUE_INSERT_METADATA_KEY => nonce) if nonce @@ -236,6 +240,15 @@ def transaction(&) end end + private def runtime_notification_connection + params, schema = @connection_class.connection_pool.with_connection do |connection| + [connection.raw_connection.conninfo_hash, connection.select_value("SELECT current_schema()")] + end + params.compact! + params[:connect_timeout] = "5" unless params[:connect_timeout].to_i.positive? + [::PG.connect(params), schema] + end + private def runtime_postgres? !@is_sqlite end diff --git a/ruby/driver/riverqueue-activerecord/spec/client_spec.rb b/ruby/driver/riverqueue-activerecord/spec/client_spec.rb index 9f2cc47d4..6ed38ce4a 100644 --- a/ruby/driver/riverqueue-activerecord/spec/client_spec.rb +++ b/ruby/driver/riverqueue-activerecord/spec/client_spec.rb @@ -6,6 +6,7 @@ require_relative "../../../spec/worker_process_shared_examples" require_relative "../../../spec/insert_notification_shared_examples" require_relative "../../../spec/row_decoding_shared_examples" +require_relative "../../../spec/notification_receiver_shared_examples" RSpec.describe "ActiveRecord client integration" do [:postgres, :sqlite].each do |adapter| @@ -24,6 +25,7 @@ end it_behaves_like "client driver end to end" + it_behaves_like "notification receiving", adapter it_behaves_like "Postgres state update races" if adapter == :postgres it_behaves_like "SQLite corrupt job runtime" if adapter == :sqlite it_behaves_like "Postgres insert notifications" if adapter == :postgres diff --git a/ruby/driver/riverqueue-sequel/lib/driver.rb b/ruby/driver/riverqueue-sequel/lib/driver.rb index a2c576e50..fcdd7f6e2 100644 --- a/ruby/driver/riverqueue-sequel/lib/driver.rb +++ b/ruby/driver/riverqueue-sequel/lib/driver.rb @@ -209,6 +209,15 @@ def transaction end end + private def runtime_notification_connection + params, schema = @db.synchronize do |connection| + [connection.conninfo_hash, connection.exec("SELECT current_schema()").getvalue(0, 0)] + end + params.compact! + params[:connect_timeout] = "5" unless params[:connect_timeout].to_i.positive? + [::PG.connect(params), schema] + end + private def runtime_postgres? !@is_sqlite end diff --git a/ruby/driver/riverqueue-sequel/spec/client_spec.rb b/ruby/driver/riverqueue-sequel/spec/client_spec.rb index 17cf40f4a..5e74b1c44 100644 --- a/ruby/driver/riverqueue-sequel/spec/client_spec.rb +++ b/ruby/driver/riverqueue-sequel/spec/client_spec.rb @@ -6,6 +6,7 @@ require_relative "../../../spec/worker_process_shared_examples" require_relative "../../../spec/insert_notification_shared_examples" require_relative "../../../spec/row_decoding_shared_examples" +require_relative "../../../spec/notification_receiver_shared_examples" RSpec.describe "Sequel client integration" do [:postgres, :sqlite].each do |adapter| @@ -24,6 +25,7 @@ end it_behaves_like "client driver end to end" + it_behaves_like "notification receiving", adapter it_behaves_like "Postgres state update races" if adapter == :postgres it_behaves_like "SQLite corrupt job runtime" if adapter == :sqlite it_behaves_like "Postgres insert notifications" if adapter == :postgres diff --git a/ruby/lib/client.rb b/ruby/lib/client.rb index 02c3dae8b..068f69336 100644 --- a/ruby/lib/client.rb +++ b/ruby/lib/client.rb @@ -310,7 +310,13 @@ def queue_update(name, metadata:) @driver.queue_update(name.to_s, metadata: metadata) || raise(NotFoundError, "queue not found: #{name}") end - # Starts polling configured queues and working jobs in background threads. + # Requests that the current leader resign. Delivery waits for the caller's + # Active Record or Sequel transaction to commit, if one is open. + def request_resign + @driver.request_resign + end + + # Starts notification receiving, polling, and working configured queues. # Returns self. def start @runtime.start diff --git a/ruby/lib/client_runtime.rb b/ruby/lib/client_runtime.rb index dbf82328b..31de6cda0 100644 --- a/ruby/lib/client_runtime.rb +++ b/ruby/lib/client_runtime.rb @@ -28,10 +28,18 @@ def initialize(client, driver, config) @condition = ConditionVariable.new @config = config @driver = driver + @fetching_queues = {} + @leader = false + @maintenance_condition = ConditionVariable.new + @maintenance_generation = 0 @mutex = Mutex.new + @notifications_disabled = false @producer_threads = {} @queue_configs = config.queues.dup + @queue_generations = Hash.new(0) + @queue_paused = {} @removed_queues = {} + @resign_requested = false @running = {} @started = false @stopped = true @@ -98,12 +106,21 @@ def perform_job(id, allow_scheduled: false) def publish_queue(kind, queue) event = Event.new(kind, nil, queue, nil) - @mutex.synchronize { @subscriptions.dup }.each { |subscription| subscription.publish(event) } + subscriptions = @mutex.synchronize do + paused = kind == EVENT_QUEUE_PAUSED + return if @queue_paused[queue.name] == paused + + @queue_paused[queue.name] = paused + @subscriptions.dup + end + subscriptions.each { |subscription| subscription.publish(event) } end def healthy? @mutex.synchronize do - @stop_requested || (@producer_threads.all? { |name, thread| @removed_queues[name] || thread.alive? } && (!@maintenance_thread || @maintenance_thread.alive?)) + @stop_requested || (@producer_threads.all? { |name, thread| @removed_queues[name] || thread.alive? } && + (!@maintenance_thread || @maintenance_thread.alive?) && + (!@notification_thread || @notifications_disabled || @notification_thread.alive?)) end end @@ -183,6 +200,7 @@ def start end begin + start_notifications queues = @mutex.synchronize { @queue_configs.dup } queues.each { |name, queue_config| start_producer(name, queue_config) } start_maintenance unless @queue_configs.empty? && @periodic_jobs.empty? && @config.maintenance_services.empty? @@ -207,6 +225,7 @@ def stop(cancel: false, wait: true) @stop_requested = true @condition.broadcast + @maintenance_condition.broadcast @running.values.each { |entry| entry[:thread].raise(Interrupted) if cancel && entry[:working] } @threads.dup end @@ -225,6 +244,9 @@ def stop(cancel: false, wait: true) @threads.clear @producer_threads.clear @maintenance_thread = nil + @notification_thread = nil + @leader = false + @resign_requested = false end self @@ -245,21 +267,23 @@ def subscribe(kinds, buffer_size: 100) end def wake - @mutex.synchronize { @condition.broadcast } + wake_queue("*") end private def begin_work(id) - should_interrupt = @mutex.synchronize do + interruption = @mutex.synchronize do entry = @running.fetch(id) if @stop_requested - true + Interrupted + elsif entry[:cancelled] + JobCancelError else entry[:working] = true - false + nil end end - raise Interrupted if should_interrupt + raise interruption if interruption end private def check_remote_cancellations(queue) @@ -338,6 +362,7 @@ def wake ensure @mutex.synchronize do @running.delete(row.id) + @queue_generations[row.queue] += 1 @condition.broadcast end end @@ -457,16 +482,39 @@ def wake end end - @mutex.synchronize { @running[row.id] = {queue: row.queue, thread: thread, working: false} } + @mutex.synchronize do + @running[row.id] = {cancelled: @fetching_queues.fetch(row.queue).include?(row.id), queue: row.queue, thread: thread, working: false} + end gate.push(true) end private def maintenance_loop leader = false + next_election = 0.0 next_schedule = next_rescue = next_cleanup = Time.at(0) until stopping? + generation, resign = @mutex.synchronize { [@maintenance_generation, @resign_requested] } + if resign && leader + @driver.leader_release(@config.id) + leader = false + @mutex.synchronize do + @leader = false + @resign_requested = false + end + # Give followers a chance to acquire the term we just released. + next_election = monotonic_now + 5 + end + if !leader && (delay = next_election - monotonic_now).positive? + wait_for_maintenance(delay, generation) + next + end + now = Time.now.utc leader = leader ? @driver.leader_renew(@config.id, now: now) : @driver.leader_acquire(@config.id, now: now) + @mutex.synchronize do + @leader = leader + @resign_requested = false unless leader + end if leader if now >= next_schedule @driver.job_schedule(now: now) @@ -496,7 +544,7 @@ def wake end end - wait(5) + wait_for_maintenance(5, generation) end rescue => error @config.logger.error("River maintenance stopped: #{error.full_message}") @@ -536,6 +584,55 @@ def wake ((count + 1 + (1 << 63)) % (1 << 64)) - (1 << 63) end + private def notification_loop(ready) + listener = nil #: untyped + retry_delay = 0.1 + begin + until notifications_stopping? + begin + unless listener + listener = @driver.notification_listener + @mutex.synchronize { @notifications_disabled = listener.nil? } + ready << true if ready + ready = nil + return unless listener + + # Poll durable state after reconnecting: Postgres messages sent + # while disconnected cannot be replayed. + wake + wake_maintenance + end + listener.poll(0.1).each { |topic, payload| receive_notification(topic, payload) } + retry_delay = 0.1 + rescue => error + if ready + ready << error + ready = nil + return + end + @config.logger.error("River notification receiver failed: #{error.full_message}") + # SQLite's cursor survives temporary read errors. Reopening it would + # skip requests committed while the database was unavailable. + unless listener.is_a?(Driver::NotificationListener::SQLite) + listener&.close + listener = nil + end + sleep(retry_delay) unless notifications_stopping? + retry_delay = [retry_delay * 2, 1.0].min + end + end + ensure + ready << true if ready + listener&.close + end + end + + private def notifications_stopping? + @mutex.synchronize do + @stop_requested && @running.empty? && @producer_threads.values.none?(&:alive?) + end + end + private def perform_work(worker, job) timeout = worker_timeout(worker, job) result = timeout ? Timeout.timeout(timeout) { worker.work(job) } : worker.work(job) @@ -545,7 +642,7 @@ def wake private def periodic_jobs_changed start_maintenance - wake + wake_maintenance end private def producer_loop(queue, queue_config) @@ -556,6 +653,7 @@ def wake last_fetch = 0.0 poll_interval = queue_config.resolved_fetch_poll_interval(@config) loop do + generation = @mutex.synchronize { @queue_generations[queue] } draining = queue_stopping?(queue) break if draining && running_count(queue).zero? @@ -573,18 +671,25 @@ def wake wait(sleep_for) end next if queue_stopping?(queue) + end - # A pause may arrive during cooldown; only read the queue immediately - # before fetching, and avoid the read entirely when all slots are busy. - unless @driver.queue_get(queue)&.paused_at + # Refresh controls even when every worker slot is occupied. A pause + # can also arrive during cooldown, so read immediately before fetching. + queue_state = @driver.queue_get(queue) + publish_queue(queue_state.paused_at ? EVENT_QUEUE_PAUSED : EVENT_QUEUE_RESUMED, queue_state) if queue_state + if capacity.positive? && !queue_state&.paused_at + @mutex.synchronize { @fetching_queues[queue] = [] } + begin jobs = @driver.job_get_available(attempted_by: @config.id, max: capacity, queue: queue, **fetch_options) last_fetch = monotonic_now jobs.each { |job| launch(job) } next unless jobs.empty? + ensure + @mutex.synchronize { @fetching_queues.delete(queue) } end end - wait(poll_interval) + wait_for_queue(queue, poll_interval, generation) end rescue => error @config.logger.error("River producer for #{queue.inspect} stopped: #{error.full_message}") @@ -615,6 +720,52 @@ def wake @mutex.synchronize { @stop_requested || @removed_queues[queue] } end + private def receive_notification(topic, payload) + message = JSON.parse(payload) + return unless message.is_a?(Hash) + + case topic + when "river_insert" + wake_queue(message["queue"]) if message["queue"].is_a?(String) + when "river_control" + queue = message["queue"] + return unless queue.is_a?(String) + + case message["action"] + when "cancel" + return unless message["job_id"].is_a?(Integer) + + @mutex.synchronize do + entry = @running[message["job_id"]] + if entry && entry[:queue] == queue + entry[:cancelled] = true + entry[:thread].raise(JobCancelError) if entry[:working] + elsif (pending = @fetching_queues[queue]) + # A claim may have committed before its workers are registered. + pending << message["job_id"] + end + end + when "metadata_changed", "pause", "resume" + wake_queue(queue) + end + when "river_leadership" + case message["action"] + when "request_resign" + @mutex.synchronize do + if @leader + @resign_requested = true + @maintenance_generation += 1 + @maintenance_condition.broadcast + end + end + when "resigned" + wake_maintenance if message["leader_id"].is_a?(String) && message["leader_id"] != @config.id + end + end + rescue JSON::ParserError, TypeError => error + @config.logger.warn("River ignored invalid notification: #{error.message}") + end + private def remove_subscription(subscription) @mutex.synchronize { @subscriptions.delete(subscription) } end @@ -680,11 +831,28 @@ def wake end end + private def start_notifications + return if @config.poll_only || !@driver.respond_to?(:notification_listener) + + ready = ::Queue.new + @mutex.synchronize do + return if @stop_requested + + thread = Thread.new { notification_loop(ready) } + @notification_thread = thread + @threads << thread + end + result = ready.pop + raise result if result.is_a?(Exception) + end + private def start_producer(name, queue_config) - @driver.queue_upsert(name) + queue = @driver.queue_upsert(name) @mutex.synchronize do return if @stop_requested || @removed_queues[name] || @producer_threads.key?(name) + @queue_paused[name] = !queue.paused_at.nil? if queue + # Shutdown must see every launched thread, including a producer whose # database setup overlapped a stop or queue removal. thread = Thread.new { producer_loop(name, queue_config) } @@ -701,6 +869,21 @@ def wake @mutex.synchronize { @condition.wait(@mutex, duration) unless @stop_requested } end + private def wait_for_maintenance(duration, generation) + @mutex.synchronize do + @maintenance_condition.wait(@mutex, duration) if !@stop_requested && @maintenance_generation == generation + end + end + + private def wait_for_queue(queue, duration, generation) + deadline = monotonic_now + duration + @mutex.synchronize do + while !@stop_requested && !@removed_queues[queue] && @queue_generations[queue] == generation && (remaining = deadline - monotonic_now).positive? + @condition.wait(@mutex, remaining) + end + end + end + private def wait_for_running_jobs(queue, duration) @mutex.synchronize do # Draining still needs a poll interval after stop is requested. Check @@ -709,6 +892,20 @@ def wake end end + private def wake_maintenance + @mutex.synchronize do + @maintenance_generation += 1 + @maintenance_condition.broadcast + end + end + + private def wake_queue(queue) + @mutex.synchronize do + @queue_configs.each_key { |name| @queue_generations[name] += 1 if queue == "*" || queue == name } + @condition.broadcast + end + end + private def worker_timeout(worker, job) timeout = worker.respond_to?(:timeout) ? worker.timeout(job) : @config.job_timeout timeout = Float(timeout) unless timeout.nil? diff --git a/ruby/lib/config.rb b/ruby/lib/config.rb index 7fd698a81..82add0547 100644 --- a/ruby/lib/config.rb +++ b/ruby/lib/config.rb @@ -47,7 +47,7 @@ class Config :discarded_job_retention_period, :error_handler, :fetch_cooldown, :fetch_only_known_kinds, :fetch_poll_interval, :id, :job_timeout, :leader_election_disabled, :logger, :maintenance_services, - :periodic_jobs, :plugins, :queues, :retry_policy, :workers + :periodic_jobs, :plugins, :poll_only, :queues, :retry_policy, :workers # Creates a client configuration. # @@ -56,6 +56,8 @@ class Config # +fetch_only_known_kinds+ leaves unregistered kinds untouched in shared # queues. +leader_election_disabled+ prevents this client's maintenance; # another eligible client must run scheduling, rescue, and cleanup. + # +poll_only+ disables notification receiving, including leadership + # resignation requests. Job and queue state are still polled. def initialize( queues: {}, workers: Workers.new, @@ -70,6 +72,7 @@ def initialize( plugins: [], maintenance_services: [], periodic_jobs: [], + poll_only: false, cancelled_job_retention_period: 86_400, completed_job_retention_period: 86_400, discarded_job_retention_period: 604_800, @@ -89,6 +92,7 @@ def initialize( @maintenance_services = maintenance_services.dup.freeze @periodic_jobs = periodic_jobs.dup.freeze @plugins = plugins.dup.freeze + @poll_only = poll_only @queues = normalize_queues(queues).freeze @retry_policy = retry_policy @workers = workers @@ -113,6 +117,7 @@ def with(**overrides) maintenance_services: maintenance_services, periodic_jobs: periodic_jobs, plugins: plugins, + poll_only: poll_only, queues: queues, retry_policy: retry_policy, workers: workers, **overrides) diff --git a/ruby/lib/driver.rb b/ruby/lib/driver.rb index 3be96401c..d88368639 100644 --- a/ruby/lib/driver.rb +++ b/ruby/lib/driver.rb @@ -50,4 +50,5 @@ def initialize( require_relative "driver/postgres_capabilities" require_relative "driver/job_row_decoder" +require_relative "driver/notification_listener" require_relative "driver/runtime" diff --git a/ruby/lib/driver/notification_listener.rb b/ruby/lib/driver/notification_listener.rb new file mode 100644 index 000000000..dd6d2c12c --- /dev/null +++ b/ruby/lib/driver/notification_listener.rb @@ -0,0 +1,61 @@ +# frozen_string_literal: true + +module River::Driver + # Internal notification transports. Each client owns its listener; SQLite + # cursors never consume or delete another client's notifications. + module NotificationListener + TOPICS = %w[river_control river_insert river_leadership].freeze + + class Postgres + def initialize(connection, schema) + @connection = connection + @channels = TOPICS.to_h { |topic| ["#{schema}.#{topic}", topic] } + @channels.each_key { |channel| @connection.exec("LISTEN #{@connection.escape_identifier(channel)}") } + rescue + close + raise + end + + def close + @connection.close unless @connection.finished? + end + + def poll(timeout) + notifications = [] #: Array[untyped] + @connection.wait_for_notify(timeout) do |channel, _pid, payload| + notifications << [@channels[channel], payload] + end + # Bound each batch so a busy sender cannot prevent shutdown. + while notifications.length < 1_000 && (notification = @connection.notifies) + notifications << [@channels[notification[:relname]], notification[:extra]] + end + notifications + end + end + + class SQLite + def initialize(query) + @query = query + # Historical requests must not resign a newly started client. + @last_id = @query.call("SELECT coalesce(max(id), 0) AS id FROM river_notification").first.transform_keys(&:to_sym).fetch(:id).to_i + end + + def close + end + + def poll(timeout) + rows = @query.call("SELECT id, topic, payload FROM river_notification WHERE id > #{@last_id} ORDER BY id LIMIT 1000") + if rows.empty? + sleep(timeout) + return [] + end + + rows.map do |row| + row = row.transform_keys(&:to_sym) + @last_id = row.fetch(:id).to_i + [row.fetch(:topic), row.fetch(:payload)] + end + end + end + end +end diff --git a/ruby/lib/driver/runtime.rb b/ruby/lib/driver/runtime.rb index 93d942981..24445a98f 100644 --- a/ruby/lib/driver/runtime.rb +++ b/ruby/lib/driver/runtime.rb @@ -244,7 +244,9 @@ def job_retry(id, now: Time.now.utc) AND attempt < #{River::MAX_ATTEMPTS_LIMIT} RETURNING id SQL - [updated_id, job_get_by_id(updated_id || id)] + job = job_get_by_id(updated_id || id) + runtime_notify("river_insert", queue: job.queue) if updated_id && job + [updated_id, job] end if !updated_id && job && job.state != "running" && (job.state != "available" || job.scheduled_at > now) && job.attempt >= River::MAX_ATTEMPTS_LIMIT raise ArgumentError, "cannot retry a job with #{River::MAX_ATTEMPTS_LIMIT} or more attempts" @@ -254,18 +256,21 @@ def job_retry(id, now: Time.now.utc) def job_schedule(now: Time.now.utc, max: 1_000) transaction do + queues = [] #: Array[String] # Hold each selected row until its transition (including uniqueness # conflict handling) finishes. A concurrent retry may change its due time. lock = runtime_postgres? ? "FOR UPDATE SKIP LOCKED" : "" - ids = runtime_query_rows(<<~SQL).map { |row| runtime_value(row, :id).to_i } - SELECT id FROM river_job + rows = runtime_query_rows(<<~SQL) + SELECT id, queue FROM river_job WHERE state IN ('retryable', 'scheduled') AND scheduled_at <= #{runtime_time(now)} ORDER BY priority, scheduled_at, id LIMIT #{Integer(max)} #{lock} SQL - ids.each do |id| + rows.each do |row| + id = runtime_value(row, :id).to_i transaction do runtime_execute("UPDATE river_job SET state = 'available' WHERE id = #{id} AND state IN ('retryable', 'scheduled')") + queues << runtime_value(row, :queue) end rescue runtime_unique_violation_class runtime_execute(<<~SQL) @@ -276,7 +281,8 @@ def job_schedule(now: Time.now.utc, max: 1_000) SQL end - ids.length + queues.uniq.each { |queue| runtime_notify("river_insert", queue: queue) } + rows.length end end @@ -365,6 +371,16 @@ def notification_delete_before(horizon:, max: 10_000) SQL end + # Opens a receiver independent of application transactions. Postgres uses + # one dedicated connection; SQLite keeps a cursor in its shared outbox. + def notification_listener + return NotificationListener::SQLite.new(method(:notification_query)) unless runtime_postgres? + return unless postgres_capabilities.supports_listen_notify + + connection, schema = runtime_notification_connection + NotificationListener::Postgres.new(connection, schema) + end + def queue_get(name) row = runtime_query_rows("SELECT #{runtime_queue_columns} FROM river_queue WHERE name = #{runtime_quote(name)}").first runtime_queue_from_row(row) @@ -431,6 +447,16 @@ def queue_upsert(name, metadata: {}, now: Time.now.utc) queue_get(name) end + # Delivery follows the caller's transaction, just like job insertion. + def request_resign + transaction { runtime_notify("river_leadership", action: "request_resign", leader_id: "") } + nil + end + + private def notification_query(sql) + runtime_query_rows(sql) + end + private def runtime_append_error(error) value = error.respond_to?(:to_h) ? error.to_h : error if runtime_postgres? diff --git a/ruby/sig/client.rbs b/ruby/sig/client.rbs index b0652d648..57abc45be 100644 --- a/ruby/sig/client.rbs +++ b/ruby/sig/client.rbs @@ -46,6 +46,7 @@ module River def queue_remove: (String | Symbol) -> self def queue_resume: (String | Symbol) -> bool def queue_update: (String | Symbol, metadata: Hash[untyped, untyped]) -> Queue + def request_resign: () -> nil def start: () -> self def started?: () -> bool def stop: (?wait: bool) -> self diff --git a/ruby/sig/driver.rbs b/ruby/sig/driver.rbs index 00853c97a..4a2a1017a 100644 --- a/ruby/sig/driver.rbs +++ b/ruby/sig/driver.rbs @@ -27,6 +27,8 @@ module River def leader_release: (String) -> untyped def leader_renew: (String, ?ttl: Integer, ?now: Time) -> bool def notification_delete_before: (horizon: Time, ?max: Integer) -> Integer + def notification_listener: () -> untyped + def request_resign: () -> nil def queue_get: (String) -> Queue? def queue_list: (?max: Integer) -> Array[Queue] def queue_pause: (String, ?now: Time) -> Array[Queue] @@ -56,6 +58,25 @@ module River def finish: (JobRow) -> JobRow end + module NotificationListener + TOPICS: Array[String] + class Postgres + @connection: untyped + @channels: Hash[String, String] + def initialize: (untyped, String) -> void + def close: () -> untyped + def poll: (Float | Integer) -> Array[untyped] + end + + class SQLite + @query: untyped + @last_id: Integer + def initialize: (untyped) -> void + def close: () -> nil + def poll: (Float | Integer) -> Array[untyped] + end + end + class PostgresCapabilities attr_reader supports_listen_notify: bool attr_reader unique_insert_mode: :metadata_nonce | :returning_old | :xmax @@ -126,6 +147,8 @@ module River def leader_release: (untyped) -> untyped def leader_renew: (untyped, ?ttl: untyped, ?now: untyped) -> untyped def notification_delete_before: (horizon: untyped, ?max: untyped) -> untyped + def notification_listener: () -> untyped + def request_resign: () -> nil def queue_get: (untyped) -> untyped def queue_list: (?max: untyped) -> untyped def queue_pause: (untyped, ?now: untyped) -> untyped @@ -150,6 +173,8 @@ module River def runtime_metadata_equals: (untyped, untyped) -> untyped def runtime_nullable_time: (untyped) -> untyped def runtime_notify: (String, untyped) -> untyped + def notification_query: (String) -> untyped + def runtime_notification_connection: () -> untyped def sqlite_insert_nonces: (Array[JobInsertParams]) -> Array[String] def runtime_parse_json: (untyped) -> untyped def runtime_parse_time: (untyped) -> untyped diff --git a/ruby/sig/runtime.rbs b/ruby/sig/runtime.rbs index ed40bf87d..aaf77f04c 100644 --- a/ruby/sig/runtime.rbs +++ b/ruby/sig/runtime.rbs @@ -160,6 +160,8 @@ module River end class Config + @poll_only: bool + attr_reader poll_only: bool @cancelled_job_retention_period: Float? @completed_job_retention_period: Float? @discarded_job_retention_period: Float? @@ -198,7 +200,7 @@ module River def initialize: (?queues: untyped, ?workers: untyped, ?id: untyped, ?fetch_cooldown: untyped, ?fetch_only_known_kinds: bool, ?fetch_poll_interval: untyped, ?job_timeout: untyped, - ?leader_election_disabled: bool, + ?leader_election_disabled: bool, ?poll_only: bool, ?retry_policy: untyped, ?error_handler: untyped, ?plugins: untyped, ?maintenance_services: untyped, ?periodic_jobs: untyped, ?cancelled_job_retention_period: untyped, ?completed_job_retention_period: untyped, @@ -341,13 +343,22 @@ module River @condition: ConditionVariable @config: Config @driver: _Driver + @fetching_queues: Hash[String, Array[Integer]] + @leader: bool + @maintenance_condition: ConditionVariable + @maintenance_generation: Integer @maintenance_thread: Thread? @mutex: Mutex - @periodic_jobs: PeriodicJobBundle + @notification_thread: Thread? + @notifications_disabled: bool @performing: bool? + @periodic_jobs: PeriodicJobBundle @producer_threads: Hash[String, Thread] @queue_configs: Hash[String, QueueConfig] + @queue_generations: Hash[String, Integer] + @queue_paused: Hash[String, bool] @removed_queues: Hash[String, bool] + @resign_requested: bool @running: Hash[Integer, untyped] @started: bool @stop_requested: bool @@ -388,22 +399,30 @@ module River def monotonic_now: () -> Float def next_retry: (JobRow, Exception, Time, ?worker: untyped) -> Time def next_snooze_count: (untyped) -> Integer + def notification_loop: (untyped) -> untyped + def notifications_stopping?: () -> bool def retry_allowed?: (untyped, Job, Exception) -> untyped def perform_work: (untyped, Job) -> untyped def periodic_jobs_changed: () -> untyped def producer_loop: (String, QueueConfig) -> untyped def publish: (Symbol, JobRow, JobTiming) -> untyped def queue_stopping?: (String) -> bool + def receive_notification: (untyped, untyped) -> untyped def remove_subscription: (Subscription) -> Subscription? def rescue_job?: (JobRow, Time) -> bool def resolve_worker: (String) -> untyped def run_periodic: (Time) -> untyped def running_count: (String) -> Integer def start_maintenance: () -> untyped + def start_notifications: () -> untyped def start_producer: (String, QueueConfig) -> untyped def stopping?: () -> bool def wait: (Float | Integer) -> untyped + def wait_for_maintenance: (Float | Integer, Integer) -> untyped + def wait_for_queue: (String, Float | Integer, Integer) -> untyped def wait_for_running_jobs: (String, Float | Integer) -> untyped + def wake_maintenance: () -> untyped + def wake_queue: (String) -> untyped def worker_timeout: (untyped, Job) -> Float? end end diff --git a/ruby/spec/client_runtime_branch_spec.rb b/ruby/spec/client_runtime_branch_spec.rb index f8f4662b5..999be512c 100644 --- a/ruby/spec/client_runtime_branch_spec.rb +++ b/ruby/spec/client_runtime_branch_spec.rb @@ -335,7 +335,10 @@ def execute(runtime, value) it "allows queue changes during startup without launching a producer twice" do driver = Object.new value = runtime(driver: driver, queues: {branch: 1}) - driver.define_singleton_method(:queue_upsert) { |name| value.queue_add("dynamic", 1) if name == "branch" } + driver.define_singleton_method(:queue_upsert) do |name| + value.queue_add("dynamic", 1) if name == "branch" + Struct.new(:name, :paused_at).new(name, nil) + end value.define_singleton_method(:producer_loop) { |*_args| } value.define_singleton_method(:start_maintenance) {} @@ -416,7 +419,7 @@ def execute(runtime, value) [] end value = runtime(driver: driver) - value.define_singleton_method(:wait) { |_duration| @stop_requested = true } + value.define_singleton_method(:wait_for_queue) { |_queue, _duration, _generation| @stop_requested = true } value.send(:producer_loop, "branch", River::QueueConfig.new(max_workers: 1)) @@ -502,7 +505,7 @@ def execute(runtime, value) fetched_at = [] job = row driver = Object.new - driver.define_singleton_method(:queue_get) { |_| Struct.new(:paused_at).new(paused_at) } + driver.define_singleton_method(:queue_get) { |_| Struct.new(:name, :paused_at).new("branch", paused_at) } driver.define_singleton_method(:job_get_available) do |**| fetched_at << now [job] @@ -516,6 +519,8 @@ def execute(runtime, value) now += duration end + value.define_singleton_method(:wait_for_queue) { |_, duration, _| wait(duration) } + value.send(:producer_loop, "branch", River::QueueConfig.new(max_workers: 1, fetch_cooldown: 1, fetch_poll_interval: 1)) expect(fetched_at).to eq([100.0]) @@ -531,7 +536,7 @@ def execute(runtime, value) checks >= 2 end - value.define_singleton_method(:wait) { |_duration| } + value.define_singleton_method(:wait_for_maintenance) { |_duration, _generation| } value.send(:maintenance_loop) @@ -546,7 +551,7 @@ def execute(runtime, value) driver.define_singleton_method(method) { |**_options| calls << method } end value = runtime(driver: driver) - value.define_singleton_method(:wait) { |_duration| @stop_requested = true } + value.define_singleton_method(:wait_for_maintenance) { |_duration, _generation| @stop_requested = true } value.send(:maintenance_loop) expect(calls).to eq([:job_schedule, :job_rescue_stuck, :job_delete_finalized]) @@ -606,7 +611,7 @@ def execute(runtime, value) healthy_service.define_singleton_method(:run) { |*_args| calls << :healthy_service } value = described_class.new(Object.new, driver, config.with(logger: Logger.new(output), maintenance_services: [service, healthy_service])) waits = 0 - value.define_singleton_method(:wait) do |_duration| + value.define_singleton_method(:wait_for_maintenance) do |_duration, _generation| waits += 1 @stop_requested = true if waits == 2 end diff --git a/ruby/spec/config_spec.rb b/ruby/spec/config_spec.rb index 82da8cc2c..ee779070c 100644 --- a/ruby/spec/config_spec.rb +++ b/ruby/spec/config_spec.rb @@ -91,6 +91,7 @@ job_timeout: River::JOB_TIMEOUT_DEFAULT, leader_election_disabled: false, plugins: [], + poll_only: false, queues: {}, workers: be_a(River::Workers) ) @@ -124,13 +125,13 @@ plugin = Object.new original = described_class.new( id: "client-one", job_timeout: nil, plugins: [plugin], - fetch_only_known_kinds: true, leader_election_disabled: true, + fetch_only_known_kinds: true, leader_election_disabled: true, poll_only: true, queues: {default: 2}, workers: workers ) copy = original.with(id: "client-two", fetch_poll_interval: 2) expect(copy).to have_attributes(id: "client-two", fetch_poll_interval: 2.0, job_timeout: nil, workers: workers) - expect(copy).to have_attributes(fetch_only_known_kinds: true, leader_election_disabled: true) + expect(copy).to have_attributes(fetch_only_known_kinds: true, leader_election_disabled: true, poll_only: true) expect(copy.queues.keys).to eq(["default"]) expect(copy.plugins).to eq([plugin]) expect(original.id).to eq("client-one") diff --git a/ruby/spec/conformance_spec.rb b/ruby/spec/conformance_spec.rb index 918c1b967..f5c72f4c7 100644 --- a/ruby/spec/conformance_spec.rb +++ b/ruby/spec/conformance_spec.rb @@ -151,9 +151,7 @@ class ConformanceRecord < ActiveRecord::Base (adapter == "activerecord") ? ConformanceRecord.remove_connection : database&.disconnect end - # Ruby emits these messages; its polling runtime does not consume - # request_resign or the other notification fixtures yet. - %w[cancel insert metadata_changed pause resigned resume].each do |name| + %w[cancel insert metadata_changed pause request_resign resigned resume].each do |name| fixture = protocol.fetch("notifications").find { |notification| notification.fetch("name") == name } raise "Missing #{name} notification fixture" unless fixture @@ -177,13 +175,15 @@ class ConformanceRecord < ActiveRecord::Base @driver.queue_pause("priority") @driver.queue_resume("priority") end + when "request_resign" + @client.request_resign when "resigned" @driver.leader_acquire("client-1") @driver.leader_release("client-1") end notification = @driver.send(:runtime_query_rows, "SELECT topic, payload FROM river_notification ORDER BY id DESC LIMIT 1").first expect(@driver.send(:runtime_value, notification, :topic)).to eq(fixture.fetch("topic")) - topic = if name == "resigned" + topic = if %w[request_resign resigned].include?(name) "leadership" else ((name == "insert") ? "insert" : "control") diff --git a/ruby/spec/driver_runtime_feature_spec.rb b/ruby/spec/driver_runtime_feature_spec.rb index 869a64101..3c793561f 100644 --- a/ruby/spec/driver_runtime_feature_spec.rb +++ b/ruby/spec/driver_runtime_feature_spec.rb @@ -36,6 +36,7 @@ postgres_driver.define_singleton_method(:runtime_execute) { |_sql| raise "unexpected notification" } params = Struct.new(:state, :queue).new("available", "default") + expect(postgres_driver.notification_listener).to be_nil expect { postgres_driver.init_driver }.not_to raise_error expect { postgres_driver.send(:postgres_notify_insert, [params]) }.not_to raise_error expect { postgres_driver.send(:runtime_notify, "river_control", {action: "cancel"}) }.not_to raise_error diff --git a/ruby/spec/notification_receiver_shared_examples.rb b/ruby/spec/notification_receiver_shared_examples.rb new file mode 100644 index 000000000..a0924c036 --- /dev/null +++ b/ruby/spec/notification_receiver_shared_examples.rb @@ -0,0 +1,372 @@ +# frozen_string_literal: true + +require "timeout" +require "stringio" + +RSpec.shared_examples "notification receiving" do |backend| + def observe_fetches + fetches = Queue.new + original = @driver.method(:job_get_available) + @driver.define_singleton_method(:job_get_available) do |**options| + result = original.call(**options) + fetches << result + result + end + fetches + end + + def read_notifications(listener, count: 1) + result = [] + Timeout.timeout(3) do + result.concat(Thread.new { listener.poll(0.01) }.value) until result.length >= count + end + result.map { |topic, payload| [topic, JSON.parse(payload)] } + end + + def receive_signal(queue) + Timeout.timeout(3) { queue.pop } + end + + def receiver_client(**options, &work) + River::Client.new(@driver, config: River::Config.new( + fetch_cooldown: 0.001, fetch_poll_interval: 30, + leader_election_disabled: true, queues: {notifications: 1}, + workers: River::Workers.new.add(:notification, &work || ->(_job) {}), **options + )) + end + + if backend == :postgres + it "bounds the default connection timeout and preserves an explicit timeout" do + @driver.transaction do + pool = @driver.send(:runtime_connection_pool) + raw = if pool.respond_to?(:with_connection) + pool.with_connection(&:raw_connection) + else + pool.synchronize { |connection| connection } + end + original = raw.method(:conninfo_hash) + connection = nil + {"0" => "5", "3" => "3"}.each do |configured, expected| + raw.define_singleton_method(:conninfo_hash) { original.call.merge(connect_timeout: configured) } + connection, = @driver.send(:runtime_notification_connection) + expect(connection.conninfo_hash[:connect_timeout]).to eq(expected) + connection.close + end + ensure + raw.singleton_class.remove_method(:conninfo_hash) + connection.close if connection && !connection.finished? + end + end + + it "drains buffered Postgres messages and closes failed subscriptions" do + listener = @driver.notification_listener + @driver.transaction do + 3.times { |id| @driver.send(:runtime_notify, "river_insert", queue: "queue-#{id}") } + end + expect(read_notifications(listener, count: 3).map(&:last)).to eq(3.times.map { |id| {"queue" => "queue-#{id}"} }) + listener.close + expect { listener.close }.not_to raise_error + + connection, schema = @driver.send(:runtime_notification_connection) + connection.exec("BEGIN") + expect { connection.exec("SELECT 1 / 0") }.to raise_error(PG::DivisionByZero) + expect { River::Driver::NotificationListener::Postgres.new(connection, schema) }.to raise_error(PG::InFailedSqlTransaction) + expect(connection.finished?).to be true + ensure + listener&.close + connection.close if connection && !connection.finished? + end + + it "reconnects after Postgres terminates the listener connection" do + fetches = observe_fetches + listeners = Queue.new + original = @driver.method(:notification_listener) + @driver.define_singleton_method(:notification_listener) do + original.call.tap { |listener| listeners << listener } + end + worked = Queue.new + client = receiver_client(logger: Logger.new(StringIO.new)) { |job| worked << job.id } + client.start + listener = receive_signal(listeners) + receive_signal(fetches) + pid = listener.instance_variable_get(:@connection).backend_pid + @driver.send(:runtime_execute, "SELECT pg_terminate_backend(#{pid})") + expect(receive_signal(listeners)).not_to equal(listener) + + row = River::Client.new(@driver).insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + expect(receive_signal(worked)).to eq(row.id) + expect(client.__runtime_healthy?).to be true + ensure + client&.stop_and_cancel + end + end + + it "broadcasts committed requests to independent listeners without replaying history" do + client = River::Client.new(@driver) + client.request_resign + first = @driver.notification_listener + second = @driver.notification_listener + expect(first.poll(0.01)).to be_empty + + @driver.transaction do + client.request_resign + expect(Thread.new { first.poll(0.01) }.value).to be_empty + end + + expected = [["river_leadership", {"action" => "request_resign", "leader_id" => ""}]] + expect(read_notifications(first)).to eq(expected) + expect(read_notifications(second)).to eq(expected) + @driver.transaction do + client.request_resign + raise @driver.rollback_exception + end + expect(first.poll(0.01)).to be_empty + ensure + first&.close + second&.close + end + + it "fails startup if subscribing fails and allows a subsequent start" do + original = @driver.method(:notification_listener) + attempts = 0 + @driver.define_singleton_method(:notification_listener) do + attempts += 1 + raise "subscription failed" if attempts == 1 + + original.call + end + client = receiver_client + expect { client.start }.to raise_error("subscription failed") + expect(client.stopped?).to be true + expect(client.start).to equal(client) + expect(client.__runtime_healthy?).to be true + ensure + client&.stop_and_cancel + end + + it "hands leadership to a waiting follower after request_resign" do + elected = Queue.new + service = Object.new + service.define_singleton_method(:run) { |client, _driver, _now| elected << client.id } + first = receiver_client(id: "notification-leader", leader_election_disabled: false, maintenance_services: [service], queues: {}) + second = receiver_client(id: "notification-follower", leader_election_disabled: false, maintenance_services: [service], queues: {}) + first.start + expect(receive_signal(elected)).to eq(first.id) + second.start + + River::Client.new(@driver).request_resign + + expect(receive_signal(elected)).to eq(second.id) + expect(@driver.leader_renew(first.id)).to be false + expect(@driver.leader_renew(second.id)).to be true + ensure + first&.stop_and_cancel + second&.stop_and_cancel + end + + it "honors a cancellation received between claiming a job and registering its worker" do + worked = Queue.new + client = receiver_client { |job| worked << job.id } + events = client.subscribe(:job_cancelled) + row = client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + runtime = client.instance_variable_get(:@runtime) + original = @driver.method(:job_get_available) + @driver.define_singleton_method(:job_get_available) do |**options| + jobs = original.call(**options) + jobs.each do |job| + # Deliver at the precise point where another connection can observe + # the committed claim, before the producer registers the attempt. + job_cancel(job.id) + runtime.send(:receive_notification, "river_control", JSON.generate(action: "cancel", queue: job.queue, job_id: job.id)) + end + jobs + end + client.start + + expect(Timeout.timeout(3) { events.pop }.job.id).to eq(row.id) + expect(worked).to be_empty + client.stop + expect(runtime.instance_variable_get(:@fetching_queues)).to be_empty + ensure + client&.stop_and_cancel + events&.close + end + + it "ignores malformed and unknown messages and continues receiving" do + fetches = observe_fetches + worked = Queue.new + client = receiver_client { |job| worked << job.id } + client.start + receive_signal(fetches) + [nil, [], {}, {queue: 42}, {queue: "notifications", action: "unknown"}].each do |payload| + @driver.send(:runtime_notify, "river_control", payload) + end + @driver.send(:runtime_notify, "river_leadership", action: "unknown") + row = client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + + expect(receive_signal(worked)).to eq(row.id) + expect(client.__runtime_healthy?).to be true + ensure + client&.stop_and_cancel + end + + it "keeps polling available when notification receiving is disabled" do + @driver.define_singleton_method(:notification_listener) { raise "unexpected notification listener" } + worked = Queue.new + client = receiver_client(poll_only: true, fetch_poll_interval: 0.01) { |job| worked << job.id } + client.start + row = client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + expect(receive_signal(worked)).to eq(row.id) + ensure + client&.stop_and_cancel + end + + it "notifies workers when another client retries a job" do + fetches = observe_fetches + worked = Queue.new + client = receiver_client { |job| worked << job.id } + row = client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications, state: :pending).job + client.start + receive_signal(fetches) + + River::Client.new(@driver).job_retry(row.id) + + expect(receive_signal(worked)).to eq(row.id) + ensure + client&.stop_and_cancel + end + + it "notifies workers when scheduled jobs become available" do + fetches = observe_fetches + worked = Queue.new + client = receiver_client { |job| worked << job.id } + row = client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications, scheduled_at: Time.now.utc - 1, state: :scheduled).job + client.start + receive_signal(fetches) + + @driver.job_schedule + + expect(receive_signal(worked)).to eq(row.id) + ensure + client&.stop_and_cancel + end + + it "observes remote queue pause and wildcard resume without waiting for a poll" do + fetches = observe_fetches + worked = Queue.new + client = receiver_client { |job| worked << job.id } + events = client.subscribe(:queue_paused, :queue_resumed) + client.start + receive_signal(fetches) + other = River::Client.new(@driver) + other.queue_pause("*") + expect(Timeout.timeout(3) { events.pop }).to have_attributes(kind: :queue_paused, queue: have_attributes(name: "notifications")) + row = other.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + other.queue_resume("*") + + expect(receive_signal(worked)).to eq(row.id) + expect(Timeout.timeout(3) { events.pop }.kind).to eq(:queue_resumed) + ensure + client&.stop_and_cancel + events&.close + end + + it "publishes queue controls while every worker slot is occupied" do + entered = Queue.new + client = receiver_client { + entered << true + sleep + } + events = client.subscribe(:queue_paused, :queue_resumed) + client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications) + client.start + receive_signal(entered) + + other = River::Client.new(@driver) + other.queue_pause(:notifications) + expect(Timeout.timeout(3) { events.pop }.kind).to eq(:queue_paused) + other.queue_resume(:notifications) + expect(Timeout.timeout(3) { events.pop }.kind).to eq(:queue_resumed) + ensure + client&.stop_and_cancel + events&.close + end + + it "receives cancellations while gracefully draining a long-running worker" do + entered = Queue.new + client = receiver_client(job_timeout: nil) do |job| + entered << job.id + sleep + end + events = client.subscribe(:job_cancelled) + row = client.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + client.start + expect(receive_signal(entered)).to eq(row.id) + client.stop(wait: false) + + River::Client.new(@driver).job_cancel(row.id) + + expect(Timeout.timeout(3) { events.pop }.job.id).to eq(row.id) + Timeout.timeout(3) { client.stop } + expect(client.stopped?).to be true + ensure + client&.stop_and_cancel + events&.close + end + + it "recovers from a listener failure and receives again after a restart" do + fetches = observe_fetches + failures = Queue.new + failed = Queue.new + opened = Queue.new + worked = Queue.new + output = StringIO.new + original = @driver.method(:notification_listener) + @driver.define_singleton_method(:notification_listener) do + listener = original.call + poll = listener.method(:poll) + listener.define_singleton_method(:poll) do |timeout| + unless failures.empty? + failures.pop + failed << true + raise "simulated listener failure" + end + poll.call(timeout) + end + opened << listener + listener + end + client = receiver_client(logger: Logger.new(output)) { |job| worked << job.id } + client.start + receive_signal(opened) + receive_signal(fetches) + failures << true + receive_signal(failed) + other = River::Client.new(@driver) + row = other.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + + expect(receive_signal(worked)).to eq(row.id) + expect(client.__runtime_healthy?).to be true + expect(output.string).to include("simulated listener failure") + client.stop + client.start + row = other.insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + expect(receive_signal(worked)).to eq(row.id) + ensure + client&.stop_and_cancel + end + + it "wakes idle workers after an insertion from another client" do + fetches = observe_fetches + worked = Queue.new + client = receiver_client { |job| worked << job.id } + client.start + expect(receive_signal(fetches)).to be_empty + + row = River::Client.new(@driver).insert(River::JobArgsHash.new(:notification, {}), queue: :notifications).job + + expect(receive_signal(worked)).to eq(row.id) + ensure + client&.stop_and_cancel + end +end diff --git a/ruby/spec/notification_receiver_spec.rb b/ruby/spec/notification_receiver_spec.rb new file mode 100644 index 000000000..c1c4d014d --- /dev/null +++ b/ruby/spec/notification_receiver_spec.rb @@ -0,0 +1,24 @@ +# frozen_string_literal: true + +require "spec_helper" +require "riverqueue-activerecord" +require "riverqueue-sequel" +require_relative "notification_receiver_shared_examples" +require_relative "support/client_test_database" + +RSpec.describe "River notification receiver" do + %i[active_record sequel].each do |adapter| + %i[postgres sqlite].each do |backend| + context "#{adapter} with #{backend}", database: backend do + around do |example| + ClientTestDatabase.public_send(:"with_#{adapter}", backend) do |driver| + @driver = driver + example.run + end + end + + it_behaves_like "notification receiving", backend + end + end + end +end diff --git a/ruby/spec/notification_runtime_spec.rb b/ruby/spec/notification_runtime_spec.rb new file mode 100644 index 000000000..8b6fcb9c3 --- /dev/null +++ b/ruby/spec/notification_runtime_spec.rb @@ -0,0 +1,135 @@ +# frozen_string_literal: true + +require "spec_helper" +require "stringio" + +RSpec.describe "River notification dispatch and recovery" do + let(:driver) { Object.new } + let(:output) { StringIO.new } + let(:runtime) do + River::ClientRuntime.new(Object.new, driver, River::Config.new( + id: "notification-client", leader_election_disabled: true, logger: Logger.new(output) + )) + end + + it "backs off repeated connection failures without indefinitely delaying shutdown" do + listener = Object.new + listener.define_singleton_method(:poll) { |_| raise "connection lost" } + listener.define_singleton_method(:close) {} + attempts = 0 + driver.define_singleton_method(:notification_listener) do + attempts += 1 + listener if attempts <= 7 + end + delays = [] + runtime.define_singleton_method(:sleep) { |duration| delays << duration } + + runtime.send(:notification_loop, Queue.new) + + expect(delays).to eq([0.1, 0.2, 0.4, 0.8, 1.0, 1.0, 1.0]) + end + + it "cancels an attempt before its worker starts but ignores another queue's job" do + entry = {queue: "one", thread: Thread.current, working: false} + runtime.instance_variable_set(:@running, 42 => entry) + runtime.send(:receive_notification, "river_control", JSON.generate(action: "cancel", queue: "two", job_id: 42)) + expect(entry[:cancelled]).to be_nil + + runtime.send(:receive_notification, "river_control", JSON.generate(action: "cancel", queue: "one", job_id: 42)) + expect { runtime.send(:begin_work, 42) }.to raise_error(River::JobCancelError) + expect(entry[:working]).to be false + end + + it "ignores invalid payloads without disrupting the receiver" do + [nil, "{", "null", "[]", "{}"].each { |payload| runtime.send(:receive_notification, "river_insert", payload) } + runtime.send(:receive_notification, "unknown", "{}") + runtime.send(:receive_notification, "river_control", JSON.generate(action: "cancel", queue: "one", job_id: "42")) + runtime.send(:receive_notification, "river_control", JSON.generate(action: "cancel", queue: "one", job_id: 42)) + + expect(runtime.instance_variable_get(:@running)).to be_empty + expect(output.string).to include("ignored invalid notification") + end + + it "keeps a client healthy when its backend cannot listen" do + driver.define_singleton_method(:notification_listener) { nil } + runtime.start + expect(runtime).to be_healthy + runtime.stop + expect(runtime).to be_stopped + ensure + runtime.stop(cancel: true) + end + + it "retries connection failures after losing an established subscription" do + closes = 0 + listener = Object.new + listener.define_singleton_method(:poll) { |_| raise "connection lost" } + listener.define_singleton_method(:close) { closes += 1 } + attempts = 0 + driver.define_singleton_method(:notification_listener) do + attempts += 1 + raise "reconnect failed" if attempts == 2 + listener if attempts == 1 + end + ready = Queue.new + + runtime.send(:notification_loop, ready) + + expect(ready.pop(true)).to be true + expect(ready).to be_empty + expect(attempts).to eq(3) + expect(closes).to eq(1) + expect(output.string).to include("connection lost", "reconnect failed") + end + + it "routes leadership messages only to the appropriate elector" do + runtime.send(:receive_notification, "river_leadership", JSON.generate(action: "request_resign")) + runtime.send(:receive_notification, "river_leadership", JSON.generate(action: "resigned", leader_id: "notification-client")) + runtime.send(:receive_notification, "river_leadership", JSON.generate(action: "resigned", leader_id: 42)) + expect(runtime.instance_variable_get(:@maintenance_generation)).to eq(0) + expect(runtime.instance_variable_get(:@resign_requested)).to be false + + runtime.send(:receive_notification, "river_leadership", JSON.generate(action: "resigned", leader_id: "another-client")) + expect(runtime.instance_variable_get(:@maintenance_generation)).to eq(1) + runtime.instance_variable_set(:@leader, true) + runtime.send(:receive_notification, "river_leadership", JSON.generate(action: "request_resign")) + expect(runtime.instance_variable_get(:@resign_requested)).to be true + expect(runtime.instance_variable_get(:@maintenance_generation)).to eq(2) + end + + it "skips reconnect backoff after stopping and closes the failed listener" do + value = runtime + closed = false + listener = Object.new + listener.define_singleton_method(:poll) do |_| + value.instance_variable_set(:@stop_requested, true) + raise "connection lost during stop" + end + listener.define_singleton_method(:close) { closed = true } + driver.define_singleton_method(:notification_listener) { listener } + + runtime.send(:notification_loop, Queue.new) + + expect(closed).to be true + expect(output.string).to include("connection lost during stop") + end + + it "unblocks startup if shutdown wins the race to start receiving" do + driver.define_singleton_method(:notification_listener) { raise "unexpected subscription" } + runtime.instance_variable_set(:@stop_requested, true) + runtime.send(:start_notifications) + expect(runtime.instance_variable_get(:@threads)).to be_empty + + ready = Queue.new + runtime.send(:notification_loop, ready) + expect(ready.pop(true)).to be true + end + + it "wakes only the queue named by an insertion" do + runtime.queue_add("one", 1) + runtime.queue_add("two", 1) + generations = runtime.instance_variable_get(:@queue_generations).dup + runtime.send(:receive_notification, "river_insert", JSON.generate(queue: "one")) + expect(runtime.instance_variable_get(:@queue_generations)).to eq(generations.merge("one" => generations["one"] + 1)) + end +end diff --git a/ruby/spec/support/client_test_database.rb b/ruby/spec/support/client_test_database.rb index d7b58fedf..7efd99b85 100644 --- a/ruby/spec/support/client_test_database.rb +++ b/ruby/spec/support/client_test_database.rb @@ -9,7 +9,8 @@ # tables: migrate an empty disposable schema with the bundled canonical SQL. module ClientTestDatabase def self.with_active_record(adapter, migrate: true, pg_catalog_last: false) - original = ActiveRecord::Base.connection_db_config.configuration_hash + # The core suite may not have configured a connection yet. + original = ActiveRecord::Base.remove_connection if adapter == :postgres ActiveRecord::Base.establish_connection(ENV["TEST_DATABASE_URL"] || "postgres://localhost/river_test") schema = "river_client_test_#{SecureRandom.hex(8)}" @@ -42,7 +43,8 @@ def self.with_active_record(adapter, migrate: true, pg_catalog_last: false) begin ActiveRecord::Base.connection.execute("DROP SCHEMA #{schema} CASCADE") if schema_created ensure - ActiveRecord::Base.establish_connection(original) + ActiveRecord::Base.remove_connection + ActiveRecord::Base.establish_connection(original) if original end end