From f0eb3cf28a64722b2beed97c2aef25e991f67783 Mon Sep 17 00:00:00 2001 From: Blake Gentry Date: Sat, 8 Jun 2024 13:27:20 -0500 Subject: [PATCH 1/2] expose DisableStaggerStart via TestOnlyConfig --- client.go | 24 ++++++++++++++++-------- client_test.go | 24 ++++++++++++------------ 2 files changed, 28 insertions(+), 20 deletions(-) diff --git a/client.go b/client.go index 88efec97..964095ec 100644 --- a/client.go +++ b/client.go @@ -200,6 +200,10 @@ type Config struct { // Defaults to DefaultRetryPolicy. RetryPolicy ClientRetryPolicy + // TestOnly holds configuration for test-only settings that should generally + // not be used outside of test suites. + TestOnly TestOnlyConfig + // Workers is a bundle of registered job workers. // // This field may be omitted for a program that's only enqueueing jobs @@ -208,12 +212,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 @@ -295,6 +293,16 @@ type QueueConfig struct { MaxWorkers int } +// TestOnlyConfig contains test-only settings that should generally not be used +// outside of test suites. +type TestOnlyConfig struct { + // DisableStaggerStart 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 +} + // Client is a single isolated instance of River. Your application may use // multiple instances operating on different databases or Postgres schemas // within a single database. @@ -447,8 +455,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 +615,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.DisableStaggerStart { client.queueMaintainer.StaggerStartupDisable(true) } } diff --git a/client_test.go b/client_test.go index 009d1a27..79cd7017 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: TestOnlyConfig{DisableStaggerStart: true}, // disables staggered start in maintenance services + Workers: workers, + schedulerInterval: riverinternaltest.SchedulerShortInterval, + time: &riverinternaltest.TimeStub{}, } } @@ -3968,7 +3968,7 @@ 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) + require.False(t, client.config.TestOnly.DisableStaggerStart) } func Test_NewClient_Overrides(t *testing.T) { @@ -3998,8 +3998,8 @@ func Test_NewClient_Overrides(t *testing.T) { MaxAttempts: 5, Queues: map[string]QueueConfig{QueueDefault: {MaxWorkers: 1}}, RetryPolicy: retryPolicy, + TestOnly: TestOnlyConfig{DisableStaggerStart: true}, // disables staggered start in maintenance services Workers: workers, - disableStaggerStart: true, }) require.NoError(t, err) @@ -4020,7 +4020,7 @@ 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) + require.True(t, client.config.TestOnly.DisableStaggerStart) } func Test_NewClient_MissingParameters(t *testing.T) { From 70a34c8f9c4265e26205e72d8ac5a08a5c21fa94 Mon Sep 17 00:00:00 2001 From: Brandur Date: Tue, 2 Jul 2024 19:39:42 -0700 Subject: [PATCH 2/2] Expose `Config.TestOnly` configuration option for disabling staggered start This one's a continuation of #385. Maintenance services have a staggered start feature that causes them to sleep for a random amount of jittered time on startup so they don't all try to work simultaneously. This is useful in production, but somewhat harmful in tests because it makes start and stop slower and thereby integration test cases slower. River has an internal flag that allows staggered start to be disabled in its own test suite, but external users of River have no way to access this functionality. Here, introduce `Config.TestOnly` that can be provided to client configuration in test suites regardless of whether the caller is internal or not. This differs slightly from #385 in that it provides only a boolean, with the idea being that if we find it useful to disable other features for tests in the future, a boolean keeps third party code for forwards compatible in that they get these disabled automatically. In case it does become important to distinguish between individual features at some later time, I figure we can add an additional `TestOnlyConfig` property that allows full configuration beyond defaults. --- CHANGELOG.md | 4 ++++ client.go | 25 +++++++++++------------- client_test.go | 10 ++++++---- internal/maintenance/queue_maintainer.go | 4 ++++ 4 files changed, 25 insertions(+), 18 deletions(-) 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 964095ec..458cf8d9 100644 --- a/client.go +++ b/client.go @@ -200,9 +200,16 @@ type Config struct { // Defaults to DefaultRetryPolicy. RetryPolicy ClientRetryPolicy - // TestOnly holds configuration for test-only settings that should generally - // not be used outside of test suites. - TestOnly TestOnlyConfig + // 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. // @@ -293,16 +300,6 @@ type QueueConfig struct { MaxWorkers int } -// TestOnlyConfig contains test-only settings that should generally not be used -// outside of test suites. -type TestOnlyConfig struct { - // DisableStaggerStart 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 -} - // Client is a single isolated instance of River. Your application may use // multiple instances operating on different databases or Postgres schemas // within a single database. @@ -615,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.TestOnly.DisableStaggerStart { + if config.TestOnly { client.queueMaintainer.StaggerStartupDisable(true) } } diff --git a/client_test.go b/client_test.go index 79cd7017..9df4a20a 100644 --- a/client_test.go +++ b/client_test.go @@ -142,7 +142,7 @@ func newTestConfig(t *testing.T, callback callbackFunc) *Config { Logger: riverinternaltest.Logger(t), MaxAttempts: MaxAttemptsDefault, Queues: map[string]QueueConfig{QueueDefault: {MaxWorkers: 50}}, - TestOnly: TestOnlyConfig{DisableStaggerStart: true}, // disables staggered start in maintenance services + 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.TestOnly.DisableStaggerStart) } func Test_NewClient_Overrides(t *testing.T) { @@ -3998,7 +3999,7 @@ func Test_NewClient_Overrides(t *testing.T) { MaxAttempts: 5, Queues: map[string]QueueConfig{QueueDefault: {MaxWorkers: 1}}, RetryPolicy: retryPolicy, - TestOnly: TestOnlyConfig{DisableStaggerStart: true}, // disables staggered start in maintenance services + TestOnly: true, // disables staggered start in maintenance services Workers: workers, }) 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.TestOnly.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 {