diff --git a/CHANGELOG.md b/CHANGELOG.md index 4fd232894..22610a510 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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). diff --git a/client.go b/client.go index f1dce2b27..a0ee88d92 100644 --- a/client.go +++ b/client.go @@ -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. @@ -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 @@ -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, @@ -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") } @@ -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 { @@ -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 // @@ -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 @@ -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 } diff --git a/client_test.go b/client_test.go index 98b44a4ef..a87e0d36e 100644 --- a/client_test.go +++ b/client_test.go @@ -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() @@ -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) @@ -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) { diff --git a/periodic_job.go b/periodic_job.go index a6fd87956..e06c99576 100644 --- a/periodic_job.go +++ b/periodic_job.go @@ -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 { @@ -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 { @@ -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() } @@ -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) } @@ -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. // @@ -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) } @@ -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) { diff --git a/plugin_test.go b/plugin_test.go index 13135c6f1..9520d27fc 100644 --- a/plugin_test.go +++ b/plugin_test.go @@ -95,7 +95,7 @@ func TestClientPilotPlugin(t *testing.T) { pluginPilot *TestPilotWithPlugin } - setup := func(t *testing.T) (*Client[pgx.Tx], *testBundle) { + setup := func(t *testing.T, leaderElectionDisabled bool) (*Client[pgx.Tx], *testBundle) { t.Helper() var ( @@ -106,6 +106,7 @@ func TestClientPilotPlugin(t *testing.T) { pluginDriver = newDriverWithPlugin(t, dbPool) pluginPilot = newPilotWithPlugin(t) ) + config.LeaderElectionDisabled = leaderElectionDisabled pluginDriver.pilot = pluginPilot client, err := NewClient(pluginDriver, config) @@ -117,10 +118,27 @@ func TestClientPilotPlugin(t *testing.T) { } } + t.Run("LeaderElectionDisabled", func(t *testing.T) { + t.Parallel() + + client, bundle := setup(t, true) + + startClient(ctx, t, client) + riversharedtest.WaitOrTimeout(t, client.baseStartStop.Started()) + riversharedtest.WaitOrTimeout(t, bundle.pluginPilot.service.Started()) + + require.False(t, bundle.pluginPilot.maintenanceServicesCalled) + select { + case <-bundle.pluginPilot.maintenanceService.Started(): + t.Fatal("plugin maintenance service should not have started") + default: + } + }) + t.Run("ServicesStart", func(t *testing.T) { t.Parallel() - client, bundle := setup(t) + client, bundle := setup(t, false) startClient(ctx, t, client) @@ -136,8 +154,9 @@ var _ pilotPlugin = &TestPilotWithPlugin{} type TestPilotWithPlugin struct { riverpilot.StandardPilot - maintenanceService startstop.Service - service startstop.Service + maintenanceService startstop.Service + maintenanceServicesCalled bool + service startstop.Service } func newPilotWithPlugin(t *testing.T) *TestPilotWithPlugin { @@ -170,6 +189,7 @@ func newPilotWithPlugin(t *testing.T) *TestPilotWithPlugin { } func (d *TestPilotWithPlugin) PluginMaintenanceServices() []startstop.Service { + d.maintenanceServicesCalled = true return []startstop.Service{d.maintenanceService} } diff --git a/riverdriver/riverdrivertest/driver_client_test.go b/riverdriver/riverdrivertest/driver_client_test.go index f1d720ebb..df850063c 100644 --- a/riverdriver/riverdrivertest/driver_client_test.go +++ b/riverdriver/riverdrivertest/driver_client_test.go @@ -1253,6 +1253,50 @@ func ExerciseClient[TTx any](ctx context.Context, t *testing.T, require.Equal(t, job.ID, listRes.Jobs[0].ID) }) + t.Run("LeaderElectionDisabled", func(t *testing.T) { + t.Parallel() + + for _, testCase := range []struct { + name string + pollOnly bool + }{ + {name: "Default"}, + {name: "PollOnly", pollOnly: true}, + } { + t.Run(testCase.name, func(t *testing.T) { + t.Parallel() + + config, bundle := setupConfig(t) + config.LeaderElectionDisabled = true + config.PollOnly = testCase.pollOnly + + client, err := river.NewClient(bundle.driver, config) + require.NoError(t, err) + + // Exercise restart as well as initial startup, including shutdown + // without a queue maintainer or a leadership lease to resign. + for range 2 { + subscribeChan := subscribe(t, client) + startClient(ctx, t, client) + + insertRes, err := client.Insert(ctx, noOpArgs{}, nil) + require.NoError(t, err) + event := riversharedtest.WaitOrTimeout(t, subscribeChan) + require.Equal(t, river.EventKindJobCompleted, event.Kind) + require.Equal(t, insertRes.Job.ID, event.Job.ID) + + _, err = bundle.exec.LeaderGetElectedLeader(ctx, &riverdriver.LeaderGetElectedLeaderParams{Schema: bundle.schema}) + require.ErrorIs(t, err, rivertype.ErrNotFound) + + stopCtx, cancelFunc := context.WithTimeout(ctx, 5*time.Second) + err = client.Stop(stopCtx) + cancelFunc() + require.NoError(t, err) + } + }) + } + }) + t.Run("QueueGet", func(t *testing.T) { t.Parallel()