diff --git a/CHANGELOG.md b/CHANGELOG.md index 55a503d9..cf96ccb2 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 + +- `Config.TestOnly` has been added. It disables various features in the River client like staggered maintenance service start that are useful in production, but may be somewhat harmful in tests because they make start/stop slower. [PR #414](https://github.com/riverqueue/river/pull/414). + ### Fixed - Pausing or resuming a queue that was already paused or not paused respectively no longer returns `rivertype.ErrNotFound`. The same goes for pausing or resuming using the all queues string (`*`) when no queues are in the database (previously that also returned `rivertype.ErrNotFound`). [PR #408](https://github.com/riverqueue/river/pull/408). diff --git a/client.go b/client.go index 88efec97..458cf8d9 100644 --- a/client.go +++ b/client.go @@ -200,6 +200,17 @@ type Config struct { // Defaults to DefaultRetryPolicy. RetryPolicy ClientRetryPolicy + // TestOnly can be set to true to disable certain features that are useful + // in production, but which may be harmful to tests, in ways like having the + // effect of making them slower. It should not be used outside of test + // suites. + // + // For example, queue maintenance services normally stagger their startup + // with a random jittered sleep so they don't all try to work at the same + // time. This is nice in production, but makes starting and stopping the + // client in a test case slower. + TestOnly bool + // Workers is a bundle of registered job workers. // // This field may be omitted for a program that's only enqueueing jobs @@ -208,12 +219,6 @@ type Config struct { // (i.e. That it wasn't forgotten by accident.) Workers *Workers - // Disables the normal random jittered sleep occurring in queue maintenance - // services to stagger their startup so they don't all try to work at the - // same time. Appropriate for use in tests to make sure that the client can - // always be started and stopped again hastily. - disableStaggerStart bool - // Scheduler run interval. Shared between the scheduler and producer/job // executors, but not currently exposed for configuration. schedulerInterval time.Duration @@ -447,8 +452,8 @@ func NewClient[TTx any](driver riverdriver.Driver[TTx], config *Config) (*Client ReindexerSchedule: config.ReindexerSchedule, RescueStuckJobsAfter: valutil.ValOrDefault(config.RescueStuckJobsAfter, rescueAfter), RetryPolicy: retryPolicy, + TestOnly: config.TestOnly, Workers: config.Workers, - disableStaggerStart: config.disableStaggerStart, schedulerInterval: valutil.ValOrDefault(config.schedulerInterval, maintenance.JobSchedulerIntervalDefault), time: config.time, } @@ -607,7 +612,7 @@ func NewClient[TTx any](driver riverdriver.Driver[TTx], config *Config) (*Client // started conditionally based on whether the client is the leader. client.queueMaintainer = maintenance.NewQueueMaintainer(archetype, maintenanceServices) - if config.disableStaggerStart { + if config.TestOnly { client.queueMaintainer.StaggerStartupDisable(true) } } diff --git a/client_test.go b/client_test.go index 009d1a27..9df4a20a 100644 --- a/client_test.go +++ b/client_test.go @@ -137,15 +137,15 @@ func newTestConfig(t *testing.T, callback callbackFunc) *Config { AddWorker(workers, &noOpWorker{}) return &Config{ - FetchCooldown: 20 * time.Millisecond, - FetchPollInterval: 50 * time.Millisecond, - Logger: riverinternaltest.Logger(t), - MaxAttempts: MaxAttemptsDefault, - Queues: map[string]QueueConfig{QueueDefault: {MaxWorkers: 50}}, - Workers: workers, - disableStaggerStart: true, // disables staggered start in maintenance services - schedulerInterval: riverinternaltest.SchedulerShortInterval, - time: &riverinternaltest.TimeStub{}, + FetchCooldown: 20 * time.Millisecond, + FetchPollInterval: 50 * time.Millisecond, + Logger: riverinternaltest.Logger(t), + MaxAttempts: MaxAttemptsDefault, + Queues: map[string]QueueConfig{QueueDefault: {MaxWorkers: 50}}, + TestOnly: true, // disables staggered start in maintenance services + Workers: workers, + schedulerInterval: riverinternaltest.SchedulerShortInterval, + time: &riverinternaltest.TimeStub{}, } } @@ -3957,9 +3957,11 @@ func Test_NewClient_Defaults(t *testing.T) { require.Equal(t, maintenance.CancelledJobRetentionPeriodDefault, jobCleaner.Config.CancelledJobRetentionPeriod) require.Equal(t, maintenance.CompletedJobRetentionPeriodDefault, jobCleaner.Config.CompletedJobRetentionPeriod) require.Equal(t, maintenance.DiscardedJobRetentionPeriodDefault, jobCleaner.Config.DiscardedJobRetentionPeriod) + require.False(t, jobCleaner.StaggerStartupIsDisabled()) enqueuer := maintenance.GetService[*maintenance.PeriodicJobEnqueuer](client.queueMaintainer) require.Zero(t, enqueuer.Config.AdvisoryLockPrefix) + require.False(t, enqueuer.StaggerStartupIsDisabled()) require.Nil(t, client.config.ErrorHandler) require.Equal(t, FetchCooldownDefault, client.config.FetchCooldown) @@ -3968,7 +3970,6 @@ func Test_NewClient_Defaults(t *testing.T) { require.NotZero(t, client.baseService.Logger) require.Equal(t, MaxAttemptsDefault, client.config.MaxAttempts) require.IsType(t, &DefaultClientRetryPolicy{}, client.config.RetryPolicy) - require.False(t, client.config.disableStaggerStart) } func Test_NewClient_Overrides(t *testing.T) { @@ -3998,8 +3999,8 @@ func Test_NewClient_Overrides(t *testing.T) { MaxAttempts: 5, Queues: map[string]QueueConfig{QueueDefault: {MaxWorkers: 1}}, RetryPolicy: retryPolicy, + TestOnly: true, // disables staggered start in maintenance services Workers: workers, - disableStaggerStart: true, }) require.NoError(t, err) @@ -4009,9 +4010,11 @@ func Test_NewClient_Overrides(t *testing.T) { require.Equal(t, 1*time.Hour, jobCleaner.Config.CancelledJobRetentionPeriod) require.Equal(t, 2*time.Hour, jobCleaner.Config.CompletedJobRetentionPeriod) require.Equal(t, 3*time.Hour, jobCleaner.Config.DiscardedJobRetentionPeriod) + require.True(t, jobCleaner.StaggerStartupIsDisabled()) enqueuer := maintenance.GetService[*maintenance.PeriodicJobEnqueuer](client.queueMaintainer) require.Equal(t, int32(123_456), enqueuer.Config.AdvisoryLockPrefix) + require.True(t, enqueuer.StaggerStartupIsDisabled()) require.Equal(t, errorHandler, client.config.ErrorHandler) require.Equal(t, 123*time.Millisecond, client.config.FetchCooldown) @@ -4020,7 +4023,6 @@ func Test_NewClient_Overrides(t *testing.T) { require.Equal(t, logger, client.baseService.Logger) require.Equal(t, 5, client.config.MaxAttempts) require.Equal(t, retryPolicy, client.config.RetryPolicy) - require.True(t, client.config.disableStaggerStart) } func Test_NewClient_MissingParameters(t *testing.T) { diff --git a/internal/maintenance/queue_maintainer.go b/internal/maintenance/queue_maintainer.go index cf08875d..e99f9ab2 100644 --- a/internal/maintenance/queue_maintainer.go +++ b/internal/maintenance/queue_maintainer.go @@ -124,6 +124,10 @@ func (s *queueMaintainerServiceBase) StaggerStartupDisable(disabled bool) { s.staggerStartupDisabled = disabled } +func (s *queueMaintainerServiceBase) StaggerStartupIsDisabled() bool { + return s.staggerStartupDisabled +} + // withStaggerStartupDisable is an interface to a service whose stagger startup // sleep can be disable. type withStaggerStartupDisable interface {