Skip to content
Merged
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Fixed

- Fixed cancelled transaction starts leaving Turso connections unusable, which could prevent maintenance from recovering after a startup failure. [PR #1347](https://github.com/riverqueue/river/pull/1347).
- Fixed maintenance startup failures leaving a client renewing leadership with maintenance stopped in poll-only mode. After exhausting startup retries, clients now request local resignation without depending on database notifications. [PR #1347](https://github.com/riverqueue/river/pull/1347).
- Fixed YugabyteDB clients relying on notifications when `LISTEN/NOTIFY` is unavailable or disabled. Clients now automatically poll for running job cancellations and queue pause, resume, and metadata changes with the default `PollOnly: false`, and skip unsupported notification broadcasts. Native notifications require YugabyteDB 2025.2.3 or later with `ysql_yb_enable_listen_notify=true` on both Masters and TServers. [PR #1347](https://github.com/riverqueue/river/pull/1347).
- Fixed job cancellations received during a fetch being lost before the fetched jobs started. Matching jobs now receive cancellation before work begins. [PR #1397](https://github.com/riverqueue/river/pull/1397).
- Fixed SQLite `JobCancel` and `JobCancelTx` notifying running workers through the shared control outbox, so their contexts are cancelled when the transaction commits. [PR #1398](https://github.com/riverqueue/river/pull/1398).
- Fixed `UniqueOpts.ByArgs` skipping distinct jobs or failing inserts when JSON keys contain path syntax (like `user.id`), are empty, or come from unnamed tags like `json:",omitempty"`. Unaffected unique keys remain unchanged; affected jobs may be inserted again after upgrading or by old and new clients during a rolling upgrade. [PR #1387](https://github.com/riverqueue/river/pull/1387).
Expand Down
40 changes: 31 additions & 9 deletions client.go
Original file line number Diff line number Diff line change
Expand Up @@ -1075,10 +1075,9 @@ func NewClient[TTx any](driver riverdriver.Driver[TTx], config *Config) (*Client
}

client.queueMaintainerLeader = maintenance.NewQueueMaintainerLeader(archetype, &maintenance.QueueMaintainerLeaderConfig{
ClientID: config.ID,
Elector: client.elector,
QueueMaintainer: client.queueMaintainer,
RequestResignFunc: client.clientNotifyBundle.RequestResign,
ClientID: config.ID,
Elector: client.elector,
QueueMaintainer: client.queueMaintainer,
})
client.services = append(client.services, client.queueMaintainerLeader)
client.testSignals.queueMaintainerLeader = &client.queueMaintainerLeader.TestSignals
Expand Down Expand Up @@ -1136,9 +1135,30 @@ func (c *Client[TTx]) Start(ctx context.Context) error {
// available, the client appears to have started even though it's completely
// non-functional. Here we try to make an initial assessment of health and
// return quickly in case of an apparent problem.
if err := c.driver.GetExecutor().Exec(fetchCtx, "SELECT 1"); err != nil {
executor := c.driver.GetExecutor()
if err := executor.Ping(fetchCtx); err != nil {
return fmt.Errorf("error making initial connection to database: %w", err)
}
if err := executor.InitDriver(fetchCtx); err != nil {
return fmt.Errorf("error initializing driver: %w", err)
}

// Database capabilities are only known after initialization. A notifier
// created by NewClient must be removed before any services start if the
// server can't deliver notifications (for example, older Yugabyte).
if c.notifier != nil && !c.driver.SupportsListener() {
c.config.Logger.InfoContext(fetchCtx, "Database does not support listener; entering poll only mode")
c.services = slices.DeleteFunc(c.services, func(service startstop.Service) bool {
return service == c.notifier
})
c.notifier = nil
if c.elector != nil {
c.elector.SetNotifier(nil)
}
for _, producer := range c.producersByQueueName {
producer.config.Notifier = nil
}
}

// Each time we start, we need a fresh completer subscribe channel to
// send job completion events on, because the completer will close it
Expand Down Expand Up @@ -1455,8 +1475,9 @@ func (c *Client[TTx]) Driver() riverdriver.Driver[TTx] {
//
// If the job is currently running, it is not immediately cancelled, but is
// instead marked for cancellation. The client running the job will also be
// notified (via LISTEN/NOTIFY) to cancel the running job's context. Although
// the job's context will be cancelled, since Go does not provide a mechanism to
// notified to cancel the running job's context. When running without a notifier,
// clients poll for cancellation requests every two seconds. Although the job's
// context will be cancelled, since Go does not provide a mechanism to
// interrupt a running goroutine the job will continue running until it returns.
// As always, it is important for workers to respect context cancellation and
// return promptly when the job context is done.
Expand Down Expand Up @@ -1511,8 +1532,9 @@ func (c *Client[TTx]) JobCancel(ctx context.Context, jobID int64) (*rivertype.Jo
//
// If the job is currently running, it is not immediately cancelled, but is
// instead marked for cancellation. The client running the job will also be
// notified (via LISTEN/NOTIFY) to cancel the running job's context. Although
// the job's context will be cancelled, since Go does not provide a mechanism to
// notified to cancel the running job's context. When running without a notifier,
// clients poll for cancellation requests every two seconds. Although the job's
// context will be cancelled, since Go does not provide a mechanism to
// interrupt a running goroutine the job will continue running until it returns.
// As always, it is important for workers to respect context cancellation and
// return promptly when the job context is done.
Expand Down
109 changes: 109 additions & 0 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6363,6 +6363,7 @@ func Test_Client_Maintenance(t *testing.T) {

// After all retries exhausted, the client should request resignation.
client.queueMaintainerLeader.TestSignals.StartRetriesExhausted.WaitOrTimeout()
client.queueMaintainerLeader.TestSignals.ElectedLeader.WaitOrTimeout()
})

t.Run("PeriodicJobEnqueuerWithInsertOpts", func(t *testing.T) {
Expand Down Expand Up @@ -8411,6 +8412,26 @@ func Test_Client_Start_Error(t *testing.T) {
require.Equal(t, pgerrcode.InvalidCatalogName, pgErr.Code)
})

t.Run("DatabaseErrorAfterSuccessfulStart", func(t *testing.T) {
t.Parallel()

dbPool := riversharedtest.DBPoolClone(ctx, t)
driver := NewDriverPollOnly(dbPool)
schema := riverdbtest.TestSchema(ctx, t, driver, nil)

client, err := NewClient(driver, newTestConfig(t, schema))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, client.Stop(ctx)) })

require.NoError(t, client.Start(ctx))
require.NoError(t, client.Stop(ctx))

dbPool.Close()

err = client.Start(ctx)
require.ErrorIs(t, err, riverdriver.ErrClosedPool)
})

t.Run("CanRestartAfterFailure", func(t *testing.T) {
t.Parallel()

Expand All @@ -8437,6 +8458,94 @@ func Test_Client_Start_Error(t *testing.T) {
})
}

func Test_Client_YugabyteQueueControl(t *testing.T) {
t.Parallel()

ctx := context.Background()
for _, testCase := range []struct {
enabled *bool
name string
}{
{enabled: new(false), name: "Disabled"},
{enabled: new(true), name: "Enabled"},
{name: "Unavailable"},
} {
t.Run(testCase.name, func(t *testing.T) {
t.Parallel()

basePool := riversharedtest.DBPool(ctx, t)
schema := riverdbtest.TestSchema(ctx, t, riverpgxv5.New(basePool), nil)
pool := riversharedtest.DBPoolWithYugabyteVersion(ctx, t, schema, testCase.enabled)
config := newTestConfig(t, schema)
config.queuePollInterval = 20 * time.Millisecond
require.False(t, config.PollOnly)

client, err := NewClient(riverpgxv5.New(pool), config)
require.NoError(t, err)
require.NotNil(t, client.notifier, "server capability isn't known until Start")
client.testSignals.Init(t)

// A separate insert-only client prevents local control delivery
// from hiding a missing notification or queue poll.
controller, err := NewClient(riverpgxv5.New(pool), &Config{Schema: schema})
require.NoError(t, err)

for run := range 2 {
events := subscribe(t, client)
startClient(ctx, t, client)
client.queueMaintainerLeader.TestSignals.ElectedLeader.WaitOrTimeout()
if testCase.enabled == nil || !*testCase.enabled {
require.Nil(t, client.notifier)
} else {
require.NotNil(t, client.notifier)
}

if run == 0 {
require.NoError(t, client.Queues().Add("added_after_start", QueueConfig{MaxWorkers: 1}))
}
for _, queue := range []string{QueueDefault, "added_after_start"} {
producer := client.producersByQueueName[queue]
// Only initialize the signals we consume to avoid filling
// unrelated signal buffers while the client is running.
if run == 0 {
producer.testSignals.MetadataChanged.Init(t)
}

require.NoError(t, controller.QueuePause(ctx, queue, nil))
event := riversharedtest.WaitOrTimeout(t, events)
require.Equal(t, EventKindQueuePaused, event.Kind)
require.Equal(t, queue, event.Queue.Name)

inserted, err := controller.Insert(ctx, noOpArgs{}, &InsertOpts{Queue: queue})
require.NoError(t, err)
job, err := controller.JobGet(ctx, inserted.Job.ID)
require.NoError(t, err)
require.Equal(t, rivertype.JobStateAvailable, job.State)

tx, err := pool.Begin(ctx)
require.NoError(t, err)
t.Cleanup(func() { _ = tx.Rollback(ctx) })
_, err = controller.QueueUpdateTx(ctx, tx, queue, &QueueUpdateParams{
Metadata: []byte(fmt.Sprintf(`{"revision":%d}`, run+1)),
})
require.NoError(t, err)
require.NoError(t, tx.Commit(ctx))
producer.testSignals.MetadataChanged.WaitOrTimeout()

require.NoError(t, controller.QueueResume(ctx, queue, nil))
event = riversharedtest.WaitOrTimeout(t, events)
require.Equal(t, EventKindQueueResumed, event.Kind)
require.Equal(t, queue, event.Queue.Name)
event = riversharedtest.WaitOrTimeout(t, events)
require.Equal(t, EventKindJobCompleted, event.Kind)
require.Equal(t, inserted.Job.ID, event.Job.ID)
}
require.NoError(t, client.Stop(ctx))
}
})
}
}

func Test_Config_WithDefaults(t *testing.T) {
t.Parallel()

Expand Down
18 changes: 18 additions & 0 deletions docs/yugabyte.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
# YugabyteDB

The PostgreSQL drivers automatically use polling when YugabyteDB's
`yb_enable_listen_notify` setting is absent or disabled. This includes YugabyteDB
2025.2.1, even with the default `PollOnly: false`. New jobs are picked up on the
`FetchPollInterval`, and queue pause, resume, and metadata changes are picked up
by polling queue settings (every two seconds by default). Running job cancellation
requests are also polled every two seconds, including requests from other clients
and requests made with `JobCancelTx` once committed.

Native notifications require YugabyteDB **2025.2.3 or later**, with
`ysql_yb_enable_listen_notify=true` on **both Masters and TServers**. The feature
is disabled by default. See [Yugabyte's LISTEN/NOTIFY documentation](https://docs.yugabyte.com/stable/api/ysql/the-sql-language/statements/cmd_listen_notify/)
for the additional replication configuration requirements. `PollOnly: true`
continues to force polling even when native notifications are enabled.

Database capabilities are cached for the lifetime of the driver. After enabling
notifications, restart the application with a new driver to detect the change.
53 changes: 34 additions & 19 deletions internal/leadership/elector.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,10 @@ const (
)

type Notification struct {
IsLeader bool
IsLeader bool
// Term identifies a local leadership term and increases on each election.
// Pass it to RequestResign to avoid resigning a subsequent term.
Term uint64
Timestamp time.Time
}

Expand Down Expand Up @@ -233,6 +236,7 @@ type Elector struct {
isLeader bool
pendingRequestResign bool
subscriptions []*Subscription
term uint64
}

type leadershipTerm struct {
Expand Down Expand Up @@ -305,6 +309,29 @@ func trySendWakeup(ctx context.Context, wakeupChan chan struct{}) {
}
}

// RequestResign requests resignation of this elector's given leadership term.
// It is safe to call concurrently and returns without waiting for resignation.
// Requests for a past term, a follower, or a cancelled context have no effect.
// A zero term requests resignation of whichever term is current, as used for
// database notifications. Local callers should use the term from Listen.
func (e *Elector) RequestResign(ctx context.Context, term uint64) {
e.mu.Lock()
defer e.mu.Unlock()

if ctx.Err() != nil || !e.isLeader || (term != 0 && term != e.term) {
return
}

e.pendingRequestResign = true
trySendWakeup(ctx, e.wakeupChan)
}

// SetNotifier changes the notifier before startup, after database capability
// detection. It must only be called while the elector is stopped.
func (e *Elector) SetNotifier(notifier *notifier.Notifier) {
e.notifier = notifier
}

func (e *Elector) Start(ctx context.Context) error {
ctx, shouldStart, started, stopped := e.StartInit(ctx)
if !shouldStart {
Expand Down Expand Up @@ -448,11 +475,7 @@ func (e *Elector) handleLeadershipNotification(ctx context.Context, topic notifi

switch notification.Action {
case DBNotificationKindRequestResign:
if !e.markPendingRequestResign() {
return
}

trySendWakeup(ctx, e.wakeupChan)
e.RequestResign(ctx, 0)
case DBNotificationKindResigned:
// If this a resignation from _this_ client, ignore the change.
if notification.LeaderID == e.config.ClientID {
Expand Down Expand Up @@ -656,6 +679,7 @@ func (e *Elector) Listen() *Subscription {

initialNotification := &Notification{
IsLeader: e.isLeader,
Term: e.term,
Timestamp: sub.creationTime,
}
sub.enqueue(initialNotification)
Expand Down Expand Up @@ -695,30 +719,21 @@ func (e *Elector) leaderTTL() time.Duration {
return e.config.ElectInterval + electIntervalTTLPaddingDefault
}

func (e *Elector) markPendingRequestResign() bool {
e.mu.Lock()
defer e.mu.Unlock()

if !e.isLeader {
return false
}

e.pendingRequestResign = true
return true
}

func (e *Elector) publishLeadershipState(isLeader bool) {
notifyTime := time.Now().UTC()
e.mu.Lock()
defer e.mu.Unlock()

e.isLeader = isLeader
if !isLeader {
if isLeader {
e.term++
} else {
e.pendingRequestResign = false
}

notification := &Notification{
IsLeader: isLeader,
Term: e.term,
Timestamp: notifyTime,
}

Expand Down
Loading
Loading