middleware

package
v0.25.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 5, 2026 License: MIT Imports: 5 Imported by: 0

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

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.

func (*FailOnException) Handle

func (m *FailOnException) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error

Handle runs the job and parks it on a failure the predicate claims.

type Func

type Func func(ctx context.Context, j *jobs.Job, next func(context.Context) error) error

Func adapts a function to Middleware.

func (Func) Handle

func (f Func) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error

Handle calls f.

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.

func (Skip) Handle

func (m Skip) Handle(ctx context.Context, _ *jobs.Job, next func(context.Context) error) error

Handle drops the job, or hands it on.

func (Skip) Unless

func (Skip) Unless(cond bool) Skip

Unless skips the job unless cond is true.

func (Skip) When

func (Skip) When(cond bool) Skip

When skips the job when cond is true.

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.

func (*SkipIfBatchCancelled) Handle

func (m *SkipIfBatchCancelled) Handle(ctx context.Context, j *jobs.Job, next func(context.Context) error) error

Handle asks the batch, and skips the job if it was cancelled.

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

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

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

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

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

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

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

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL