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
6 changes: 2 additions & 4 deletions internal/jobcompleter/job_completer.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
"github.com/riverqueue/river/rivershared/riverpilot"
"github.com/riverqueue/river/rivershared/startstop"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivertype"
)

Expand Down Expand Up @@ -683,10 +684,7 @@ func withRetries[T any](logCtx context.Context, baseService *baseservice.BaseSer
for attempt := 1; attempt <= numRetries; attempt++ {
// I've found that we want at least ten seconds for a large batch,
// although it usually doesn't need that long.
ctx, cancel := context.WithTimeout(uncancelledCtx, rivercommon.HotOperationTimeout)
defer cancel()

retVal, err := retryFunc(ctx)
retVal, err := timeoututil.WithTimeoutV(uncancelledCtx, rivercommon.HotOperationTimeout, baseService.Name+".withRetries", retryFunc)
if err != nil {
// A cancelled context or a closed pool will never succeed.
if isNonRetryableCompleterError(err) {
Expand Down
98 changes: 48 additions & 50 deletions internal/leadership/elector.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"github.com/riverqueue/river/rivershared/util/randutil"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivertype"
)

Expand Down Expand Up @@ -609,25 +610,24 @@ func (e *Elector) attemptResign(ctx context.Context, attempt int, term leadershi
// Wait one second longer each time we try to resign:
timeout := time.Duration(attempt) * time.Second

ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()

resigned, err := e.exec.LeaderResign(ctx, &riverdriver.LeaderResignParams{
ElectedAt: term.electedAt,
LeaderID: term.clientID,
LeadershipTopic: string(notifier.NotificationTopicLeadership),
Schema: e.config.Schema,
})
if err != nil {
return err
}
return timeoututil.WithTimeout(ctx, timeout, e.Name+".attemptResign", func(ctx context.Context) error {
resigned, err := e.exec.LeaderResign(ctx, &riverdriver.LeaderResignParams{
ElectedAt: term.electedAt,
LeaderID: term.clientID,
LeadershipTopic: string(notifier.NotificationTopicLeadership),
Schema: e.config.Schema,
})
if err != nil {
return err
}

if resigned {
e.Logger.DebugContext(ctx, e.Name+": Resigned leadership successfully", "client_id", e.config.ClientID)
e.testSignals.ResignedLeadership.Signal(struct{}{})
}
if resigned {
e.Logger.DebugContext(ctx, e.Name+": Resigned leadership successfully", "client_id", e.config.ClientID)
e.testSignals.ResignedLeadership.Signal(struct{}{})
}

return nil
return nil
})
}

// Produces a common set of key/value pairs for logging when an error occurs.
Expand Down Expand Up @@ -767,44 +767,42 @@ func attemptElect(ctx context.Context, exec riverdriver.Executor, params *riverd
}

func attemptElectWithTimeout(ctx context.Context, exec riverdriver.Executor, params *riverdriver.LeaderElectParams, timeout time.Duration) (*riverdriver.Leader, error) {
ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()

execTx, err := exec.Begin(ctx)
if err != nil {
var additionalDetail string
if errors.Is(err, context.DeadlineExceeded) {
additionalDetail = " (a common cause of this is a database pool that's at its connection limit; you may need to increase maximum connections)"
}
return timeoututil.WithTimeoutV(ctx, timeout, "leadership.attemptElect", func(ctx context.Context) (*riverdriver.Leader, error) {
execTx, err := exec.Begin(ctx)
if err != nil {
var additionalDetail string
if errors.Is(err, context.DeadlineExceeded) {
additionalDetail = " (a common cause of this is a database pool that's at its connection limit; you may need to increase maximum connections)"
}

return nil, fmt.Errorf("error beginning transaction: %w%s", err, additionalDetail)
}
defer dbutil.RollbackWithoutCancel(ctx, execTx)
return nil, fmt.Errorf("error beginning transaction: %w%s", err, additionalDetail)
}
defer dbutil.RollbackWithoutCancel(ctx, execTx)

if _, err := execTx.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Now: params.Now,
Schema: params.Schema,
}); err != nil {
return nil, err
}
if _, err := execTx.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Now: params.Now,
Schema: params.Schema,
}); err != nil {
return nil, err
}

leader, err := execTx.LeaderAttemptElect(ctx, params)
if err != nil && !errors.Is(err, rivertype.ErrNotFound) {
return nil, err
}
if err := execTx.Commit(ctx); err != nil {
return nil, fmt.Errorf("error committing transaction: %w", err)
}
if err != nil {
return nil, err
}
leader, err := execTx.LeaderAttemptElect(ctx, params)
if err != nil && !errors.Is(err, rivertype.ErrNotFound) {
return nil, err
}
if err := execTx.Commit(ctx); err != nil {
return nil, fmt.Errorf("error committing transaction: %w", err)
}
if err != nil {
return nil, err
}

return leader, nil
return leader, nil
})
}

func attemptReelectWithTimeout(ctx context.Context, exec riverdriver.Executor, params *riverdriver.LeaderReelectParams, timeout time.Duration) (*riverdriver.Leader, error) {
ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()

return exec.LeaderAttemptReelect(ctx, params)
return timeoututil.WithTimeoutV(ctx, timeout, "leadership.attemptReelect", func(ctx context.Context) (*riverdriver.Leader, error) {
return exec.LeaderAttemptReelect(ctx, params)
})
}
9 changes: 3 additions & 6 deletions internal/maintenance/job_cleaner.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"github.com/riverqueue/river/rivershared/util/randutil"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivershared/util/timeutil"
)

Expand Down Expand Up @@ -183,8 +184,7 @@ func (s *JobCleaner) runOnce(ctx context.Context) (*jobCleanerRunOnceResult, err
res := &jobCleanerRunOnceResult{}

for {
// Wrapped in a function so that defers run as expected.
numDeleted, err := func() (int, error) {
numDeleted, err := timeoututil.WithTimeoutV(ctx, s.Config.Timeout, s.Name+".runOnce", func(ctx context.Context) (int, error) {
// In the special case that all retentions are indefinite, don't
// bother issuing the query at all as an optimization.
if s.Config.CompletedJobRetentionPeriod == -1 &&
Expand All @@ -193,9 +193,6 @@ func (s *JobCleaner) runOnce(ctx context.Context) (*jobCleanerRunOnceResult, err
return 0, nil
}

ctx, cancelFunc := context.WithTimeout(ctx, s.Config.Timeout)
defer cancelFunc()

numDeleted, err := s.exec.JobDeleteBefore(ctx, &riverdriver.JobDeleteBeforeParams{
CancelledDoDelete: s.Config.CancelledJobRetentionPeriod != -1,
CancelledFinalizedAtHorizon: time.Now().Add(-s.Config.CancelledJobRetentionPeriod),
Expand All @@ -214,7 +211,7 @@ func (s *JobCleaner) runOnce(ctx context.Context) (*jobCleanerRunOnceResult, err
s.reducedBatchSizeBreaker.ResetIfNotOpen()

return numDeleted, nil
}()
})
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
s.reducedBatchSizeBreaker.Trip()
Expand Down
20 changes: 10 additions & 10 deletions internal/maintenance/job_rescuer.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"github.com/riverqueue/river/rivershared/util/randutil"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivershared/util/timeutil"
"github.com/riverqueue/river/rivertype"
)
Expand Down Expand Up @@ -294,24 +295,23 @@ func (s *JobRescuer) runOnce(ctx context.Context) (*rescuerRunOnceResult, error)
}

func (s *JobRescuer) getStuckJobs(ctx context.Context, afterID int64, batchSize int, stuckHorizon time.Time) ([]*rivertype.JobRow, error) {
ctx, cancelFunc := context.WithTimeout(ctx, riversharedmaintenance.TimeoutDefault)
defer cancelFunc()

params := &riverdriver.JobGetStuckParams{
AfterID: afterID,
Max: batchSize,
Schema: s.Config.Schema,
StuckHorizon: stuckHorizon,
}

if pilot, ok := s.Config.Pilot.(riverpilot.PilotJobRescuer); ok {
return pilot.JobGetStuck(ctx, s.exec, params)
}
return timeoututil.WithTimeoutV(ctx, riversharedmaintenance.TimeoutDefault, s.Name+".getStuckJobs", func(ctx context.Context) ([]*rivertype.JobRow, error) {
if pilot, ok := s.Config.Pilot.(riverpilot.PilotJobRescuer); ok {
return pilot.JobGetStuck(ctx, s.exec, params)
}

// Compatibility fallback for Pilot implementations from before
// PilotJobRescuer. Once Pilot embeds PilotJobRescuer, replace the assertion
// above and this fallback with a direct call to s.Config.Pilot.JobGetStuck.
return s.exec.JobGetStuck(ctx, params)
// Compatibility fallback for Pilot implementations from before
// PilotJobRescuer. Once Pilot embeds PilotJobRescuer, replace the assertion
// above and this fallback with a direct call to s.Config.Pilot.JobGetStuck.
return s.exec.JobGetStuck(ctx, params)
})
}

// jobRetryDecision is a signal from makeRetryDecision as to what to do with a
Expand Down
9 changes: 3 additions & 6 deletions internal/maintenance/job_scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"github.com/riverqueue/river/rivershared/util/randutil"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivershared/util/timeutil"
)

Expand Down Expand Up @@ -163,11 +164,7 @@ func (s *JobScheduler) runOnce(ctx context.Context) (*schedulerRunOnceResult, er
res := &schedulerRunOnceResult{}

for {
// Wrapped in a function so that defers run as expected.
numScheduled, err := func() (int, error) {
ctx, cancelFunc := context.WithTimeout(ctx, riversharedmaintenance.TimeoutDefault)
defer cancelFunc()

numScheduled, err := timeoututil.WithTimeoutV(ctx, riversharedmaintenance.TimeoutDefault, s.Name+".runOnce", func(ctx context.Context) (int, error) {
execTx, err := s.exec.Begin(ctx)
if err != nil {
return 0, fmt.Errorf("error starting transaction: %w", err)
Expand Down Expand Up @@ -212,7 +209,7 @@ func (s *JobScheduler) runOnce(ctx context.Context) (*schedulerRunOnceResult, er
}

return len(scheduledJobResults), execTx.Commit(ctx)
}()
})
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
s.reducedBatchSizeBreaker.Trip()
Expand Down
9 changes: 3 additions & 6 deletions internal/maintenance/queue_cleaner.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"github.com/riverqueue/river/rivershared/util/randutil"
"github.com/riverqueue/river/rivershared/util/serviceutil"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivershared/util/timeutil"
)

Expand Down Expand Up @@ -156,11 +157,7 @@ func (s *QueueCleaner) runOnce(ctx context.Context) (*queueCleanerRunOnceResult,
res := &queueCleanerRunOnceResult{QueuesDeleted: make([]string, 0, 10)}

for {
// Wrapped in a function so that defers run as expected.
queuesDeleted, err := func() ([]string, error) {
ctx, cancelFunc := context.WithTimeout(ctx, riversharedmaintenance.TimeoutDefault)
defer cancelFunc()

queuesDeleted, err := timeoututil.WithTimeoutV(ctx, riversharedmaintenance.TimeoutDefault, s.Name+".runOnce", func(ctx context.Context) ([]string, error) {
queuesDeleted, err := s.exec.QueueDeleteExpired(ctx, &riverdriver.QueueDeleteExpiredParams{
Max: s.batchSize(),
Schema: s.Config.Schema,
Expand All @@ -173,7 +170,7 @@ func (s *QueueCleaner) runOnce(ctx context.Context) (*queueCleanerRunOnceResult,
s.reducedBatchSizeBreaker.ResetIfNotOpen()

return queuesDeleted, nil
}()
})
if err != nil {
if errors.Is(err, context.DeadlineExceeded) {
s.reducedBatchSizeBreaker.Trip()
Expand Down
28 changes: 14 additions & 14 deletions internal/maintenance/sqlite_notification_cleaner.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"github.com/riverqueue/river/rivershared/startstop"
"github.com/riverqueue/river/rivershared/testsignal"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/timeoututil"
"github.com/riverqueue/river/rivershared/util/timeutil"
)

Expand Down Expand Up @@ -133,20 +134,19 @@ type sqliteNotificationCleanerRunOnceResult struct {
}

func (s *SQLiteNotificationCleaner) runOnce(ctx context.Context) (*sqliteNotificationCleanerRunOnceResult, error) {
ctx, cancelFunc := context.WithTimeout(ctx, s.Config.Timeout)
defer cancelFunc()

numDeleted, err := s.exec.NotificationDeleteBefore(ctx, &riverdriver.NotificationDeleteBeforeParams{
CreatedAtHorizon: time.Now().Add(-s.Config.RetentionPeriod),
Schema: s.Config.Schema,
})
if err != nil {
return nil, err
}
return timeoututil.WithTimeoutV(ctx, s.Config.Timeout, s.Name+".runOnce", func(ctx context.Context) (*sqliteNotificationCleanerRunOnceResult, error) {
numDeleted, err := s.exec.NotificationDeleteBefore(ctx, &riverdriver.NotificationDeleteBeforeParams{
CreatedAtHorizon: time.Now().Add(-s.Config.RetentionPeriod),
Schema: s.Config.Schema,
})
if err != nil {
return nil, err
}

s.TestSignals.DeletedBatch.Signal(struct{}{})
s.TestSignals.DeletedBatch.Signal(struct{}{})

return &sqliteNotificationCleanerRunOnceResult{
NumNotificationsDeleted: numDeleted,
}, nil
return &sqliteNotificationCleanerRunOnceResult{
NumNotificationsDeleted: numDeleted,
}, nil
})
}
11 changes: 11 additions & 0 deletions internal/maintenance/sqlite_notification_cleaner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,4 +102,15 @@ func TestSQLiteNotificationCleaner(t *testing.T) {

startstoptest.Stress(ctx, t, cleaner)
})

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

cleaner, _ := setup(t)
cleaner.Config.Timeout = time.Nanosecond

_, err := cleaner.runOnce(ctx)
require.ErrorContains(t, err, cleaner.Name+".runOnce timed out after 1ns")
require.ErrorIs(t, err, context.DeadlineExceeded)
})
}
Loading