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

- `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).
Expand Down
21 changes: 13 additions & 8 deletions client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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,
}
Expand Down Expand Up @@ -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)
}
}
Expand Down
26 changes: 14 additions & 12 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{},
}
}

Expand Down Expand Up @@ -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)
Expand All @@ -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) {
Expand Down Expand Up @@ -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)

Expand All @@ -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)
Expand All @@ -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) {
Expand Down
4 changes: 4 additions & 0 deletions internal/maintenance/queue_maintainer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down