Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions ruby/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
30 changes: 26 additions & 4 deletions ruby/docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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.

Expand Down
16 changes: 8 additions & 8 deletions ruby/docs/conformance.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand All @@ -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
Expand Down
5 changes: 3 additions & 2 deletions ruby/docs/yugabyte.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 13 additions & 0 deletions ruby/driver/riverqueue-activerecord/lib/driver.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions ruby/driver/riverqueue-activerecord/spec/client_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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|
Expand All @@ -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
Expand Down
9 changes: 9 additions & 0 deletions ruby/driver/riverqueue-sequel/lib/driver.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions ruby/driver/riverqueue-sequel/spec/client_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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|
Expand All @@ -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
Expand Down
8 changes: 7 additions & 1 deletion ruby/lib/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading