queue

package
v0.4.0 Latest Latest
Warning

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

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 26 Imported by: 0

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

Constants

This section is empty.

Variables

View Source
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.

View Source
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

func CountFailed(ctx context.Context, st Store) (int64, error)

CountFailed returns the number of failed jobs of st: at once if it is a FailedCounter, else by reading them.

func CreateTables

func CreateTables(s *migrate.Schema, table, failedTable string) error

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

func DispatchFunc(ctx context.Context, name string, payload any, opts ...DispatchOption) error

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

func IsPermanent(err error) bool

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

func Migrations(table, failedTable string) *migrate.Set

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

func Permanent(err error) error

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

func Register[J Job](q *Queue, opts ...JobOption) error

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.

func WithQueue

func WithQueue(ctx context.Context, q *Queue) context.Context

WithQueue returns ctx with q, for Dispatch. ForApp makes the queue available in every context the app creates.

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.

func LoadConfig

func LoadConfig(src config.Source) (Config, error)

LoadConfig reads the QUEUE_* settings.

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) Clear

func (s *DatabaseStore) Clear(ctx context.Context, queue string) (int64, error)

Clear implements Store.

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) Failed

func (s *DatabaseStore) Failed(ctx context.Context, offset, limit int) ([]FailedJob, error)

Failed implements Store.

func (*DatabaseStore) FindFailed added in v0.3.0

func (s *DatabaseStore) FindFailed(ctx context.Context, id string) (FailedJob, bool, error)

FindFailed implements FailedFinder.

func (*DatabaseStore) Flush

func (s *DatabaseStore) Flush(ctx context.Context) (int64, error)

Flush implements Store.

func (*DatabaseStore) Forget

func (s *DatabaseStore) Forget(ctx context.Context, id string) (bool, error)

Forget implements Store.

func (*DatabaseStore) Push

func (s *DatabaseStore) Push(ctx context.Context, m Message, delay time.Duration) error

Push implements Store. It joins the context's transaction.

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.

func (*DatabaseStore) Retry

func (s *DatabaseStore) Retry(ctx context.Context, id string) (bool, error)

Retry implements Store.

func (*DatabaseStore) Size

func (s *DatabaseStore) Size(ctx context.Context, queue string) (int64, error)

Size 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

func FindFailed(ctx context.Context, st Store, id string) (FailedJob, bool, error)

FindFailed returns the failed job id of st, and whether there is one: at once if st is a FailedFinder, else by reading them.

func (FailedJob) Job

func (f FailedJob) Job() string

Job returns the name the job was registered with, from its payload.

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.

func Current

func Current(ctx context.Context) (Info, bool)

Current returns the job running with ctx, in its Handle and Failed methods.

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

func Backoff(delays ...time.Duration) JobOption

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

func Name(name string) JobOption

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.

func Timeout

func Timeout(d time.Duration) JobOption

Timeout sets how long one attempt may run: then its context is canceled, and the attempt fails. Default QUEUE_TIMEOUT (1m).

func Tries

func Tries(n int) JobOption

Tries sets how many times a job runs before it fails for good: 1 for no retries. Default QUEUE_TRIES (3).

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 NewMemoryStore

func NewMemoryStore() *MemoryStore

NewMemoryStore returns an empty store.

func (*MemoryStore) Clear

func (s *MemoryStore) Clear(_ context.Context, queue string) (int64, error)

Clear implements Store.

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) Failed

func (s *MemoryStore) Failed(_ context.Context, offset, limit int) ([]FailedJob, error)

Failed implements Store.

func (*MemoryStore) FindFailed added in v0.3.0

func (s *MemoryStore) FindFailed(_ context.Context, id string) (FailedJob, bool, error)

FindFailed implements FailedFinder.

func (*MemoryStore) Flush

func (s *MemoryStore) Flush(context.Context) (int64, error)

Flush implements Store.

func (*MemoryStore) Forget

func (s *MemoryStore) Forget(_ context.Context, id string) (bool, error)

Forget implements Store.

func (*MemoryStore) Push

func (s *MemoryStore) Push(_ context.Context, m Message, delay time.Duration) error

Push 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.

func (*MemoryStore) Retry

func (s *MemoryStore) Retry(_ context.Context, id string) (bool, error)

Retry implements Store.

func (*MemoryStore) Size

func (s *MemoryStore) Size(_ context.Context, queue string) (int64, error)

Size 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

func WithLogger(l *slog.Logger) Option

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

func ForApp(app *anetos.App, drivers ...Driver) (*Queue, error)

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 From

func From(ctx context.Context) (*Queue, error)

From returns the queue in ctx.

func New

func New(store Store, cfg Config, opts ...Option) *Queue

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) Config

func (q *Queue) Config() Config

Config returns the queue's settings, with defaults filled in.

func (*Queue) Dispatch

func (q *Queue) Dispatch(ctx context.Context, job Job, opts ...DispatchOption) error

Dispatch dispatches job; see the function Dispatch.

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) NameOf

func (q *Queue) NameOf(job Job) (string, error)

NameOf returns the name job's type is registered with (Register).

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) Store

func (q *Queue) Store() Store

Store returns the queue's store; nil for the sync driver.

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.

Directories

Path Synopsis
Package queuetest is a conformance suite for queue stores.
Package queuetest is a conformance suite for queue stores.

Jump to

Keyboard shortcuts

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