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

## [Unreleased]

### Added

- Added `Config.LeaderElectionDisabled` to let a client work jobs without participating in leader election or running maintenance services. Other eligible clients in the same database and schema continue handling scheduling, retries, periodic enqueueing, rescue, and cleanup. [PR #1382](https://github.com/riverqueue/river/pull/1382).

### Changed

- `UniqueOpts.ByPeriod` now derives a job's period from its effective scheduled time (`InsertOpts.ScheduledAt` when set, otherwise the insertion time), so scheduled jobs are deduplicated against other jobs scheduled in the same period rather than against jobs inserted in the same period. Periods are also now always measured in UTC, so processes and `ScheduledAt` values in different time zones produce the same unique key for the same period. Unique keys for scheduled `ByPeriod` jobs, and for any `ByPeriod` job inserted from a process whose local time zone isn't UTC, differ from those produced by previous versions. During a rolling upgrade, old and new clients may therefore each insert one job for such a period; jobs that aren't scheduled and are inserted from UTC processes are unaffected. [PR #1377](https://github.com/riverqueue/river/pull/1377).
Expand Down
73 changes: 52 additions & 21 deletions client.go
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,22 @@ type Config struct {
// both a hook and middleware.
Hooks []rivertype.Hook

// LeaderElectionDisabled prevents this client from participating in leader
// election and running maintenance services. It still fetches and executes
// jobs from its configured queues normally.
//
// At least one other started client in the same database and schema must
// remain eligible to lead for scheduled jobs, retries, periodic enqueueing,
// stuck job rescue, and cleanup to progress. This client will never become
// leader, even if no other eligible client is running.
//
// PeriodicJobs must be empty when this is true. Periodic jobs cannot be
// configured on this client, though it can still execute periodic jobs
// enqueued by another client.
//
// Defaults to false.
LeaderElectionDisabled bool

// Logger is the structured logger to use for logging purposes. If none is
// specified, logs will be emitted to STDOUT with messages at warn level
// or higher.
Expand Down Expand Up @@ -304,6 +320,8 @@ type Config struct {

// PeriodicJobs are a set of periodic jobs to run at the specified intervals
// in the client.
//
// Must be empty when LeaderElectionDisabled is true.
PeriodicJobs []*PeriodicJob

// PollOnly starts the client in "poll only" mode, which avoids issuing
Expand Down Expand Up @@ -522,6 +540,7 @@ func (c *Config) WithDefaults() *Config {
JobStuckHandler: c.JobStuckHandler,
JobStuckThreshold: cmp.Or(c.JobStuckThreshold, JobStuckThresholdDefault),
JobTimeout: cmp.Or(c.JobTimeout, JobTimeoutDefault),
LeaderElectionDisabled: c.LeaderElectionDisabled,
Logger: logger,
MaxAttempts: cmp.Or(c.MaxAttempts, MaxAttemptsDefault),
Middleware: c.Middleware,
Expand Down Expand Up @@ -575,6 +594,9 @@ func (c *Config) validate() error {
if c.JobStuckThreshold < 0 {
return errors.New("JobStuckThreshold cannot be less than zero")
}
if c.LeaderElectionDisabled && len(c.PeriodicJobs) > 0 {
return errors.New("PeriodicJobs must be empty when LeaderElectionDisabled is true")
}
if c.MaxAttempts < 0 {
return errors.New("MaxAttempts cannot be less than zero")
}
Expand Down Expand Up @@ -920,11 +942,13 @@ func NewClient[TTx any](driver riverdriver.Driver[TTx], config *Config) (*Client
config.Logger.Info("Driver does not support listener; entering poll only mode")
}

client.elector = leadership.NewElector(archetype, driver.GetExecutor(), client.notifier, &leadership.Config{
ClientID: config.ID,
Schema: config.Schema,
})
client.services = append(client.services, client.elector)
if !config.LeaderElectionDisabled {
client.elector = leadership.NewElector(archetype, driver.GetExecutor(), client.notifier, &leadership.Config{
ClientID: config.ID,
Schema: config.Schema,
})
client.services = append(client.services, client.elector)
}

for queue, queueConfig := range config.Queues {
if _, err := client.producerAdd(queue, queueConfig); err != nil {
Expand All @@ -938,7 +962,9 @@ func NewClient[TTx any](driver riverdriver.Driver[TTx], config *Config) (*Client
if pluginPilot != nil {
client.services = append(client.services, pluginPilot.PluginServices()...)
}
}

if config.willExecuteJobs() && !config.LeaderElectionDisabled {
//
// Maintenance services
//
Expand Down Expand Up @@ -1235,20 +1261,22 @@ func (c *Client[TTx]) Start(ctx context.Context) error {
c.workCancel(rivercommon.ErrStop)

// Stop all mainline services where stop order isn't important.
startstop.StopAllParallel(append(
// This list of services contains the completer, which should always
// stop after the producers so that any remaining work that was enqueued
// will have a chance to have its state completed as it finishes.
//
// TODO: there's a risk here that the completer is stuck on a job that
// won't complete. We probably need a timeout or way to move on in those
// cases.
c.services,
// This list of services contains the completer, which should always
// stop after the producers so that any remaining work that was enqueued
// will have a chance to have its state completed as it finishes.
//
// TODO: there's a risk here that the completer is stuck on a job that
// won't complete. We probably need a timeout or way to move on in those
// cases.
servicesToStop := c.services

if c.queueMaintainer != nil {
// Will only be started if this client was leader, but can tolerate a
// stop without having been started.
c.queueMaintainer,
)...)
servicesToStop = append(servicesToStop, c.queueMaintainer)
}

startstop.StopAllParallel(servicesToStop...)
}()

return nil
Expand Down Expand Up @@ -2565,15 +2593,18 @@ func (c *ClientNotifyBundle[TTx]) requestResignTx(ctx context.Context, execTx ri
// PeriodicJobs returns the currently configured set of periodic jobs for the
// client, and can be used to add new or remove existing ones.
//
// This function should only be invoked on clients capable of running perioidc
// jobs. Running periodic jobs requires that the client be electable as leader
// to run maintenance services, and being electable as leader requires that a
// client be started. To be startable, a client must have Queues and Workers
// configured. Invoking this function will panic if these conditions aren't met.
// This function should only be invoked on clients capable of enqueueing
// periodic jobs: Queues and Workers must be configured and
// LeaderElectionDisabled must be false. Otherwise, invoking this function
// will panic. The client must be started and elected leader for periodic jobs
// to be enqueued.
func (c *Client[TTx]) PeriodicJobs() *PeriodicJobBundle {
if !c.config.willExecuteJobs() {
panic("client Queues and Workers must be configured to modify periodic jobs (otherwise, they'll have no effect because a client not configured to work jobs can't be started)")
}
if c.config.LeaderElectionDisabled {
panic("cannot modify periodic jobs when LeaderElectionDisabled is true")
}

return c.periodicJobs
}
Expand Down
107 changes: 107 additions & 0 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2397,6 +2397,100 @@ func (w *workerWithMiddleware[T]) Work(ctx context.Context, job *Job[T]) error {
return w.workFunc(ctx, job)
}

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

ctx := context.Background()

type testBundle struct {
driver *riverpgxv5.Driver
schema string
}

setup := func(t *testing.T) (*Client[pgx.Tx], *testBundle) {
t.Helper()

driver := riverpgxv5.New(riversharedtest.DBPool(ctx, t))
schema := riverdbtest.TestSchema(ctx, t, driver, nil)
config := newTestConfig(t, schema)
config.LeaderElectionDisabled = true

client, err := NewClient(driver, config)
require.NoError(t, err)

return client, &testBundle{driver: driver, schema: schema}
}

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

client, _ := setup(t)

require.PanicsWithValue(t, "cannot modify periodic jobs when LeaderElectionDisabled is true", func() {
client.PeriodicJobs()
})
})

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

client, _ := setup(t)

require.Nil(t, client.elector)
require.Nil(t, client.periodicJobs)
require.Nil(t, client.queueMaintainer)
require.Nil(t, client.queueMaintainerLeader)
require.NotNil(t, client.completer)
require.NotNil(t, client.notifier)
require.NotNil(t, client.subscriptionManager)
})

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

client, bundle := setup(t)
subscribeChan := subscribe(t, client)
startClient(ctx, t, client)

leaderConfig := newTestConfig(t, bundle.schema)
leaderConfig.Queues = map[string]QueueConfig{"leader": {MaxWorkers: 1}}
leaderConfig.PeriodicJobs = []*PeriodicJob{
NewPeriodicJob(NeverSchedule(), func() (JobArgs, *InsertOpts) {
return noOpArgs{}, nil
}, &PeriodicJobOpts{RunOnStart: true}),
}
leader, err := NewClient(bundle.driver, leaderConfig)
require.NoError(t, err)
leader.testSignals.Init(t)
startClient(ctx, t, leader)
leader.testSignals.queueMaintainerLeader.ElectedLeader.WaitOrTimeout()

// The leader enqueues a periodic job on the default queue, which only
// the client with election disabled consumes.
event := riversharedtest.WaitOrTimeout(t, subscribeChan)
require.Equal(t, EventKindJobCompleted, event.Kind)
require.Equal(t, []string{client.ID()}, event.Job.AttemptedBy)

getLeaderParams := &riverdriver.LeaderGetElectedLeaderParams{Schema: bundle.schema}
elected, err := bundle.driver.GetExecutor().LeaderGetElectedLeader(ctx, getLeaderParams)
require.NoError(t, err)
require.Equal(t, leader.ID(), elected.LeaderID)

stopCtx, cancelFunc := context.WithTimeout(ctx, 5*time.Second)
defer cancelFunc()
require.NoError(t, leader.Stop(stopCtx))

insertRes, err := client.Insert(ctx, noOpArgs{}, nil)
require.NoError(t, err)
event = riversharedtest.WaitOrTimeout(t, subscribeChan)
require.Equal(t, EventKindJobCompleted, event.Kind)
require.Equal(t, insertRes.Job.ID, event.Job.ID)

_, err = bundle.driver.GetExecutor().LeaderGetElectedLeader(ctx, getLeaderParams)
require.ErrorIs(t, err, rivertype.ErrNotFound)
})
}

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

Expand Down Expand Up @@ -8457,6 +8551,7 @@ func Test_NewClient_Defaults(t *testing.T) {
require.NoError(t, err)

require.Zero(t, client.config.AdvisoryLockPrefix)
require.False(t, client.config.LeaderElectionDisabled)

jobCleaner := maintenance.GetService[*maintenance.JobCleaner](client.queueMaintainer)
require.Equal(t, riversharedmaintenance.CancelledJobRetentionPeriodDefault, jobCleaner.Config.CancelledJobRetentionPeriod)
Expand Down Expand Up @@ -8925,6 +9020,18 @@ func Test_NewClient_Validations(t *testing.T) {
config.JobTimeout = 7 * 24 * time.Hour
},
},
{
name: "LeaderElectionDisabled rejects periodic jobs",
configFunc: func(config *Config) {
config.LeaderElectionDisabled = true
config.PeriodicJobs = []*PeriodicJob{
NewPeriodicJob(PeriodicInterval(time.Minute), func() (JobArgs, *InsertOpts) {
return noOpArgs{}, nil
}, nil),
}
},
wantErr: errors.New("PeriodicJobs must be empty when LeaderElectionDisabled is true"),
},
{
name: "MaxAttempts cannot be less than zero",
configFunc: func(config *Config) {
Expand Down
21 changes: 14 additions & 7 deletions periodic_job.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,8 @@ func newPeriodicJobBundle(config *Config, periodicJobEnqueuer *maintenance.Perio
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
func (b *PeriodicJobBundle) Add(periodicJob *PeriodicJob) rivertype.PeriodicJobHandle {
handle, err := b.periodicJobEnqueuer.AddSafely(b.mapper.toInternal(periodicJob))
if err != nil {
Expand All @@ -153,7 +154,8 @@ func (b *PeriodicJobBundle) AddSafely(periodicJob *PeriodicJob) (rivertype.Perio
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
func (b *PeriodicJobBundle) AddMany(periodicJobs []*PeriodicJob) []rivertype.PeriodicJobHandle {
handles, err := b.periodicJobEnqueuer.AddManySafely(sliceutil.Map(periodicJobs, b.mapper.toInternal))
if err != nil {
Expand All @@ -173,7 +175,8 @@ func (b *PeriodicJobBundle) AddManySafely(periodicJobs []*PeriodicJob) ([]rivert
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
func (b *PeriodicJobBundle) Clear() {
b.periodicJobEnqueuer.Clear()
}
Expand All @@ -186,7 +189,8 @@ func (b *PeriodicJobBundle) Clear() {
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
func (b *PeriodicJobBundle) Remove(periodicJobHandle rivertype.PeriodicJobHandle) {
b.periodicJobEnqueuer.Remove(periodicJobHandle)
}
Expand All @@ -196,7 +200,8 @@ func (b *PeriodicJobBundle) Remove(periodicJobHandle rivertype.PeriodicJobHandle
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
//
// Has no effect if no jobs with the given ID is configured.
//
Expand All @@ -214,7 +219,8 @@ func (b *PeriodicJobBundle) RemoveByID(id string) bool {
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
func (b *PeriodicJobBundle) RemoveMany(periodicJobHandles []rivertype.PeriodicJobHandle) {
b.periodicJobEnqueuer.RemoveMany(periodicJobHandles)
}
Expand All @@ -225,7 +231,8 @@ func (b *PeriodicJobBundle) RemoveMany(periodicJobHandles []rivertype.PeriodicJo
// Adding or removing periodic jobs has no effect unless this client is elected
// leader because only the leader enqueues periodic jobs. To make sure that a
// new periodic job is fully enabled or disabled, it should be added or removed
// from _every_ active River client across all processes.
// from every active River client eligible for leader election across all
// processes.
//
// Has no effect if no jobs with the given IDs are configured.
func (b *PeriodicJobBundle) RemoveManyByID(ids []string) {
Expand Down
Loading
Loading