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 client.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import (
"github.com/riverqueue/river/internal/notifylimiter"
"github.com/riverqueue/river/internal/pluginconfig"
"github.com/riverqueue/river/internal/pluginlookup"
"github.com/riverqueue/river/internal/retrypolicy"
"github.com/riverqueue/river/internal/rivercommon"
"github.com/riverqueue/river/internal/riverplugin"
"github.com/riverqueue/river/internal/workunit"
Expand Down Expand Up @@ -823,6 +824,9 @@ func NewClient[TTx any](driver riverdriver.Driver[TTx], config *Config) (*Client
archetype.Time = &baseservice.TimeGeneratorWithStubWrapper{TimeGenerator: config.Test.Time}
}
}
if _, ok := config.RetryPolicy.(*DefaultClientRetryPolicy); ok {
config.RetryPolicy = retrypolicy.NewDefault(archetype.Time)
}

var (
middleware = pluginconfig.CombinedMiddleware(config.Middleware, config.JobInsertMiddleware, config.WorkerMiddleware)
Expand Down
112 changes: 110 additions & 2 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import (
"github.com/riverqueue/river/internal/maintenance"
"github.com/riverqueue/river/internal/notifier"
"github.com/riverqueue/river/internal/pluginlookup"
"github.com/riverqueue/river/internal/retrypolicy"
"github.com/riverqueue/river/internal/rivercommon"
"github.com/riverqueue/river/internal/riverinternaltest"
"github.com/riverqueue/river/internal/riverinternaltest/retrypolicytest"
Expand Down Expand Up @@ -994,6 +995,75 @@ func Test_Client_Common(t *testing.T) {
require.WithinDuration(t, time.Now(), *updatedJob.FinalizedAt, 2*time.Second)
})

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

config, bundle := setupConfig(t)
configuredNow := time.Now().UTC().Add(-10 * time.Minute).Truncate(time.Microsecond)
timeStub := &riversharedtest.TimeStub{}
timeStub.StubNow(configuredNow)
config.RetryPolicy = &retrypolicytest.RetryPolicyInvalid{}
config.Test.Time = timeStub
client := newTestClient(t, bundle.dbPool, config)

type JobArgs struct {
testutil.JobArgsReflectKind[JobArgs]
}

AddWorker(client.config.Workers, WorkFunc(func(ctx context.Context, job *Job[JobArgs]) error {
return errors.New("retry using configured fallback time")
}))

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

insertRes, err := client.Insert(ctx, &JobArgs{}, nil)
require.NoError(t, err)

event := riversharedtest.WaitOrTimeout(t, subscribeChan)
require.Equal(t, EventKindJobFailed, event.Kind)
require.Equal(t, rivertype.JobStateRetryable, event.Job.State)
require.WithinDuration(t, configuredNow.Add(time.Second), event.Job.ScheduledAt, 150*time.Millisecond)

updatedJob, err := client.JobGet(ctx, insertRes.Job.ID)
require.NoError(t, err)
require.WithinDuration(t, configuredNow.Add(time.Second), updatedJob.ScheduledAt, 150*time.Millisecond)
})

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

config, bundle := setupConfig(t)
configuredNow := time.Now().UTC().Add(-10 * time.Minute).Truncate(time.Microsecond)
timeStub := &riversharedtest.TimeStub{}
timeStub.StubNow(configuredNow)
config.Test.Time = timeStub
client := newTestClient(t, bundle.dbPool, config)

type JobArgs struct {
testutil.JobArgsReflectKind[JobArgs]
}

AddWorker(client.config.Workers, WorkFunc(func(ctx context.Context, job *Job[JobArgs]) error {
return errors.New("retry using configured time")
}))

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

insertRes, err := client.Insert(ctx, &JobArgs{}, nil)
require.NoError(t, err)

event := riversharedtest.WaitOrTimeout(t, subscribeChan)
require.Equal(t, EventKindJobFailed, event.Kind)
require.Equal(t, rivertype.JobStateRetryable, event.Job.State)
require.WithinDuration(t, configuredNow.Add(time.Second), event.Job.ScheduledAt, 150*time.Millisecond)

updatedJob, err := client.JobGet(ctx, insertRes.Job.ID)
require.NoError(t, err)
require.WithinDuration(t, configuredNow.Add(time.Second), updatedJob.ScheduledAt, 150*time.Millisecond)
})

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

Expand Down Expand Up @@ -1025,6 +1095,39 @@ func Test_Client_Common(t *testing.T) {
require.WithinDuration(t, time.Now().Add(15*time.Minute), updatedJob.ScheduledAt, 2*time.Second)
})

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

config, bundle := setupConfig(t)
configuredNow := time.Now().UTC().Add(-10 * time.Minute).Truncate(time.Microsecond)
timeStub := &riversharedtest.TimeStub{}
timeStub.StubNow(configuredNow)
config.Test.Time = timeStub
client := newTestClient(t, bundle.dbPool, config)

type JobArgs struct {
testutil.JobArgsReflectKind[JobArgs]
}

AddWorker(client.config.Workers, WorkFunc(func(ctx context.Context, job *Job[JobArgs]) error {
return JobSnooze(15 * time.Minute)
}))

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

insertRes, err := client.Insert(ctx, &JobArgs{}, nil)
require.NoError(t, err)

event := riversharedtest.WaitOrTimeout(t, subscribeChan)
require.Equal(t, EventKindJobSnoozed, event.Kind)
require.Equal(t, configuredNow.Add(15*time.Minute), event.Job.ScheduledAt)

updatedJob, err := client.JobGet(ctx, insertRes.Job.ID)
require.NoError(t, err)
require.Equal(t, configuredNow.Add(15*time.Minute), updatedJob.ScheduledAt)
})

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

Expand Down Expand Up @@ -8218,7 +8321,7 @@ func Test_NewClient_Defaults(t *testing.T) {
require.NotZero(t, client.baseService.Logger)
require.Equal(t, MaxAttemptsDefault, client.config.MaxAttempts)
require.Equal(t, maintenance.ReindexerTimeoutDefault, client.config.ReindexerTimeout)
require.IsType(t, &DefaultClientRetryPolicy{}, client.config.RetryPolicy)
require.IsType(t, &retrypolicy.Default{}, client.config.RetryPolicy)
require.False(t, client.config.SkipUnknownJobCheck)
require.IsType(t, nil, client.config.Test.Time)
require.IsType(t, &baseservice.UnStubbableTimeGenerator{}, client.baseService.Time)
Expand All @@ -8240,11 +8343,13 @@ func Test_NewClient_Overrides(t *testing.T) {
return JobStuckHandlerResult{}
})
logger := slog.New(slog.NewTextHandler(os.Stderr, nil))
timeStub := &riversharedtest.TimeStub{}
timeStub.StubNow(time.Now().UTC())

workers := NewWorkers()
AddWorker(workers, &noOpWorker{})

retryPolicy := &DefaultClientRetryPolicy{}
retryPolicy := &retrypolicytest.RetryPolicyNoJitter{}

type noOpHook struct {
HookDefaults
Expand Down Expand Up @@ -8280,6 +8385,7 @@ func Test_NewClient_Overrides(t *testing.T) {
RetryPolicy: retryPolicy,
Schema: schema,
SkipUnknownJobCheck: true,
Test: TestConfig{Time: timeStub},
TestOnly: true, // disables staggered start in maintenance services
Workers: workers,
WorkerMiddleware: []rivertype.WorkerMiddleware{&noOpWorkerMiddleware{}},
Expand Down Expand Up @@ -8316,6 +8422,8 @@ func Test_NewClient_Overrides(t *testing.T) {
require.Equal(t, 5, client.config.MaxAttempts)
require.Equal(t, 125*time.Millisecond, client.config.ReindexerTimeout)
require.Equal(t, retryPolicy, client.config.RetryPolicy)
require.Equal(t, logger, retryPolicy.Logger)
require.Same(t, timeStub, retryPolicy.Time)
require.Equal(t, schema, client.config.Schema)
require.True(t, client.config.SkipUnknownJobCheck)
require.Len(t, client.config.WorkerMiddleware, 1)
Expand Down
2 changes: 1 addition & 1 deletion internal/jobexecutor/job_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -383,7 +383,7 @@ func (e *JobExecutor) reportResult(ctx context.Context, jobRow *rivertype.JobRow
slog.String("job_kind", jobRow.Kind),
slog.Duration("duration", snoozeErr.Duration),
)
nextAttemptScheduledAt := time.Now().Add(snoozeErr.Duration)
nextAttemptScheduledAt := e.Time.Now().Add(snoozeErr.Duration)

snoozesValue := gjson.GetBytes(jobRow.Metadata, "snoozes").Int()
if res.MetadataUpdates == nil {
Expand Down
77 changes: 77 additions & 0 deletions internal/retrypolicy/default.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
// Package retrypolicy contains River's internal retry policy implementations.
package retrypolicy

import (
"math"
"math/rand/v2"
"time"

"github.com/riverqueue/river/rivershared/util/timeutil"
"github.com/riverqueue/river/rivertype"
)

// Default is River's clock-aware default retry policy for internal use.
type Default struct {
timeGenerator rivertype.TimeGenerator
}

// NewDefault returns a default retry policy that derives retries from the
// given time generator.
func NewDefault(timeGenerator rivertype.TimeGenerator) *Default {
return &Default{timeGenerator: timeGenerator}
}

// NextRetry calculates when the next retry for a failed job should take place.
func (p *Default) NextRetry(job *rivertype.JobRow) time.Time {
return NextRetryAt(p.timeGenerator.Now().UTC(), job)
}

// NextRetryAt calculates when the next retry for a failed job should take
// place relative to now.
func NextRetryAt(now time.Time, job *rivertype.JobRow) time.Time {
// In modern versions of River `len(job.Errors)` is the same number as
// `attempt`. However, in older versions snoozing a job wouldn't restore its
// attempt count to the pre-fetch value, and that would lead to incorrect
// retry durations when jobs are first snoozed, then retried. To avoid this
// and keep backward compatibility, the number of errors are used instead.
errorCount := len(job.Errors) + 1

return now.Add(timeutil.SecondsAsDuration(retrySeconds(errorCount)))
}

// The maximum value of a duration before it overflows. About 292 years.
const maxDuration time.Duration = 1<<63 - 1

// Same as the above, but changed to a float represented in seconds.
var maxDurationSeconds = maxDuration.Seconds() //nolint:gochecknoglobals

// Gets a number of retry seconds for the given attempt, random jitter included.
func retrySeconds(attempt int) float64 {
retrySeconds := retrySecondsWithoutJitter(attempt)

// After hitting maximum retry durations jitter is no longer applied because
// it might overflow time.Duration. That's okay though because so much
// jitter will already have been applied up to this point (jitter measured
// in decades) that jobs will no longer run anywhere near contemporaneously
// unless there's been considerable manual intervention.
if retrySeconds == maxDurationSeconds {
return maxDurationSeconds
}

// Jitter number of seconds +/- 10%.
retrySeconds += retrySeconds * (rand.Float64()*0.2 - 0.1)

// Cap retrySeconds once more in case adding random jitter pushed it over
// maxDurationSeconds. (This should never realistically happen, but protect
// against it just in case.)
return min(retrySeconds, maxDurationSeconds)
}

// Gets a base number of retry seconds for the given attempt, jitter excluded.
// If the number of seconds returned would overflow time.Duration if it were to
// be made one, returns the maximum number of seconds that can fit in a
// time.Duration instead, approximately 292 years.
func retrySecondsWithoutJitter(attempt int) float64 {
retrySeconds := math.Pow(float64(attempt), 4)
return min(retrySeconds, maxDurationSeconds)
}
Loading
Loading