Documentation
¶
Overview ¶
Package queue runs jobs in the background: typed structs, dispatched from a request (or anywhere), kept in a store and run by workers, with retries, backoff, timeouts and a record of the jobs that failed.
Set it up at startup, register the job types, and start the workers:
q, err := queue.ForApp(app, redis.QueueDriver()) // QUEUE_DRIVER picks the store
err = queue.Register[jobs.SendWelcome](q, queue.Tries(5))
err = q.Work(queue.Queues("emails", "default"), queue.Concurrency(10))
Then dispatch jobs with the context of a request or another job:
err := queue.Dispatch(ctx, jobs.SendWelcome{UserID: u.ID}, queue.OnQueue("emails"))
The drivers: sync (runs jobs at once, for development and tests), memory (in the process), database (the app's database; dispatches join its transaction), and redis (in drivers/redis).
Delivery is at-least-once: a job may run more than once, so jobs must be idempotent. A failed attempt is retried after a backoff until the job has used its tries; then it is kept as failed, and the queue:failed, queue:retry, queue:forget and queue:flush commands list, retry and delete failed jobs.
Index ¶
- Variables
- func CountFailed(ctx context.Context, st Store) (int64, error)
- func CreateTables(s *migrate.Schema, table, failedTable string) error
- func Dispatch(ctx context.Context, job Job, opts ...DispatchOption) error
- func DispatchFunc(ctx context.Context, name string, payload any, opts ...DispatchOption) error
- func IsPermanent(err error) bool
- func Migrations(table, failedTable string) *migrate.Set
- func Permanent(err error) error
- func Register[J Job](q *Queue, opts ...JobOption) error
- func RegisterFunc[T any](q *Queue, name string, fn func(ctx context.Context, payload T) error, ...) error
- func WithQueue(ctx context.Context, q *Queue) context.Context
- type Config
- type DatabaseStore
- func (s *DatabaseStore) Clear(ctx context.Context, queue string) (int64, error)
- func (s *DatabaseStore) Close() error
- func (s *DatabaseStore) CountFailed(ctx context.Context) (int64, error)
- func (s *DatabaseStore) Delete(ctx context.Context, r *Reservation) error
- func (s *DatabaseStore) Fail(ctx context.Context, r *Reservation, errMsg string) error
- func (s *DatabaseStore) Failed(ctx context.Context, offset, limit int) ([]FailedJob, error)
- func (s *DatabaseStore) FindFailed(ctx context.Context, id string) (FailedJob, bool, error)
- func (s *DatabaseStore) Flush(ctx context.Context) (int64, error)
- func (s *DatabaseStore) Forget(ctx context.Context, id string) (bool, error)
- func (s *DatabaseStore) Push(ctx context.Context, m Message, delay time.Duration) error
- func (s *DatabaseStore) Release(ctx context.Context, r *Reservation, delay time.Duration, refund bool) error
- func (s *DatabaseStore) Reserve(ctx context.Context, queue string, lease time.Duration) (*Reservation, error)
- func (s *DatabaseStore) Retry(ctx context.Context, id string) (bool, error)
- func (s *DatabaseStore) Size(ctx context.Context, queue string) (int64, error)
- type DispatchOption
- type Dispatched
- type Driver
- type FailedCounter
- type FailedFinder
- type FailedJob
- type Failer
- type Info
- type Job
- type JobOption
- type MemoryStore
- func (s *MemoryStore) Clear(_ context.Context, queue string) (int64, error)
- func (s *MemoryStore) Close() error
- func (s *MemoryStore) CountFailed(context.Context) (int64, error)
- func (s *MemoryStore) Delete(_ context.Context, r *Reservation) error
- func (s *MemoryStore) Fail(_ context.Context, r *Reservation, errMsg string) error
- func (s *MemoryStore) Failed(_ context.Context, offset, limit int) ([]FailedJob, error)
- func (s *MemoryStore) FindFailed(_ context.Context, id string) (FailedJob, bool, error)
- func (s *MemoryStore) Flush(context.Context) (int64, error)
- func (s *MemoryStore) Forget(_ context.Context, id string) (bool, error)
- func (s *MemoryStore) Push(_ context.Context, m Message, delay time.Duration) error
- func (s *MemoryStore) Release(_ context.Context, r *Reservation, delay time.Duration, refund bool) error
- func (s *MemoryStore) Reserve(_ context.Context, queue string, lease time.Duration) (*Reservation, error)
- func (s *MemoryStore) Retry(_ context.Context, id string) (bool, error)
- func (s *MemoryStore) Size(_ context.Context, queue string) (int64, error)
- type Message
- type Option
- type Queue
- func (q *Queue) Config() Config
- func (q *Queue) Dispatch(ctx context.Context, job Job, opts ...DispatchOption) error
- func (q *Queue) DispatchFunc(ctx context.Context, name string, payload any, opts ...DispatchOption) error
- func (q *Queue) Fake()
- func (q *Queue) NameOf(job Job) (string, error)
- func (q *Queue) Observe(fn func(ctx context.Context, d Dispatched))
- func (q *Queue) Run(ctx context.Context, opts ...WorkOption) error
- func (q *Queue) Store() Store
- func (q *Queue) Work(opts ...WorkOption) error
- type Reservation
- type Store
- type WorkOption
Constants ¶
This section is empty.
Variables ¶
var ErrLeaseLost = errors.New("queue: the job's reservation has ended")
ErrLeaseLost is returned by a Store when a reservation has ended: its lease ran out (and another worker may have reserved the job), or the job was removed.
var ErrNoQueue = errors.New("queue: no queue in the context: call queue.ForApp at startup, or queue.WithQueue")
ErrNoQueue is returned by From (and Dispatch) when the context has no queue.
Functions ¶
func CountFailed ¶ added in v0.3.0
CountFailed returns the number of failed jobs of st: at once if it is a FailedCounter, else by reading them.
func CreateTables ¶
CreateTables creates the tables of a DatabaseStore: table for the jobs and failedTable for the failed ones. Times are Unix milliseconds.
func Dispatch ¶
func Dispatch(ctx context.Context, job Job, opts ...DispatchOption) error
Dispatch dispatches job on the queue in ctx (from ForApp):
err := queue.Dispatch(ctx, jobs.SendWelcome{UserID: u.ID}, queue.OnQueue("emails"))
With the database driver, a dispatch inside a transaction is part of it. With the sync driver, the job runs at once and its error is returned.
func DispatchFunc ¶
DispatchFunc dispatches a job of the function registered as name (see RegisterFunc) with payload, on the queue in ctx. payload must have the function's payload type (or point to a value of it).
func IsPermanent ¶
IsPermanent reports whether err, or an error it wraps, has a Permanent() bool method returning true: errors from Permanent, and from pubsub.Permanent, so code shared by jobs and listeners can use either.
func Migrations ¶
Migrations returns the migration creating the database driver's tables (default "jobs" and "failed_jobs"; pass the QUEUE_TABLE and QUEUE_FAILED_TABLE values if you set them) with CreateTables. Pass it to migrate.ForApp with the app's own:
migrate.ForApp(app, []*migrate.Set{migrations.All, queue.Migrations("", "")})
func Permanent ¶
Permanent wraps err so that the job fails for good, without retries: for errors that won't go away, such as a record that doesn't exist.
func Register ¶
Register adds the job type J, so the queue can dispatch and run it. Register every job type at startup, in the app that dispatches it and in the workers' (usually the same program):
err := queue.Register[jobs.SendWelcome](q, queue.Tries(5))
J is a struct type (or a pointer to one) implementing Job. A worker fails jobs whose type it doesn't know, so deploy workers before the code that dispatches new job types.
func RegisterFunc ¶
func RegisterFunc[T any](q *Queue, name string, fn func(ctx context.Context, payload T) error, opts ...JobOption) error
RegisterFunc registers fn as a job type named name, whose jobs carry a payload of type T, stored as JSON. Dispatch them with DispatchFunc:
err := queue.RegisterFunc(q, "reports.send", func(ctx context.Context, r ReportRequest) error { … })
err = queue.DispatchFunc(ctx, "reports.send", ReportRequest{Month: "2026-09"})
It suits jobs that need dependencies a struct's fields can't carry (fn can be a closure), and packages building on the queue (events.OnQueued). The name works like Name's, and the other options like Register's.
Types ¶
type Config ¶
type Config struct {
// Driver is where jobs are kept: sync (run at once), memory, database,
// or one passed to ForApp (redis). QUEUE_DRIVER, default sync.
Driver string `env:"QUEUE_DRIVER" default:"sync"`
// Default is the queue jobs go on without [OnQueue], and the one
// workers take jobs from without [Queues]. QUEUE_DEFAULT, default
// "default".
Default string `env:"QUEUE_DEFAULT" default:"default"`
// Tries is how many times a job runs before it fails for good, unless
// its type says otherwise ([Tries]). QUEUE_TRIES, default 3.
Tries int `env:"QUEUE_TRIES" default:"3"`
// Timeout is how long an attempt may run ([Timeout]). QUEUE_TIMEOUT,
// default 1m.
Timeout time.Duration `env:"QUEUE_TIMEOUT" default:"1m"`
// Backoff is the wait before the first retry, doubling for each
// later one ([Backoff]). QUEUE_BACKOFF, default 10s.
Backoff time.Duration `env:"QUEUE_BACKOFF" default:"10s"`
// MaxBackoff caps the doubling. QUEUE_BACKOFF_MAX, default 10m.
MaxBackoff time.Duration `env:"QUEUE_BACKOFF_MAX" default:"10m"`
// Poll is how long an idle worker waits before looking for jobs
// again. QUEUE_POLL, default 1s.
Poll time.Duration `env:"QUEUE_POLL" default:"1s"`
// Table is the database driver's table of jobs. QUEUE_TABLE, default
// jobs.
Table string `env:"QUEUE_TABLE" default:"jobs"`
// FailedTable is the database driver's table of failed jobs.
// QUEUE_FAILED_TABLE, default failed_jobs.
FailedTable string `env:"QUEUE_FAILED_TABLE" default:"failed_jobs"`
// Prefix starts the keys of the Redis driver, so apps (and the cache
// and sessions) can share a server. QUEUE_PREFIX, default APP_NAME
// followed by ":queue:" ("blog:queue:").
Prefix string `env:"QUEUE_PREFIX"`
}
Config selects and configures the app's queue.
type DatabaseStore ¶
type DatabaseStore struct {
// contains filtered or unexported fields
}
DatabaseStore keeps jobs in a table of the app's database (and failed jobs in another), so every instance of the app shares them without another server. Create the tables with Migrations. Delays and leases follow the database server's clock.
Store.Push runs in the context's transaction, if there is one: the job is dispatched if and only if the transaction commits, and workers see it only then. The workers' queries use their own connections (except with SQLite, which has one writer at a time). Workers reserve jobs with SELECT … FOR UPDATE SKIP LOCKED, which needs PostgreSQL, MySQL 8.0 or MariaDB 10.6 or later.
func NewDatabaseStore ¶
func NewDatabaseStore(d *db.DB, table, failedTable string) *DatabaseStore
NewDatabaseStore returns a store in the tables table (default "jobs") and failedTable (default "failed_jobs") of d.
func (*DatabaseStore) Close ¶
func (s *DatabaseStore) Close() error
Close implements Store; the database belongs to the app, so it does nothing.
func (*DatabaseStore) CountFailed ¶ added in v0.3.0
func (s *DatabaseStore) CountFailed(ctx context.Context) (int64, error)
CountFailed implements FailedCounter.
func (*DatabaseStore) Delete ¶
func (s *DatabaseStore) Delete(ctx context.Context, r *Reservation) error
Delete implements Store.
func (*DatabaseStore) Fail ¶
func (s *DatabaseStore) Fail(ctx context.Context, r *Reservation, errMsg string) error
Fail implements Store.
func (*DatabaseStore) FindFailed ¶ added in v0.3.0
FindFailed implements FailedFinder.
func (*DatabaseStore) Flush ¶
func (s *DatabaseStore) Flush(ctx context.Context) (int64, error)
Flush implements Store.
func (*DatabaseStore) Release ¶
func (s *DatabaseStore) Release(ctx context.Context, r *Reservation, delay time.Duration, refund bool) error
Release implements Store.
func (*DatabaseStore) Reserve ¶
func (s *DatabaseStore) Reserve(ctx context.Context, queue string, lease time.Duration) (*Reservation, error)
Reserve implements Store.
type DispatchOption ¶
type DispatchOption func(*dispatchOptions)
DispatchOption configures one dispatch.
func AfterCommit ¶
func AfterCommit() DispatchOption
AfterCommit dispatches the job only once the database transaction in the context commits, and not at all if it rolls back, so the job never looks for data that isn't there (yet). Since Dispatch has returned by then, a failure to dispatch is logged, not returned. Without a transaction, the job is dispatched at once, and Dispatch returns the error. The database driver writes the job in the transaction instead, when it is on the queue's database: then the job is dispatched if and only if the transaction commits.
func Delay ¶
func Delay(d time.Duration) DispatchOption
Delay makes the job available only after d. The sync driver ignores it.
func OnDispatched ¶
func OnDispatched(fn func(ctx context.Context, d Dispatched)) DispatchOption
OnDispatched calls fn once the job is dispatched, when Queue.Observe functions are: once the store has it (after the commit with AfterCommit); not if the store refuses it or the transaction rolls back. With the sync driver, fn is called before the job runs, whatever its outcome. Each OnDispatched option adds a function. mailer.Queue uses it to record queued email.
func OnQueue ¶
func OnQueue(name string) DispatchOption
OnQueue puts the job on the queue name: workers take jobs from the queues they were given, in order. Default QUEUE_DEFAULT ("default").
type Dispatched ¶
type Dispatched struct {
// ID is the job's ID.
ID string
// Job is the job type's name ("SendWelcome", "mail:send").
Job string
// Queue is the queue it went on.
Queue string
// Delay is its delay ([Delay]).
Delay time.Duration
// Data is the job, or a function job's payload, as JSON.
Data json.RawMessage
}
Dispatched describes a dispatched job, for Queue.Observe.
func (Dispatched) Decode ¶
func (d Dispatched) Decode(v any) error
Decode decodes the job's data into v: a pointer to the job type, or to a function job's payload type.
type Driver ¶
type Driver struct {
// Name is the value of QUEUE_DRIVER that selects the driver.
Name string
// Open returns the store for the app. It may add providers to the
// app, for example to check a server when the app boots. A nil store
// runs jobs at once.
Open func(app *anetos.App, cfg Config) (Store, error)
}
Driver opens a store for ForApp. The sync, memory and database drivers are built in; driver modules provide others (redis.QueueDriver()).
func DatabaseDriver ¶
func DatabaseDriver() Driver
DatabaseDriver is the database store's driver (QUEUE_DRIVER=database), in the tables QUEUE_TABLE and QUEUE_FAILED_TABLE. It uses the app's database: call db.Connect before queue.ForApp.
func MemoryDriver ¶
func MemoryDriver() Driver
MemoryDriver keeps jobs in the process's memory (QUEUE_DRIVER=memory): workers in the same process run them, and they are lost when it stops. For development and tests.
func SyncDriver ¶
func SyncDriver() Driver
SyncDriver runs each job at once, in Dispatch, which returns its error (QUEUE_DRIVER=sync): for development and tests. Jobs run once, without a delay, and workers do nothing.
type FailedCounter ¶ added in v0.3.0
type FailedCounter interface {
// CountFailed returns the number of failed jobs.
CountFailed(ctx context.Context) (int64, error)
}
FailedCounter is a Store that counts its failed jobs at once; the memory, database and Redis stores are.
type FailedFinder ¶ added in v0.3.0
type FailedFinder interface {
// FindFailed returns the failed job id, and whether there is one.
FindFailed(ctx context.Context, id string) (FailedJob, bool, error)
}
FailedFinder is a Store that finds a failed job by ID at once; the memory, database and Redis stores are.
type FailedJob ¶
type FailedJob struct {
// ID is the job's ID.
ID string
// Queue is the queue it was on.
Queue string
// Payload is the job's encoding.
Payload []byte
// Error is the error of its last attempt.
Error string
// Attempts is the number of times it ran.
Attempts int
// FailedAt is when it failed, by the store's clock.
FailedAt time.Time
}
FailedJob is a job that failed for good: its last attempt failed, it used all its tries, or it returned a Permanent error.
func FindFailed ¶ added in v0.3.0
FindFailed returns the failed job id of st, and whether there is one: at once if st is a FailedFinder, else by reading them.
type Failer ¶
type Failer interface {
// Failed handles the job's failure, say by telling someone.
Failed(ctx context.Context, err error)
}
Failer is implemented by jobs that need to know when they fail for good: Failed runs once, after the last attempt, with its error. If the last attempt didn't record its outcome (its worker died, or it ran past its lease), Failed runs instead of another attempt, although that attempt's work may have been done.
type Info ¶
type Info struct {
// ID identifies the job across its attempts: use it to make the job
// idempotent.
ID string
// Job is the job type's name.
Job string
// Queue is the queue it came from.
Queue string
// Attempt is this attempt's number, from 1.
Attempt int
// Tries is the number of attempts the job may have.
Tries int
}
Info describes the job running with a context.
type Job ¶
type Job interface {
// Handle does the work. An error retries the job, after a backoff,
// until it has used its tries; a [Permanent] error fails it at once.
Handle(ctx context.Context) error
}
Job is work to do in the background. Its exported fields are encoded as JSON when it is dispatched, so keep them small: IDs, not records.
type SendWelcome struct {
UserID int64 `json:"user_id"`
}
func (j SendWelcome) Handle(ctx context.Context) error { … }
Delivery is at-least-once: a job can run more than once (if a worker dies, or its lease runs out), so Handle must be idempotent. ctx carries what the app provides (the database, the cache, the queue), the job's Info, and its timeout.
type JobOption ¶
type JobOption func(*jobType) error
JobOption configures a job type for Register.
func Backoff ¶
Backoff sets the waits before retries: the first before the second attempt, and so on; the last is used for the rest. Each varies by up to 20%, so failed jobs don't all come back at once. Default QUEUE_BACKOFF (10s), doubling each time up to QUEUE_BACKOFF_MAX (10m).
func Name ¶
Name sets the name a job type is stored under. Default its Go type, such as "jobs.SendWelcome". Set it to keep jobs already queued working when the type is renamed or moved.
type MemoryStore ¶
type MemoryStore struct {
// contains filtered or unexported fields
}
MemoryStore keeps jobs in the process's memory: only workers in the same process see them, and they are lost when it stops. For development and tests.
func (*MemoryStore) Close ¶
func (s *MemoryStore) Close() error
Close implements Store; it does nothing.
func (*MemoryStore) CountFailed ¶ added in v0.3.0
func (s *MemoryStore) CountFailed(context.Context) (int64, error)
CountFailed implements FailedCounter.
func (*MemoryStore) Delete ¶
func (s *MemoryStore) Delete(_ context.Context, r *Reservation) error
Delete implements Store.
func (*MemoryStore) Fail ¶
func (s *MemoryStore) Fail(_ context.Context, r *Reservation, errMsg string) error
Fail implements Store.
func (*MemoryStore) FindFailed ¶ added in v0.3.0
FindFailed implements FailedFinder.
func (*MemoryStore) Flush ¶
func (s *MemoryStore) Flush(context.Context) (int64, error)
Flush implements Store.
func (*MemoryStore) Release ¶
func (s *MemoryStore) Release(_ context.Context, r *Reservation, delay time.Duration, refund bool) error
Release implements Store.
func (*MemoryStore) Reserve ¶
func (s *MemoryStore) Reserve(_ context.Context, queue string, lease time.Duration) (*Reservation, error)
Reserve implements Store.
type Message ¶
type Message struct {
// ID identifies the job, also once it has failed: a UUID.
ID string
// Queue is the name of the queue.
Queue string
// Payload is the job's encoding.
Payload []byte
}
Message is a job to push.
type Option ¶
type Option func(*Queue)
Option configures a Queue made with New.
func WithLogger ¶
WithLogger sets the logger of the queue's workers. Default slog.Default().
type Queue ¶
type Queue struct {
// contains filtered or unexported fields
}
Queue dispatches jobs to a Store and runs them in workers. Create it with ForApp (or New), and register each job type with Register.
func ForApp ¶
ForApp sets up the app's queue from the QUEUE_* settings: it opens the store with the driver QUEUE_DRIVER names (sync, memory and database are built in; pass others, such as redis.QueueDriver()), makes the queue available in every context the app creates (for Dispatch) and to anetos.Resolve, closes the store at shutdown, and adds the queue:failed, queue:retry, queue:forget, queue:flush and queue:clear commands.
q, err := queue.ForApp(app, redis.QueueDriver())
The database driver needs db.Connect first, and the tables from Migrations. Then register the job types (Register) and start the workers (Queue.Work).
func New ¶
New returns a queue keeping its jobs in store. Zero fields of cfg take their defaults (those of the QUEUE_* settings). A nil store runs jobs at once, like the sync driver. Run its workers with Queue.Run; for a queue made with ForApp, Queue.Work adds them to the app.
func (*Queue) DispatchFunc ¶
func (q *Queue) DispatchFunc(ctx context.Context, name string, payload any, opts ...DispatchOption) error
DispatchFunc dispatches a job of a function; see the function DispatchFunc.
func (*Queue) Fake ¶
func (q *Queue) Fake()
Fake makes the queue, from now on, pass dispatched jobs to the Queue.Observe functions only: they are neither stored nor run. For tests (anetostest.FakeQueue); it can't be undone.
func (*Queue) Observe ¶
func (q *Queue) Observe(fn func(ctx context.Context, d Dispatched))
Observe calls fn with each job dispatched from now on: once the store has it (with the database driver, inside a transaction, once that commits), or, with the sync driver, before it runs; with AfterCommit, after the commit. A job the store refuses isn't passed. In a transaction of db.WithTx, whose commit db doesn't see, the database driver's jobs aren't passed either. For tests and instrumentation. fn must be quick and safe for concurrent use. anetostest uses it to record jobs.
func (*Queue) Run ¶
func (q *Queue) Run(ctx context.Context, opts ...WorkOption) error
Run runs workers until ctx is canceled, then lets running jobs finish (see ShutdownGrace) and returns. Jobs get ctx's values, such as what an app adds to its contexts. Queue.Work runs it in the app.
func (*Queue) Work ¶
func (q *Queue) Work(opts ...WorkOption) error
Work adds workers to the app, as a component with the role "workers" (so `run --only=workers` runs only them) that stops after the HTTP server and listeners, so it finishes the jobs they dispatched:
err := q.Work(queue.Queues("emails", "default"), queue.Concurrency(10))
With the sync driver it does nothing. It needs a queue made with ForApp.
type Reservation ¶
type Reservation struct {
// ID is the job's ID.
ID string
// Queue is the name of its queue.
Queue string
// Payload is the job's encoding.
Payload []byte
// Attempts counts the reservations of the job, this one included.
Attempts int
// Token identifies this reservation to the store.
Token string
}
Reservation is a job a worker has leased.
type Store ¶
type Store interface {
// Push adds a job to m.Queue, available after delay (0 for now).
Push(ctx context.Context, m Message, delay time.Duration) error
// Reserve leases the next available job of queue for lease, adding
// one to its attempts. It returns nil, nil when none is available.
// Jobs become available in order of the time they became available,
// as far as the store's clock allows.
Reserve(ctx context.Context, queue string, lease time.Duration) (*Reservation, error)
// Delete removes a reserved job, once it is done. It returns
// [ErrLeaseLost] if the job isn't reserved with r's token anymore.
Delete(ctx context.Context, r *Reservation) error
// Release ends a reservation: the job becomes available again after
// delay. With refund, the attempt doesn't count (the worker stopped
// before the job could finish). It returns [ErrLeaseLost] if the job
// isn't reserved with r's token anymore.
Release(ctx context.Context, r *Reservation, delay time.Duration, refund bool) error
// Fail removes a reserved job and keeps it as failed, with the error
// message (replacing a failed job with the same ID). It returns
// [ErrLeaseLost], and keeps nothing, if the job isn't reserved with
// r's token anymore.
Fail(ctx context.Context, r *Reservation, errMsg string) error
// Size returns the number of jobs in queue, reserved ones included.
Size(ctx context.Context, queue string) (int64, error)
// Clear removes every job of queue, reserved ones included, and
// returns how many there were.
Clear(ctx context.Context, queue string) (int64, error)
// Failed returns up to limit failed jobs, the latest first, after
// skipping offset of them.
Failed(ctx context.Context, offset, limit int) ([]FailedJob, error)
// Retry pushes the failed job id back to its queue, with no attempts,
// and stops keeping it as failed. It reports whether there was one.
Retry(ctx context.Context, id string) (bool, error)
// Forget removes the failed job id, and reports whether there was one.
Forget(ctx context.Context, id string) (bool, error)
// Flush removes every failed job, and returns how many there were.
Flush(ctx context.Context) (int64, error)
// Close releases the store's resources.
Close() error
}
Store keeps jobs for a Queue: the memory and database stores are built in, and driver modules provide others (redis.QueueDriver()).
A reserved job is leased: it is hidden from other workers until the lease ends, then becomes available again (at-least-once delivery if a worker dies). Each reservation has its own token, and Delete, Release and Fail act only if the job is still reserved with it, so a worker whose lease ran out can't remove the job from under the worker that reserved it next.
Every method must be safe for concurrent use. The queue/queuetest package is a conformance suite for stores.
type WorkOption ¶
type WorkOption func(*workOptions) error
WorkOption configures workers.
func Concurrency ¶
func Concurrency(n int) WorkOption
Concurrency sets how many jobs the workers run at once. Default 1.
func Queues ¶
func Queues(names ...string) WorkOption
Queues sets the queues the workers take jobs from, in order of priority: a job on the first is taken before any on the second, and so on. Default QUEUE_DEFAULT.
func ShutdownGrace ¶
func ShutdownGrace(d time.Duration) WorkOption
ShutdownGrace sets how long running jobs may take to finish once the workers stop: then their contexts are canceled, and the jobs that return because of it go back on the queue, without the attempt counting (unless they return a Permanent error). Default: half of APP_SHUTDOWN_TIMEOUT with Queue.Work, and never more than the app's shutdown budget leaves (minus 2s to put the jobs back; with a budget of a few seconds, running jobs are stopped at once); 15s with Queue.Run.