Documentation
¶
Overview ¶
Package middleware wraps the handling of a job.
A middleware sits between the worker and the handler and can decide the work should not happen now -- another worker holds the lock, the budget for this minute is spent, the batch was cancelled -- by putting the job back on its queue instead of running it.
Middleware.Handle gets the context, the job and the rest of the chain, and what it does with the job is release it, delete it, fail it, or hand it on.
RateLimited and ThrottlesExceptions count in a cache.Store, so which store they count in is wiring rather than a second type: the RESP-backed version of RateLimited is RateLimited over a RESP-backed store.
Every key is scoped by tenant ¶
The lock a job takes and the counter a job spends are named after the tenant the job belongs to. Without it, one customer's slow import would rate limit every other customer's.
Index ¶
- type FailOnException
- type Func
- type Middleware
- type RateLimited
- type Skip
- type SkipIfBatchCancelled
- type ThrottlesExceptions
- func (m *ThrottlesExceptions) Backoff(d time.Duration) *ThrottlesExceptions
- func (m *ThrottlesExceptions) By(key string) *ThrottlesExceptions
- func (m *ThrottlesExceptions) ByJob() *ThrottlesExceptions
- func (m *ThrottlesExceptions) DeleteWhen(f func(error) bool) *ThrottlesExceptions
- func (m *ThrottlesExceptions) FailWhen(f func(error) bool) *ThrottlesExceptions
- func (m *ThrottlesExceptions) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error
- func (m *ThrottlesExceptions) Report(f func(error) bool) *ThrottlesExceptions
- func (m *ThrottlesExceptions) RetryAfter(d time.Duration) *ThrottlesExceptions
- func (m *ThrottlesExceptions) ShouldReport(err error) bool
- func (m *ThrottlesExceptions) When(f func(error) bool) *ThrottlesExceptions
- func (m *ThrottlesExceptions) WithPrefix(prefix string) *ThrottlesExceptions
- type WithoutOverlapping
- func (m *WithoutOverlapping) DontRelease() *WithoutOverlapping
- func (m *WithoutOverlapping) ExpireAfter(d time.Duration) *WithoutOverlapping
- func (m *WithoutOverlapping) GetLockKey(j *jobs.Job) string
- func (m *WithoutOverlapping) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error
- func (m *WithoutOverlapping) ReleaseAfter(d time.Duration) *WithoutOverlapping
- func (m *WithoutOverlapping) Shared() *WithoutOverlapping
- func (m *WithoutOverlapping) WithPrefix(prefix string) *WithoutOverlapping
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type FailOnException ¶
type FailOnException struct {
// contains filtered or unexported fields
}
FailOnException parks a job on a failure that retrying cannot fix.
The default schedule retries everything the same number of times, which is right for a timeout and wrong for a payload that will never parse: five deliveries of a job that cannot succeed is five copies of the same error, an hour apart.
m := middleware.NewFailOnException(func(err error) bool {
return errors.Is(err, database.ErrNotFound)
})
The failure is still returned to the worker, so the log line and the recorded job are the same as any other failure. The job is already parked by then, and the worker does not park it twice.
func NewFailOnException ¶
func NewFailOnException(when func(error) bool) *FailOnException
NewFailOnException returns the middleware.
The set of failures it parks on is errors.Is inside the predicate, written where the compiler can read it.
type Middleware ¶
type Middleware interface {
Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error
}
Middleware wraps the handling of one job.
Call next to hand the job on. Not calling it is how a middleware says the work should not happen: release the job to try again later, delete it to drop it, or return without touching it -- which the worker reads as "handled", and the job is deleted.
type RateLimited ¶
type RateLimited struct {
// contains filtered or unexported fields
}
RateLimited runs a job only while the budget for it lasts.
The job over budget is released and comes back when the window rolls, so the work is delayed rather than dropped.
m := middleware.NewRateLimited(limiter, cache.PerMinute(30))
The limit's key is scoped by tenant, so one customer cannot spend another's budget. An empty key means the job's name, which is the "thirty of these a minute, per customer" case.
There is no RESP-specific variant: the limiter counts in whatever cache.Store it was built over.
func NewRateLimited ¶
func NewRateLimited(limiter *cache.RateLimiter, limit cache.Limit) *RateLimited
NewRateLimited returns the middleware.
func (*RateLimited) DontRelease ¶
func (m *RateLimited) DontRelease() *RateLimited
DontRelease drops the job over budget instead of putting it back.
func (*RateLimited) Handle ¶
func (m *RateLimited) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error
Handle spends one attempt, and runs the job if it fit.
func (*RateLimited) ReleaseAfter ¶
func (m *RateLimited) ReleaseAfter(d time.Duration) *RateLimited
ReleaseAfter fixes how long a job over budget waits, instead of asking the limiter when the window rolls.
type Skip ¶
type Skip struct {
// contains filtered or unexported fields
}
Skip drops a job without running it.
The job is not released and not failed: the worker reads a middleware that returned without touching the job as "handled", so the job is deleted. Skipping means the work is not wanted, not that it should be tried later -- for later, release it.
w := queue.NewWorker(q, queue.WorkerOptions{
Middleware: []middleware.Middleware{middleware.Skip{}.Unless(featureOn)},
})
Skip.When and Skip.Unless are methods on the zero value, so `middleware.Skip{}.When(cond)` reads as one phrase.
type SkipIfBatchCancelled ¶
type SkipIfBatchCancelled struct {
// contains filtered or unexported fields
}
SkipIfBatchCancelled drops a job whose batch was cancelled.
Cancelling a batch cannot recall the jobs already on the queue, so they are delivered, they ask, and they skip -- which is what makes cancelling mean anything.
A job that belongs to no batch is handed on and the store is never consulted, so this middleware is safe to put on a worker that also runs jobs pushed on their own.
func NewSkipIfBatchCancelled ¶
func NewSkipIfBatchCancelled(batches bus.BatchRepository) *SkipIfBatchCancelled
NewSkipIfBatchCancelled returns the middleware over the repository the batches live in.
type ThrottlesExceptions ¶
type ThrottlesExceptions struct {
// contains filtered or unexported fields
}
ThrottlesExceptions stops hammering a dependency that is already failing.
It counts failures rather than attempts: once a job has failed more than the limit allows inside the window, the ones that follow are released without being run at all, so a third-party API that is down gets a pause instead of the whole queue retrying against it.
m := middleware.NewThrottlesExceptions(limiter, cache.PerMinute(10)). RetryAfter(5 * time.Minute)
A failure it catches is not returned to the worker: the job is released and tries again, which is the point -- the worker's own attempt counter is for jobs that are wrong, and this is for dependencies that are down.
func NewThrottlesExceptions ¶
func NewThrottlesExceptions(limiter *cache.RateLimiter, limit cache.Limit) *ThrottlesExceptions
NewThrottlesExceptions returns the middleware.
func (*ThrottlesExceptions) Backoff ¶
func (m *ThrottlesExceptions) Backoff(d time.Duration) *ThrottlesExceptions
Backoff is how long the jobs that arrive while the limit is spent wait before they are eligible again.
It answers backoff(). It is the pause on the closed circuit, where RetryAfter is the wait after the failure that closed it: the first says how long the dependency gets to recover, the second how soon the job that noticed tries again.
func (*ThrottlesExceptions) By ¶
func (m *ThrottlesExceptions) By(key string) *ThrottlesExceptions
By counts against a key of your own instead of the job's name.
It answers by(). Use it when the thing that is failing is neither the job nor the record but something they share: an account, a region, a tenant of the third party.
func (*ThrottlesExceptions) ByJob ¶
func (m *ThrottlesExceptions) ByJob() *ThrottlesExceptions
ByJob counts each job separately instead of counting the whole name together.
It answers byJob(). Use it when the dependency is per-record rather than shared: one broken invoice should not throttle the other nine hundred.
func (*ThrottlesExceptions) DeleteWhen ¶
func (m *ThrottlesExceptions) DeleteWhen(f func(error) bool) *ThrottlesExceptions
DeleteWhen drops the job instead of releasing it, for the failures the predicate claims.
It answers deleteWhen(). It is for the failure that means the work is pointless rather than early: the record was deleted while the job waited, and releasing it would put it back for four more rounds of the same answer.
func (*ThrottlesExceptions) FailWhen ¶
func (m *ThrottlesExceptions) FailWhen(f func(error) bool) *ThrottlesExceptions
FailWhen parks the job instead of releasing it, for the failures the predicate claims.
It answers failWhen(). It is DeleteWhen's louder sibling: the work is pointless and somebody should see that it was.
func (*ThrottlesExceptions) Handle ¶
func (m *ThrottlesExceptions) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error
Handle refuses while the failures are piling up, and counts the ones it sees.
func (*ThrottlesExceptions) Report ¶
func (m *ThrottlesExceptions) Report(f func(error) bool) *ThrottlesExceptions
Report narrows which failures are worth reporting.
A dependency that is down produces one failure per job, and reporting every one of them is how the report becomes noise: the predicate is what says "only the first", or "only the ones that are not timeouts". A nil predicate reports everything.
It does not report anything itself. What it decides is what ThrottlesExceptions.ShouldReport answers, and the caller's error reporter is what acts on it -- this package has no business knowing where a report goes.
func (*ThrottlesExceptions) RetryAfter ¶
func (m *ThrottlesExceptions) RetryAfter(d time.Duration) *ThrottlesExceptions
RetryAfter is how long a job waits after a failure this middleware caught.
func (*ThrottlesExceptions) ShouldReport ¶
func (m *ThrottlesExceptions) ShouldReport(err error) bool
ShouldReport reports whether a failure is worth reporting.
It is the read half of ThrottlesExceptions.Report. This package does not report anything itself, so the caller asks and acts.
func (*ThrottlesExceptions) When ¶
func (m *ThrottlesExceptions) When(f func(error) bool) *ThrottlesExceptions
When narrows which failures are counted. A failure it says no to is returned to the worker untouched, which parks the job on the usual schedule.
func (*ThrottlesExceptions) WithPrefix ¶
func (m *ThrottlesExceptions) WithPrefix(prefix string) *ThrottlesExceptions
WithPrefix puts a prefix in front of the counter's key.
It answers withPrefix(). Two middlewares over the same job name would otherwise share one counter, and the failures of one would throttle the other.
type WithoutOverlapping ¶
type WithoutOverlapping struct {
// contains filtered or unexported fields
}
WithoutOverlapping runs at most one job with a given key at a time.
The usual key is the id of the thing being worked on -- an account, an invoice, a report -- so two jobs about the same row never run at once while two jobs about different rows run in parallel.
m := middleware.NewWithoutOverlapping(locks, accountID).ReleaseAfter(10 * time.Second)
A job that cannot take the lock is released and tries again later. Call WithoutOverlapping.DontRelease to drop it instead, which is right when the work is "make this true" rather than "do this once".
func NewWithoutOverlapping ¶
func NewWithoutOverlapping(locks *cache.Locks, key string) *WithoutOverlapping
NewWithoutOverlapping returns the middleware.
An empty key means the job's name, which is the "one of these at a time" case.
func (*WithoutOverlapping) DontRelease ¶
func (m *WithoutOverlapping) DontRelease() *WithoutOverlapping
DontRelease drops the job instead of putting it back when the lock is taken.
func (*WithoutOverlapping) ExpireAfter ¶
func (m *WithoutOverlapping) ExpireAfter(d time.Duration) *WithoutOverlapping
ExpireAfter sets the lock's ttl, which is the deadlock protection: a worker that dies holding it releases it when the ttl passes.
func (*WithoutOverlapping) GetLockKey ¶
func (m *WithoutOverlapping) GetLockKey(j *jobs.Job) string
GetLockKey is the name of the lock this job takes.
func (*WithoutOverlapping) Handle ¶
func (m *WithoutOverlapping) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error
Handle takes the lock, runs the job, and gives the lock back.
func (*WithoutOverlapping) ReleaseAfter ¶
func (m *WithoutOverlapping) ReleaseAfter(d time.Duration) *WithoutOverlapping
ReleaseAfter sets how long a job waits before trying for the lock again.
func (*WithoutOverlapping) Shared ¶
func (m *WithoutOverlapping) Shared() *WithoutOverlapping
Shared makes jobs with different names share the lock, so long as they share the key.
func (*WithoutOverlapping) WithPrefix ¶
func (m *WithoutOverlapping) WithPrefix(prefix string) *WithoutOverlapping
WithPrefix names the lock's namespace.