queue

package
v0.13.0 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 26 Imported by: 0

Documentation

Overview

Package queue is work that happens after the response: the Queue contract, the drivers that ship with the collection, and the Worker that drains them.

The surface is Push, PushOn, Later, Bulk, Pop, Size and Clear, and the job itself is github.com/arandu-io/hesape/queue/jobs.

The contract lives in the collection and so do the drivers that need nothing installed, for the same reason as database.Repository: Push takes an auth.Grant, and the tenant comes from it. Moving that into an optional package would make the guarantee optional, and an optional guarantee is not one.

DatabaseQueue    the application's own database (the default)
SyncQueue        runs the job at Push, for tests and for a laptop
DeferredQueue    runs the job after the response, in this process
BackgroundQueue  runs the job in another process
FailoverQueue    writes to the first connection that accepts
NullQueue        accepts everything and keeps nothing
RedisQueue       github.com/arandu-io/hesape/queue/connectors/redis

RedisQueue is a separate module because in Go there is no optional dependency and a collection that carried a Redis client would put it in every project's go.sum.

The outbox is the mechanism, not the name

A job pushed inside database.Transaction is committed by the same transaction as the row it is about: it exists if and only if the write did. That is the outbox guarantee -- the one the events package uses for events -- applied to work, and it is the reason DatabaseQueue is the default driver rather than the fallback one. The relay that drains it is the Worker.

The name does not change because of it. Calling it Outbox would name the mechanism instead of the thing, and hide it from everyone who came looking for the driver that keeps jobs in the application's database.

At-least-once

A handler that cannot run twice safely is a handler with a bug. The process can die between doing the work and deleting the job, and no queue anywhere solves that.

Where the rest of it lives

queue/attributes  the per-job settings: tries, backoff, timeout
queue/connectors  opening a connection lazily
queue/console     queue:work, queue:retry, queue:pause and the rest
queue/events      what the queue announces about itself
queue/failed      a dead letter list that outlives the queue
queue/jobs        the job itself
queue/middleware  what wraps the handling of one job

Index

Constants

View Source
const (
	// Interrupted is a signal: SIGTERM from an orchestrator, SIGINT from a
	// person. The worker finished the job it had and then stopped, which is
	// what makes a deploy not lose work.
	Interrupted = events.Interrupted

	// MaxJobsExceeded is WorkerOptions.MaxJobs reached. A worker that stops
	// after n jobs is how a leak in a handler stays survivable.
	MaxJobsExceeded = events.MaxJobsExceeded

	// MaxMemoryExceeded is WorkerOptions.Memory reached.
	MaxMemoryExceeded = events.MaxMemoryExceeded

	// MaxTimeExceeded is WorkerOptions.MaxTime reached.
	MaxTimeExceeded = events.MaxTimeExceeded

	// QueueEmpty is WorkerOptions.StopWhenEmpty with nothing left to do. It is
	// what a batch job in a pipeline waits for.
	QueueEmpty = events.QueueEmpty

	// ReceivedRestartSignal is `aru queue:restart`: a timestamp went into the
	// cache and every worker that reads it stops, so the next deploy's binary
	// is the one running.
	ReceivedRestartSignal = events.ReceivedRestartSignal
)
View Source
const (
	// ExitSuccess is a worker that stopped because it was asked to.
	ExitSuccess = 0
	// ExitError is a worker that stopped because it could not continue.
	ExitError = 1
	// ExitMemoryLimit is a worker that stopped because it was using too much
	// memory. A supervisor that restarts on this one is doing the right thing.
	ExitMemoryLimit = 12
)

Exit statuses a worker returns, and what a supervisor reads.

Twelve for memory is not one of the shell's own statuses, so a supervisor can tell "the worker decided to stop" from "something killed it".

Variables

View Source
var ErrHandlerPanicked = errors.New("queue: the job's handler panicked")

ErrHandlerPanicked is what HandlerPanicked matches under errors.Is, for a caller that wants to tell a panic from a returned error without reaching for the type.

A job that matches it was parked on its first delivery rather than retried. A panic is a defect in the handler and it is deterministic: delivering it again reproduces it, and the worker that takes it is the next one to lose the process -- which is the loop this exists to cut.

View Source
var ErrInvalidPayload = errors.New("queue: the job payload cannot be encoded")

ErrInvalidPayload is returned when a job's arguments cannot be encoded.

The value that failed to encode is not carried on it: the wrapped json error already names the type and the field.

View Source
var ErrManuallyFailed = errors.New("queue: the job failed and will not be retried")

ErrManuallyFailed is what a handler returns to park its own job.

A handler that knows the work can never succeed -- the customer is gone, the file is malformed -- wraps it, and the worker parks the job on the first delivery instead of retrying four more times on the way to the same place.

return fmt.Errorf("%w: the invoice was voided", queue.ErrManuallyFailed)
View Source
var ErrMaxAttemptsExceeded = errors.New("queue: the job has been attempted too many times")

ErrMaxAttemptsExceeded is what MaxAttemptsExceeded and TimeoutExceeded both match under errors.Is, for a caller that only wants to know the job gave up rather than why.

View Source
var ErrMissingModel = fmt.Errorf("queue: the record this job is about no longer exists")

ErrMissingModel is returned when a record a job's payload names is gone.

A job whose record was deleted cannot succeed on any retry. A handler that wants that treated as success rather than as failure says so with attributes.Attributes.DeleteWhenMissingModels.

View Source
var ErrNoBackgroundRunner = errNoBackgroundRunner{}

ErrNoBackgroundRunner is returned when a background queue was built without a way to start a process.

View Source
var ErrNoCache = fmt.Errorf("queue: no cache is wired, so there is nowhere to record this. Call SetCache in bootstrap/app.go")

ErrNoCache is returned when the manager is asked to do something that needs somewhere to write and has nowhere.

It is an error rather than a silent no-op, because "the queue is paused" is a thing an operator believes after the command returns, and believing it wrongly is how work keeps going out during an incident.

View Source
var ErrNoConnection = fmt.Errorf("queue: no such connection")

ErrNoConnection is returned when a connection nobody registered is asked for.

View Source
var ErrNoFailoverConnections = errors.New("queue: this failover queue has no connections. Name at least two in NewFailoverQueue")

ErrNoFailoverConnections is returned when a failover queue was built with none.

View Source
var ErrNoListenerCommand = errors.New("queue: this listener has no command to run. Pass one to NewListener")

ErrNoListenerCommand is returned when a listener was built with no command to run.

Functions

func CreatePayload

func CreatePayload(g auth.Grant, connection, queue, name string, data any, delay time.Duration) (jobs.Job, error)

CreatePayload builds the record a driver stores.

It takes what the caller pushed and returns the thing that goes on the wire, with the identity, the tenant and the settings filled in. The return is a jobs.Job rather than an encoded envelope, because the record has columns and the arguments are one of them.

delay is folded into RunAt rather than kept beside it: a delay and an availability instant carried side by side can disagree, and one field cannot.

func CreatePayloadUsing

func CreatePayloadUsing(hook PayloadHook)

CreatePayloadUsing registers a hook that runs while a job's record is built.

Passing nil clears the registered hooks rather than adding a nil one: a test that registers a hook has no other way to take it back.

It is process-wide. A hook is registered once at boot -- tracing, an audit field -- and never per request; anything that varies per job belongs on the job.

func ExponentialBackoff

func ExponentialBackoff(attempt int) time.Duration

ExponentialBackoff doubles the wait each attempt, capped at an hour.

Capped, because unbounded doubling means the eleventh attempt is next year -- and a job nobody will ever see fail is worse than one that parks.

func GetDisplayName

func GetDisplayName(data any) string

GetDisplayName is what to show for a job's arguments, or empty when the value has nothing to say.

func GetJobBackoff

func GetJobBackoff(data any) []time.Duration

GetJobBackoff is how long a job value asks to wait between deliveries, or nil when the worker's own schedule decides.

It is the slice attributes.Attributes.Backoff holds, not an encoding of it.

func GetJobExpiration

func GetJobExpiration(data any) time.Time

GetJobExpiration is the deadline after which a job value stops being retried, or the zero time when there is none.

It is a time.Time rather than a Unix timestamp: the same instant with its unit attached.

func GetJobTries

func GetJobTries(data any) int

GetJobTries is how many deliveries a job value asks for, or zero when it asks for nothing and the worker's own limit decides.

It reads attributes.Attributes.Tries, which is the one place a value says it.

func LaterOn

func LaterOn(ctx context.Context, q Queue, g auth.Grant, name string, delay time.Duration, j jobs.Job) error

LaterOn adds a job to a named queue, eligible after delay.

It is a package function rather than a method on the contract because it cannot be anything but Later with the queue overwritten. PushOn is on the contract instead, because a driver can do that one in a single round trip.

queue.LaterOn(ctx, q, g, "reports", time.Hour, j)

func PopUsing

func PopUsing(workerName string, callback PopCallback)

PopUsing registers how the worker called workerName chooses its next batch.

Passing nil forgets the callback rather than registering an empty one.

func Stderr

func Stderr(line string)

Stderr is where a listener with no output handler sends the child's output. It is here so a caller can say `l.SetOutputHandler(queue.Stderr)`.

Types

type AddCreatedAtToJobsTable added in v0.5.0

type AddCreatedAtToJobsTable struct{ migrations.BaseMigration }

AddCreatedAtToJobsTable adds the column Pop takes jobs in the order of, which used to be run_at.

The queue is first in, first out over whatever is eligible. Ordering by run_at is not: a job pushed with a ten second delay waits out the delay and then goes behind everything queued while it waited, again on every pass, and on a queue that is never empty it is a job that never runs.

The id cannot do the job here because it is a random uuid and sorts by nothing (see database.NewID), so the column that carries the order is the one that says when the row was written.

func (AddCreatedAtToJobsTable) Down added in v0.5.0

Down drops the column.

func (AddCreatedAtToJobsTable) GetName added in v0.5.0

func (AddCreatedAtToJobsTable) GetName() string

GetName returns the migration's name.

func (AddCreatedAtToJobsTable) Up added in v0.5.0

Up adds the column and backfills it.

Nullable, and with no default, because SQLite refuses a non-constant default on ADD COLUMN: the rows already in the table are backfilled here, and the ones the previous release's binary writes during the rollout arrive NULL, which is what the COALESCE in Pop is for.

No index of its own. The filter is still served by idx_jobs_ready, and what is left is a top-N sort under a LIMIT; an index on (queue, failed_at, created_at) would remove that sort, and it would need a DROP INDEX in Down, which MySQL spells "DROP INDEX x ON jobs" and the other two spell "DROP INDEX x". One portable migration is worth more than the sort.

type AddJobAttributesToJobsTable added in v0.5.0

type AddJobAttributesToJobsTable struct{ migrations.BaseMigration }

AddJobAttributesToJobsTable adds the per-job settings, which used to live nowhere: a job pushed with five tries got the worker's five whatever it asked for, because the number never reached the table.

They travel as one JSON column rather than one column each -- see attributes.Attributes for why they are one struct -- and the two names beside it are the two things about a job that are neither its arguments nor a setting.

func (AddJobAttributesToJobsTable) Down added in v0.5.0

Down drops the three columns.

func (AddJobAttributesToJobsTable) GetName added in v0.5.0

GetName returns the migration's name.

func (AddJobAttributesToJobsTable) Up added in v0.5.0

Up adds the three columns.

Nullable with a default, so the previous release's binary keeps inserting without them during a rollout.

type BackgroundQueue

type BackgroundQueue struct {
	// contains filtered or unexported fields
}

BackgroundQueue runs each job in a separate operating system process.

Work that must not block the request usually belongs on DeferredQueue -- which is this without the process. What is left for this driver is the case a goroutine cannot serve: a job that must survive the request's process crashing, or one that has to be isolated from it.

It is a wrapper over a command rather than a fork, because Go cannot fork: Run is `aru work --once` or whatever the application names it, and the job is handed over on its argument list. An application that has not set Run gets an error at push, not silence.

func NewBackgroundQueue

func NewBackgroundQueue(run func(ctx context.Context, j jobs.Job) error) *BackgroundQueue

NewBackgroundQueue returns the queue over a way to start a process.

run is given the job and must start a process for it and return; it must not wait for the work, because the point of this driver is that the caller does not.

func (*BackgroundQueue) Bulk

func (q *BackgroundQueue) Bulk(ctx context.Context, g auth.Grant, js []jobs.Job) error

Bulk starts each job in order, stopping at the first failure to start.

func (*BackgroundQueue) Clear

Clear removes nothing.

func (*BackgroundQueue) CreationTimeOfOldestPendingJob

func (q *BackgroundQueue) CreationTimeOfOldestPendingJob(context.Context, string) (time.Time, error)

CreationTimeOfOldestPendingJob is the zero time: nothing waits.

func (*BackgroundQueue) DelayedSize

func (q *BackgroundQueue) DelayedSize(context.Context, string) (int, error)

DelayedSize is zero. It answers delayedSize().

func (*BackgroundQueue) DeleteJob

func (q *BackgroundQueue) DeleteJob(context.Context, *jobs.Job) error

DeleteJob does nothing.

func (*BackgroundQueue) FailJob

FailJob does nothing.

func (*BackgroundQueue) Failed

func (q *BackgroundQueue) Failed(context.Context, int) ([]jobs.Job, error)

Failed lists nothing: a job that failed failed in another process, and its error went to that process's log.

func (*BackgroundQueue) GetConnectionName

func (c *BackgroundQueue) GetConnectionName() string

GetConnectionName is the name this queue was registered under.

It is empty on a queue built directly and never handed to a QueueManager, which is what a test does.

func (*BackgroundQueue) Later

func (q *BackgroundQueue) Later(ctx context.Context, g auth.Grant, delay time.Duration, j jobs.Job) error

Later starts the job with its RunAt set, and leaves the waiting to the process that runs it.

Sleeping here would hold the caller for the delay, which is the one thing this driver exists to avoid.

func (*BackgroundQueue) PendingSize

func (q *BackgroundQueue) PendingSize(context.Context, string) (int, error)

PendingSize is zero.

func (*BackgroundQueue) Pop

Pop returns nothing: the job is already in another process.

func (*BackgroundQueue) Push

func (q *BackgroundQueue) Push(ctx context.Context, g auth.Grant, j jobs.Job) error

Push starts the job in another process.

The job is authorized first, for the reason DeferredQueue.Push gives: a refusal has to reach the caller.

func (*BackgroundQueue) PushOn

func (q *BackgroundQueue) PushOn(ctx context.Context, g auth.Grant, queue string, j jobs.Job) error

PushOn starts the job on a named queue.

func (*BackgroundQueue) PushRaw

func (q *BackgroundQueue) PushRaw(ctx context.Context, g auth.Grant, name string, payload []byte, queue string) error

PushRaw starts a job whose arguments are already encoded. It answers pushRaw().

func (*BackgroundQueue) ReleaseJob

ReleaseJob does nothing. There is no queue to put the job back on.

func (*BackgroundQueue) ReservedSize

func (q *BackgroundQueue) ReservedSize(context.Context, string) (int, error)

ReservedSize is zero. It answers reservedSize().

func (*BackgroundQueue) Retry

Retry has nothing to retry.

func (*BackgroundQueue) SetConnectionName

func (c *BackgroundQueue) SetConnectionName(name string)

SetConnectionName names the connection this queue answers to.

It returns nothing rather than the queue, because an embedded struct cannot return the type that embeds it.

func (*BackgroundQueue) Size

Size is zero. Nothing is stored.

type Cache

type Cache interface {
	Get(ctx context.Context, key string) ([]byte, error)
	Put(ctx context.Context, key string, value []byte, ttl time.Duration) error
	Forget(ctx context.Context, key string) error
}

Cache is the little a worker asks of a cache store: whether somebody asked for a restart, and whether this queue is paused.

It is declared here rather than imported so the queue does not depend on a store to run. cache.Store satisfies it. There is no auth.Grant on it and that is deliberate: a restart signal and a pause are operational state about the process, not data about a customer, so there is no tenant to scope them by.

type CallQueuedClosure

type CallQueuedClosure struct {
	// contains filtered or unexported fields
}

CallQueuedClosure is a job whose work is a function rather than a registered name.

Go cannot serialize a function, so the closure stays in the binary and what travels is the name it was registered under.

c := queue.NewCallQueuedClosure(func(ctx context.Context, g auth.Grant, j *jobs.Job) error {
	return warmTheCache(ctx, g)
}).Name("cache.warm").OnFailure(func(err error) { alert(err) })

w.Handle(c.JobName(), c)
q.Push(ctx, g, c.Job(g))

The registration is the price of not having serializable closures, and it is the honest one: a closure that is not in the binary the worker is running cannot be called by it.

func NewCallQueuedClosure

func NewCallQueuedClosure(closure HandlerFunc) *CallQueuedClosure

NewCallQueuedClosure returns a job for a function.

func (*CallQueuedClosure) DisplayName

func (c *CallQueuedClosure) DisplayName() string

DisplayName is what a console and the failed job list show.

It is the name the caller gave, not the one the runtime has: runtime.FuncForPC would answer pkg.glob..func3, which tells a person less.

func (*CallQueuedClosure) Failed

func (c *CallQueuedClosure) Failed(_ context.Context, _ auth.Grant, _ *jobs.Job, cause error)

Failed calls every failure callback with the cause.

func (*CallQueuedClosure) Handle

func (c *CallQueuedClosure) Handle(ctx context.Context, g auth.Grant, j *jobs.Job) error

Handle runs the closure.

func (*CallQueuedClosure) Job

func (c *CallQueuedClosure) Job(g auth.Grant) (jobs.Job, error)

Job is the record to push for this closure.

The closure and the record are separate: the closure is registered with a worker once, and this is what goes on the queue each time.

func (*CallQueuedClosure) JobName

func (c *CallQueuedClosure) JobName() string

JobName is the name this closure is registered and pushed under.

A closure has no name of its own, so an unnamed one gets a name that says so rather than an empty string that would collide with every other unnamed closure.

func (*CallQueuedClosure) Name

Name assigns a name to the job, and returns it so the call chains.

func (*CallQueuedClosure) OnFailure

func (c *CallQueuedClosure) OnFailure(callback func(error)) *CallQueuedClosure

OnFailure adds a callback for when the job gives up, and returns the job so the call chains.

type CallQueuedHandler

type CallQueuedHandler struct {
	// contains filtered or unexported fields
}

CallQueuedHandler turns a job on the wire into a call to the code that runs it.

The name on the wire is a string and the registry is the Worker. What this does is run the handler through the job's middleware and settle the job afterwards, the same way whoever called it would have.

h := queue.NewCallQueuedHandler(w).Through(middleware.Skip{}.When(readOnly))
err := h.Call(ctx, j)

Worker.Process is the caller that matters, and it does this inline. What this type is for is everything else that has a job and needs it run -- SyncQueue, a test, a command that replays one job by hand -- so the settling is written once.

func NewCallQueuedHandler

func NewCallQueuedHandler(h Handlers) *CallQueuedHandler

NewCallQueuedHandler returns the handler over a registry.

func (*CallQueuedHandler) Call

func (h *CallQueuedHandler) Call(ctx context.Context, j *jobs.Job) error

Call runs the job and settles it.

The settling is the part worth reading, and the order is:

  • a job the middleware already released or deleted is left alone, because settling it twice is how a job runs twice
  • a job whose handler failed is left for the caller to release or park, because only the caller knows how many tries are left
  • anything else is deleted, which is what "it worked" means to a queue

func (*CallQueuedHandler) Failed

func (h *CallQueuedHandler) Failed(ctx context.Context, j *jobs.Job, cause error)

Failed runs whatever the job wants done when it gives up for good.

It is called after the job has been parked, not instead of parking it: the dead letter row is what an operator sees, and this is the compensating action -- refund the hold, mark the import abandoned, tell the customer.

func (*CallQueuedHandler) Through

Through wraps every call in middleware, outermost first, and returns the handler so the call chains.

A job here is a name and a payload, with nowhere on it to hang a middleware list, so the list comes off whoever is running it.

type Config added in v0.2.0

type Config struct {
	// Default is the connection a Dispatch goes to when it names none, and
	// empty means "database".
	Default string

	// Connections is every connection the application knows, by name.
	Connections map[string]ConnectionConfig

	// Failed is where a job goes when it has run out of tries.
	Failed FailedConfig
}

Config is the queue's settings: which connection the application dispatches to, and how each one behaves.

It is declared here, in the package that reads it, and not in a configuration package of its own. Nothing looks a value up by key, so the component is handed its settings and the compiler checks the field.

type ConnectionConfig added in v0.2.0

type ConnectionConfig struct {
	// Driver is "sync", "database" or "redis".
	//
	// There is no "beanstalkd", "sqs" or "null", and none of the three is
	// pending: the first two are stores this ecosystem does not speak to, and
	// a driver that silently discards work is a thing to reach for in a test,
	// which "sync" and a recorder already answer without pretending a job ran.
	Driver string

	// Name is the connection's own name, so that a value carries it without the
	// map key having to travel alongside.
	Name string

	// Queue is the queue a job lands on when it names none.
	Queue string

	// RetryAfter is how long a reserved job may be held before another worker
	// may take it.
	//
	// It has to be longer than the longest job actually takes. Shorter, and a
	// slow job is picked up a second time while the first run is still going --
	// which is at-least-once delivery arriving as a duplicate that nobody asked
	// for.
	RetryAfter time.Duration

	// Table is the table a database connection reads, and empty means "jobs".
	Table string

	// Connection is the database or redis connection this queue rides on, empty
	// meaning the application's default.
	Connection string

	// AfterCommit delays the dispatch until the surrounding transaction commits.
	//
	// The failure it guards against is this: a job dispatched inside a
	// transaction can be picked up by a worker before the transaction commits,
	// and it then reads a row that does not exist yet.
	//
	// It is false by default, because the outbox already has the better answer:
	// it writes the job in the same transaction as the change that produced it,
	// so there is no window at all. AfterCommit narrows the window; the outbox
	// closes it.
	AfterCommit bool
}

ConnectionConfig is one entry of Config.Connections: how one connection behaves.

type CreateJobsTable added in v0.5.0

type CreateJobsTable struct{ migrations.BaseMigration }

CreateJobsTable creates the table DatabaseQueue pushes to and pops from.

func (CreateJobsTable) Down added in v0.5.0

Down drops the jobs table, and both indexes with it.

func (CreateJobsTable) GetName added in v0.5.0

func (CreateJobsTable) GetName() string

GetName returns the migration's name.

func (CreateJobsTable) Up added in v0.5.0

Up creates the jobs table and the two indexes its queries are.

Portable types only: TEXT, INTEGER and TIMESTAMP mean the same thing on SQLite, Postgres and MySQL.

type DatabaseQueue

type DatabaseQueue struct {
	// contains filtered or unexported fields
}

DatabaseQueue is the queue backed by the application's own database.

It is the default driver because it needs nothing installed: the jobs table sits in the database the application already has.

What it offers that no other driver can is the outbox guarantee. A job pushed inside database.Transaction is committed by the same transaction as the row it is about, so it exists if and only if the write did -- the mechanism the events package uses for events, applied to work. The name is the driver and not the guarantee: naming it after the guarantee would hide it from everyone looking for the queue that lives in the application's database.

func NewDatabaseQueue

func NewDatabaseQueue(db *database.DB) *DatabaseQueue

NewDatabaseQueue returns the queue over db.

func (*DatabaseQueue) Bulk

func (q *DatabaseQueue) Bulk(ctx context.Context, g auth.Grant, js []jobs.Job) error

Bulk adds many jobs.

One statement each rather than a multi-row insert, because inside database.Transaction the difference is a round trip and outside it a multi-row insert would be the only place in the package where a partial failure leaves half a batch behind with no way to say which half.

func (*DatabaseQueue) Clear

func (q *DatabaseQueue) Clear(ctx context.Context, queue string) (int, error)

Clear removes every job waiting or in flight on a queue, and returns how many went.

Parked jobs are not cleared: a job that gave up is no longer on a queue, it is in the dead letter list, and DatabaseQueue.Failed and DatabaseQueue.Retry are how it is dealt with. The RESP driver draws the line in the same place.

func (*DatabaseQueue) CreationTimeOfOldestPendingJob

func (q *DatabaseQueue) CreationTimeOfOldestPendingJob(ctx context.Context, queue string) (time.Time, error)

CreationTimeOfOldestPendingJob is when the oldest waiting job became eligible, or the zero time when nothing is waiting.

func (*DatabaseQueue) DelayedSize

func (q *DatabaseQueue) DelayedSize(ctx context.Context, queue string) (int, error)

DelayedSize is how many jobs are waiting for a time that has not come.

It is the number DatabaseQueue.PendingSize leaves out: a job scheduled for tomorrow is not a backlog, and counting it as one is how a health check pages somebody at three in the morning about a report due at nine.

func (*DatabaseQueue) DeleteAndRelease

func (q *DatabaseQueue) DeleteAndRelease(ctx context.Context, queue string, j *jobs.Job, delay time.Duration) error

DeleteAndRelease removes a reserved job and queues it again.

A release here is one UPDATE -- the row never left -- so this is that update, under the name a driver that claims by deletion would need.

func (*DatabaseQueue) DeleteJob

func (q *DatabaseQueue) DeleteJob(ctx context.Context, j *jobs.Job) error

DeleteJob removes a finished job.

Deleted rather than marked done. A jobs table that keeps every job ever run is a table that needs its own cleanup job, and the history that matters -- what ran, how long it took, what it queried -- is on the console.

func (*DatabaseQueue) DeleteReserved

func (q *DatabaseQueue) DeleteReserved(ctx context.Context, queue, id string) error

DeleteReserved removes a reserved job by its id.

It is what a command reaches for when a job is stuck: the row is gone and no worker will pick it up when the lease expires.

func (*DatabaseQueue) FailJob

func (q *DatabaseQueue) FailJob(ctx context.Context, j *jobs.Job, cause error) error

FailJob parks the job.

func (*DatabaseQueue) Failed

func (q *DatabaseQueue) Failed(ctx context.Context, limit int) ([]jobs.Job, error)

Failed lists the jobs that gave up, most recent failure first.

func (*DatabaseQueue) GetConnectionName

func (c *DatabaseQueue) GetConnectionName() string

GetConnectionName is the name this queue was registered under.

It is empty on a queue built directly and never handed to a QueueManager, which is what a test does.

func (*DatabaseQueue) GetDatabase

func (q *DatabaseQueue) GetDatabase() *database.DB

GetDatabase is the connection this queue writes to.

The queue and the application share it, which is what makes a job pushed inside database.Transaction part of that transaction.

func (*DatabaseQueue) GetQueue

func (q *DatabaseQueue) GetQueue(name string) string

GetQueue resolves a queue name, turning empty into the default.

It is exported because every method here starts with it and a driver in another module has the same first line.

func (*DatabaseQueue) Later

func (q *DatabaseQueue) Later(ctx context.Context, g auth.Grant, delay time.Duration, j jobs.Job) error

Later adds a job that becomes eligible after delay.

func (*DatabaseQueue) Migrations

func (q *DatabaseQueue) Migrations() []migrations.Migration

Migrations returns the jobs table.

The schema is this driver's and only this one's. Module collects it, so an application wired to another driver declares no schema for a table it will never read.

func (*DatabaseQueue) PendingSize

func (q *DatabaseQueue) PendingSize(ctx context.Context, queue string) (int, error)

PendingSize is how many jobs are waiting.

func (*DatabaseQueue) Pop

func (q *DatabaseQueue) Pop(ctx context.Context, queue string, n int, lease time.Duration) ([]*jobs.Job, error)

Pop takes jobs off the queue and hides them for the lease.

Two statements rather than one, and the reason is portability: the tight form is UPDATE ... RETURNING with a FOR UPDATE SKIP LOCKED subquery, which Postgres has and SQLite does not. Selecting the candidates and then claiming each one by its id -- with reserved_until in the WHERE -- is correct on every engine, because the claim itself is the compare-and-set.

The cost is that two workers can pick the same candidate and one of them loses the claim. It gets nothing back, which is exactly right.

The candidates come back in the order they were queued. It used to be ORDER BY run_at, which put a job pushed with a delay behind everything queued while it waited -- forever, on a queue that is never empty. See the created_at migration.

COALESCE, because a row written by the previous release's binary during a rollout has no created_at and its run_at is the closest thing to one.

func (*DatabaseQueue) Push

func (q *DatabaseQueue) Push(ctx context.Context, g auth.Grant, j jobs.Job) error

Push adds a job.

Inside database.Transaction it joins it, which is the property this driver exists for: the job is committed by the same transaction as the row it describes, so it cannot refer to a write that rolled back.

func (*DatabaseQueue) PushOn

func (q *DatabaseQueue) PushOn(ctx context.Context, g auth.Grant, queue string, j jobs.Job) error

PushOn adds a job to a named queue.

func (*DatabaseQueue) PushRaw

func (q *DatabaseQueue) PushRaw(ctx context.Context, g auth.Grant, name string, payload []byte, queue string) error

PushRaw adds a job whose arguments are already encoded.

The envelope is the record's columns, so what is raw is the arguments: bytes the caller already has, that would only be unmarshalled and marshalled again on the way through jobs.New.

It goes through jobs.Authorized like every other push. Raw is about the encoding, never about the authorization.

func (*DatabaseQueue) Release

func (q *DatabaseQueue) Release(ctx context.Context, queue string, j *jobs.Job, delay time.Duration) error

Release puts a job back on its queue, eligible again after delay.

It is the driver-side half of jobs.Job.Release: a handler and a middleware call the job, and the job calls this.

func (*DatabaseQueue) ReleaseJob

func (q *DatabaseQueue) ReleaseJob(ctx context.Context, j *jobs.Job, delay time.Duration) error

ReleaseJob puts the job back on its queue, eligible again after delay.

It writes the exception count, and that is the whole reason a released job can ever be parked by MaxExceptions: the count has to cross deliveries. While the column went unwritten every delivery read zero, so a job with MaxExceptions 2 and MaxTries 5 was delivered all five times.

created_at moves to now, because a released job goes to the back of the queue: leaving the original there would put a job that has failed twice in front of work nobody has looked at yet.

func (*DatabaseQueue) ReservedSize

func (q *DatabaseQueue) ReservedSize(ctx context.Context, queue string) (int, error)

ReservedSize is how many jobs a worker is holding right now.

A number that stays high while PendingSize stays high is a worker that took work and stopped: the leases have not expired yet, so nothing is visibly wrong, and this is what shows it.

func (*DatabaseQueue) Retry

func (q *DatabaseQueue) Retry(ctx context.Context, uuid string) error

Retry puts a failed job back in line with its attempts reset.

The exception count is reset with them: a retry that kept the count would be parked again by the failure that parked it the first time, without the handler ever having run. created_at moves to now for the reason it moves on a release -- the job goes to the back of the queue.

func (*DatabaseQueue) SetConnectionName

func (c *DatabaseQueue) SetConnectionName(name string)

SetConnectionName names the connection this queue answers to.

It returns nothing rather than the queue, because an embedded struct cannot return the type that embeds it.

func (*DatabaseQueue) Size

func (q *DatabaseQueue) Size(ctx context.Context, queue string) (int, error)

Size is how many jobs the queue holds, waiting or in flight.

type DeferredQueue

type DeferredQueue struct {
	// contains filtered or unexported fields
}

DeferredQueue runs each job after the response has been sent, in this process.

It is SyncQueue with the work put off: nothing is stored, so nothing is retried and a restart loses whatever was in flight -- but the request does not wait for the work, which is the one thing SyncQueue cannot offer.

What it is for is the work that is genuinely optional: a webhook nobody waits for, a cache warm. Anything that must survive a deploy belongs on DatabaseQueue, and the difference between the two is exactly "may this be lost".

It is a distinct type rather than an option on SyncQueue because the choice is a wiring decision made once, and an option would be a second way to spell a connection name.

func NewDeferredQueue

func NewDeferredQueue(h Handlers) *DeferredQueue

NewDeferredQueue returns the queue over the same handler registry a worker uses.

The callback runs on its own goroutine, which is what "after the response" means here: the handler that pushed has already returned by the time the response is written.

func (*DeferredQueue) Bulk

func (q *DeferredQueue) Bulk(ctx context.Context, g auth.Grant, js []jobs.Job) error

Bulk defers each job in order.

func (*DeferredQueue) Clear

func (q *DeferredQueue) Clear(context.Context, string) (int, error)

Clear removes nothing.

func (*DeferredQueue) CreationTimeOfOldestPendingJob

func (q *DeferredQueue) CreationTimeOfOldestPendingJob(context.Context, string) (time.Time, error)

CreationTimeOfOldestPendingJob is the zero time: nothing waits.

func (*DeferredQueue) DeferUsing

func (q *DeferredQueue) DeferUsing(defer_ func(func())) *DeferredQueue

DeferUsing replaces how the work is put off, and returns the queue so the call chains.

It is what a test calls with func(fn func()) { fn() } to make the queue synchronous, which is the only way to assert on what a deferred job did without sleeping.

func (*DeferredQueue) DelayedSize

func (q *DeferredQueue) DelayedSize(context.Context, string) (int, error)

DelayedSize is zero. It answers delayedSize().

func (*DeferredQueue) DeleteJob

func (q *DeferredQueue) DeleteJob(context.Context, *jobs.Job) error

DeleteJob does nothing.

func (*DeferredQueue) FailJob

func (q *DeferredQueue) FailJob(context.Context, *jobs.Job, error) error

FailJob does nothing: nobody is left to hand the error to.

func (*DeferredQueue) Failed

func (q *DeferredQueue) Failed(context.Context, int) ([]jobs.Job, error)

Failed lists nothing. A deferred job that failed failed on a goroutine nobody is holding, and its error went to the log.

func (*DeferredQueue) GetConnectionName

func (c *DeferredQueue) GetConnectionName() string

GetConnectionName is the name this queue was registered under.

It is empty on a queue built directly and never handed to a QueueManager, which is what a test does.

func (*DeferredQueue) Later

func (q *DeferredQueue) Later(ctx context.Context, g auth.Grant, delay time.Duration, j jobs.Job) error

Later runs the job after the caller has moved on, and after delay.

The delay is honoured here where SyncQueue drops it, because the work is already off the request's path and sleeping costs nobody anything.

func (*DeferredQueue) PendingSize

func (q *DeferredQueue) PendingSize(context.Context, string) (int, error)

PendingSize is zero.

func (*DeferredQueue) Pop

Pop returns nothing: a deferred job runs in this process, not off a queue.

func (*DeferredQueue) Push

func (q *DeferredQueue) Push(ctx context.Context, g auth.Grant, j jobs.Job) error

Push runs the job after the caller has moved on.

The job is authorized before the deferral, not inside it: an unauthorized push has to fail the caller, and a failure that happens on another goroutine after the response is a failure nobody sees.

func (*DeferredQueue) PushOn

func (q *DeferredQueue) PushOn(ctx context.Context, g auth.Grant, queue string, j jobs.Job) error

PushOn runs the job on a named queue, after the caller has moved on.

func (*DeferredQueue) PushRaw

func (q *DeferredQueue) PushRaw(ctx context.Context, g auth.Grant, name string, payload []byte, queue string) error

PushRaw runs a job whose arguments are already encoded. It answers pushRaw().

func (*DeferredQueue) ReleaseJob

func (q *DeferredQueue) ReleaseJob(context.Context, *jobs.Job, time.Duration) error

ReleaseJob does nothing. There is no queue to put the job back on.

func (*DeferredQueue) ReservedSize

func (q *DeferredQueue) ReservedSize(context.Context, string) (int, error)

ReservedSize is zero. It answers reservedSize().

func (*DeferredQueue) Retry

Retry has nothing to retry.

func (*DeferredQueue) SetConnectionName

func (c *DeferredQueue) SetConnectionName(name string)

SetConnectionName names the connection this queue answers to.

It returns nothing rather than the queue, because an embedded struct cannot return the type that embeds it.

func (*DeferredQueue) Size

Size is zero. Nothing is stored.

type Dispatcher

type Dispatcher interface {
	// Dispatch sends an event to every listener and returns what they returned.
	Dispatch(event any, payload ...any) []any
	// Until sends an event and stops at the first listener that answers
	// something other than nil.
	Until(event any, payload ...any) any
}

Dispatcher is the two things the queue asks of an event dispatcher.

It is declared here rather than imported so the queue does not depend on the events package to run, and so a test can watch what a worker announced with eight lines and no wiring. events.Dispatcher satisfies it.

Until is on it and not just Dispatch because of one listener: a Looping listener that returns false stops the worker taking work, and halting on the first answer is the only way to hear it.

type FailedConfig added in v0.2.0

type FailedConfig struct {
	// Driver is "database" or empty for none.
	Driver string

	// Database is the connection the failed table lives on, empty meaning the
	// application's default.
	Database string

	// Table is the table itself, and empty means "failed_jobs".
	Table string
}

FailedConfig is where a job goes when it has run out of tries.

type FailoverQueue

type FailoverQueue struct {
	// contains filtered or unexported fields
}

FailoverQueue writes to the first connection that accepts a job, and reads from the first one only.

It is what stands between "the broker is down" and "the checkout returned a 500": the push tries each connection in order, and a job that could not go on Redis goes in the database instead.

The asymmetry is deliberate. Only the push fails over; every read -- Pop, Size, the oldest pending job -- goes to the first connection. A worker that drained all of them would be a worker whose lease and whose ordering mean different things per job, and two workers, one per connection, is the arrangement that keeps both simple. The jobs that landed on the fallback are drained by a worker pointed at the fallback.

It is not a way to spread load. Every job goes to the first connection that works, so as long as the first one is up the others are empty.

func NewFailoverQueue

func NewFailoverQueue(m *QueueManager, names ...string) *FailoverQueue

NewFailoverQueue returns the queue over connections, in the order they should be tried.

The first is the one that is read from, so it is the one a worker drains and the one the health check measures.

func (*FailoverQueue) Bulk

func (q *FailoverQueue) Bulk(ctx context.Context, g auth.Grant, js []jobs.Job) error

Bulk writes each job to the first connection that accepts it.

Per job rather than per batch, because a connection that goes down halfway through would otherwise leave half the batch on one store and half on another, with nothing recording which half.

func (*FailoverQueue) Clear

func (q *FailoverQueue) Clear(ctx context.Context, queue string) (int, error)

Clear empties every connection, and returns how many went in total.

Every one rather than the first: clearing is the operation whose whole point is that nothing is left, and a clear that emptied one store while jobs sat on the fallback would be the most misleading command in the collection.

func (*FailoverQueue) CreationTimeOfOldestPendingJob

func (q *FailoverQueue) CreationTimeOfOldestPendingJob(ctx context.Context, queue string) (time.Time, error)

CreationTimeOfOldestPendingJob is the first connection's. It answers creationTimeOfOldestPendingJob().

func (*FailoverQueue) DelayedSize

func (q *FailoverQueue) DelayedSize(ctx context.Context, queue string) (int, error)

DelayedSize is the first connection's, when it can answer. It answers delayedSize().

func (*FailoverQueue) DeleteJob

func (q *FailoverQueue) DeleteJob(context.Context, *jobs.Job) error

DeleteJob does nothing, for the reason FailoverQueue.ReleaseJob gives.

func (*FailoverQueue) FailJob

func (q *FailoverQueue) FailJob(context.Context, *jobs.Job, error) error

FailJob does nothing, for the reason FailoverQueue.ReleaseJob gives.

func (*FailoverQueue) Failed

func (q *FailoverQueue) Failed(ctx context.Context, limit int) ([]jobs.Job, error)

Failed lists the parked jobs on every connection, for the same reason Clear touches every one: a job that failed on the fallback is still a job somebody has to see.

func (*FailoverQueue) GetConnectionName

func (c *FailoverQueue) GetConnectionName() string

GetConnectionName is the name this queue was registered under.

It is empty on a queue built directly and never handed to a QueueManager, which is what a test does.

func (*FailoverQueue) Later

func (q *FailoverQueue) Later(ctx context.Context, g auth.Grant, delay time.Duration, j jobs.Job) error

Later writes a delayed job to the first connection that accepts it.

func (*FailoverQueue) PendingSize

func (q *FailoverQueue) PendingSize(ctx context.Context, queue string) (int, error)

PendingSize is the first connection's. It answers pendingSize().

func (*FailoverQueue) Pop

func (q *FailoverQueue) Pop(ctx context.Context, queue string, n int, lease time.Duration) ([]*jobs.Job, error)

Pop takes jobs off the first connection. It answers pop().

func (*FailoverQueue) Push

func (q *FailoverQueue) Push(ctx context.Context, g auth.Grant, j jobs.Job) error

Push writes the job to the first connection that accepts it.

func (*FailoverQueue) PushOn

func (q *FailoverQueue) PushOn(ctx context.Context, g auth.Grant, queue string, j jobs.Job) error

PushOn writes the job to a named queue on the first connection that accepts it.

func (*FailoverQueue) PushRaw

func (q *FailoverQueue) PushRaw(ctx context.Context, g auth.Grant, name string, payload []byte, queue string) error

PushRaw writes a job whose arguments are already encoded. It answers pushRaw().

func (*FailoverQueue) ReleaseJob

func (q *FailoverQueue) ReleaseJob(context.Context, *jobs.Job, time.Duration) error

ReleaseJob settles the job against the queue it actually came off.

A popped job carries its own driver -- jobs.Popped attached it -- so these three are only reached by a job built by hand, and there is nothing to settle.

func (*FailoverQueue) ReservedSize

func (q *FailoverQueue) ReservedSize(ctx context.Context, queue string) (int, error)

ReservedSize is the first connection's, when it can answer. It answers reservedSize().

func (*FailoverQueue) Retry

func (q *FailoverQueue) Retry(ctx context.Context, uuid string) error

Retry puts a failed job back on whichever connection has it.

func (*FailoverQueue) SetConnectionName

func (c *FailoverQueue) SetConnectionName(name string)

SetConnectionName names the connection this queue answers to.

It returns nothing rather than the queue, because an embedded struct cannot return the type that embeds it.

func (*FailoverQueue) SetEvents

func (q *FailoverQueue) SetEvents(d Dispatcher) *FailoverQueue

SetEvents gives the queue somewhere to send QueueFailedOver, and returns it so the call chains.

func (*FailoverQueue) Size

func (q *FailoverQueue) Size(ctx context.Context, queue string) (int, error)

Size is the first connection's. It answers size().

type Handler

type Handler interface {
	Handle(ctx context.Context, g auth.Grant, j *jobs.Job) error
}

Handler does the work.

The Grant is rebuilt from the job's tenant and action, so a handler reaches repositories the same way a service does. There is no unauthorized path into the database from a worker, which is the whole point of the Grant existing.

The record names the work and something has to turn that name into code: the name is a string and the registry is the Worker, because Go cannot reach a type from a name in a string.

type HandlerFunc

type HandlerFunc func(ctx context.Context, g auth.Grant, j *jobs.Job) error

HandlerFunc adapts a function to Handler.

func (HandlerFunc) Handle

func (f HandlerFunc) Handle(ctx context.Context, g auth.Grant, j *jobs.Job) error

Handle calls f.

type HandlerPanicked added in v0.8.0

type HandlerPanicked struct {
	// Job is the job whose handler panicked.
	Job *jobs.Job
	// Value is what was passed to panic.
	Value any
	// Stack is the traceback of the goroutine that panicked.
	Stack []byte
}

HandlerPanicked is why a job was parked: its handler panicked instead of returning an error.

It carries the recovered value and the stack captured where the panic was recovered, because the value alone says what went wrong and only the stack says where. Both end up in the job's LastError, which is the row somebody reads days later with no process left to attach a debugger to.

func (*HandlerPanicked) Error added in v0.8.0

func (e *HandlerPanicked) Error() string

Error is the message, with the traceback after it.

func (HandlerPanicked) ForJob added in v0.8.0

func (HandlerPanicked) ForJob(j *jobs.Job, value any, stack []byte) *HandlerPanicked

ForJob returns the error for a job whose handler panicked, in the shape the other two errors here are built in.

func (*HandlerPanicked) Is added in v0.8.0

func (e *HandlerPanicked) Is(target error) bool

Is makes errors.Is(err, ErrHandlerPanicked) true for this error.

func (*HandlerPanicked) Unwrap added in v0.8.0

func (e *HandlerPanicked) Unwrap() error

Unwrap returns the panicked value when it was an error, so errors.Is and errors.As reach whatever the handler panicked with. A panic carrying anything else has nothing underneath it to match.

type Handlers

type Handlers interface {
	Handler(name string) (Handler, bool)
}

Handlers resolves a job name to the handler that runs it.

Worker implements it, and it is declared as an interface so SyncQueue can borrow the worker's registry without the two constructors having to be called in a circle.

type HandlesFailure

type HandlesFailure interface {
	// Failed is called once, after the job has been parked.
	Failed(ctx context.Context, g auth.Grant, j *jobs.Job, cause error)
}

HandlesFailure is a handler that wants to know when its job gave up.

A handler that does not implement it is not called.

type HasDisplayName

type HasDisplayName interface {
	// DisplayName is what a console, a log line and the failed job list show.
	DisplayName() string
}

HasDisplayName is a job value that names itself for a person.

A value that does not implement it is named by the job name it was pushed under.

type Identifiable

type Identifiable interface {
	// ModelIdentifier is what finds this record again.
	ModelIdentifier() ModelIdentifier
}

Identifiable is a record that can say what finds it again.

It is what a domain type implements so SerializesModels can reduce it: a record says what finds it again, rather than a base class being recognized.

type InteractsWithQueue

type InteractsWithQueue struct {
	// Job is the job being run, or nil outside a worker.
	//
	// It is exported so a handler that wants the payload reads it here rather
	// than being handed a fourth accessor for it.
	Job *jobs.Job
	// contains filtered or unexported fields
}

InteractsWithQueue is what a handler embeds to reach the job it is running.

A handler is already given the job -- Handler.Handle takes it -- so this is for the other shape: a handler that is a struct with state, whose methods settle the job on themselves rather than threading it through.

type SendInvoice struct {
	queue.InteractsWithQueue
	Invoices *invoice.Repository
}

func (h *SendInvoice) Handle(ctx context.Context, g auth.Grant, j *jobs.Job) error {
	h.SetJob(j)
	if !h.Invoices.Ready(ctx, g) {
		return h.Release(ctx, time.Minute)
	}
	...
}

A zero value with no job is safe: every method answers as though the job had already been settled.

func (*InteractsWithQueue) AssertDeleted

func (i *InteractsWithQueue) AssertDeleted(t TestingT) *InteractsWithQueue

AssertDeleted fails the test unless the job was deleted.

func (*InteractsWithQueue) AssertFailed

func (i *InteractsWithQueue) AssertFailed(t TestingT) *InteractsWithQueue

AssertFailed fails the test unless the job was parked.

func (*InteractsWithQueue) AssertFailedWith

func (i *InteractsWithQueue) AssertFailedWith(t TestingT, want error) *InteractsWithQueue

AssertFailedWith fails the test unless the job was parked with a cause that matches want under errors.Is.

errors.Is rather than equality, because a handler that wrapped the cause with context still failed with it.

func (*InteractsWithQueue) AssertNotDeleted

func (i *InteractsWithQueue) AssertNotDeleted(t TestingT) *InteractsWithQueue

AssertNotDeleted fails the test if the job was deleted.

func (*InteractsWithQueue) AssertNotFailed

func (i *InteractsWithQueue) AssertNotFailed(t TestingT) *InteractsWithQueue

AssertNotFailed fails the test if the job was parked.

func (*InteractsWithQueue) AssertNotReleased

func (i *InteractsWithQueue) AssertNotReleased(t TestingT) *InteractsWithQueue

AssertNotReleased fails the test if the job was released.

func (*InteractsWithQueue) AssertReleased

func (i *InteractsWithQueue) AssertReleased(t TestingT, delay time.Duration) *InteractsWithQueue

AssertReleased fails the test unless the job was released.

A negative delay asserts only that the job was released: a test that cares when it comes back says so, and one that does not should not have to.

func (*InteractsWithQueue) Attempts

func (i *InteractsWithQueue) Attempts() int

Attempts is how many times this job has been delivered, counting now.

A handler with no job answers 1: code that branches on "is this the first try" must behave outside a worker as it does on the first try.

func (*InteractsWithQueue) Delete

func (i *InteractsWithQueue) Delete(ctx context.Context) error

Delete removes the job from its queue.

func (*InteractsWithQueue) Fail

func (i *InteractsWithQueue) Fail(ctx context.Context, cause error) error

Fail parks the job.

A nil cause parks it with ErrManuallyFailed: the record has to say something, and "the handler asked for it" is the truth.

func (*InteractsWithQueue) Release

func (i *InteractsWithQueue) Release(ctx context.Context, delay time.Duration) error

Release puts the job back on its queue, eligible again after delay.

func (*InteractsWithQueue) SetJob

SetJob gives the handler the job it is running, and returns the receiver so the call chains.

func (*InteractsWithQueue) String

func (i *InteractsWithQueue) String() string

String names the job this handler is running, for a test failure to print.

func (*InteractsWithQueue) WithFakeQueueInteractions

func (i *InteractsWithQueue) WithFakeQueueInteractions() *InteractsWithQueue

WithFakeQueueInteractions makes release, delete and fail record what was asked for instead of doing it, and returns the receiver so the call chains.

It is what a test calls before running the handler directly, so the handler can be asked what it decided without a store, a worker or a transaction.

type Listener

type Listener struct {
	// contains filtered or unexported fields
}

Listener runs a worker in a child process and restarts it when it exits.

It is what `aru queue:listen` runs, and it is not the way to run a queue in production -- Worker is. What it is for is development: a rebuilt binary is picked up without anybody remembering to restart anything.

The child is started with Command, which is the path of the binary and the arguments before the ones this adds. An application whose worker command is not `aru work` says so there.

func NewListener

func NewListener(dir string, command ...string) *Listener

NewListener returns a listener that starts command in dir.

The binary is the thing being run, so the caller says what it is:

queue.NewListener(".", os.Args[0], "work")

func (*Listener) Listen

func (l *Listener) Listen(ctx context.Context, connectionName, queue string, options ListenerOptions) error

Listen starts a worker for connection and queue, and restarts it whenever it exits.

It returns when the context is cancelled or Listener.Stop is called.

A child that exits with ExitMemoryLimit is started again, because that is what it stopped for. A child that cannot be started at all is an error: a loop that retries a binary that does not exist is a loop that fills a disk with the same line.

func (*Listener) MakeProcess

func (l *Listener) MakeProcess(ctx context.Context, connectionName, queue string, options ListenerOptions) *exec.Cmd

MakeProcess builds the child worker's command.

It answers makeProcess(). The flags are the worker's options turned back into the arguments `aru work` parses, which is the round trip that makes a listener and a worker configurable in one place.

func (*Listener) MemoryExceeded

func (l *Listener) MemoryExceeded(limitMB int) bool

MemoryExceeded reports whether this process is holding more than limit megabytes.

It reads MemStats.Sys, for the reason Worker.MemoryExceeded gives.

func (*Listener) RunProcess

func (l *Listener) RunProcess(process *exec.Cmd, memory int) error

RunProcess runs one child worker to completion.

The memory check afterwards is about this process, the listener: it has grown, and something outside it should start it again with a clean heap.

func (*Listener) SetOutputHandler

func (l *Listener) SetOutputHandler(handler func(line string))

SetOutputHandler sets what to do with each line the child writes.

Without one the child's output goes nowhere, which is what a test wants; `aru queue:listen` passes one that writes to the terminal.

func (*Listener) Stop

func (l *Listener) Stop()

Stop ends the loop after the child that is running.

It cancels rather than ending the process: a listener that killed the process would take the child with it, and the child is holding a job.

type ListenerOptions

type ListenerOptions struct {
	// WorkerOptions is what each child worker runs with. The listener turns
	// them back into the flags it starts the child with.
	WorkerOptions

	// Environment is the environment the child workers run in, and empty means
	// they inherit this process's.
	Environment string
}

ListenerOptions configures a Listener.

It is WorkerOptions with one field added: the child workers run with the embedded options, and the listener turns them back into flags.

type MaxAttemptsExceeded

type MaxAttemptsExceeded struct {
	// Job is the job that ran out. It is never nil on an error built by
	// [MaxAttemptsExceeded.ForJob].
	Job *jobs.Job
	// Attempts is how many deliveries it had, and Max is how many it was
	// allowed. Both are on the error because "attempted too many times" without
	// the two numbers is the least useful sentence a log can carry.
	Attempts int
	Max      int
}

MaxAttemptsExceeded is why a job was parked: it ran out of deliveries.

It keeps the job, so a listener on the failed-job event can say which one and on which queue without decoding anything.

func (*MaxAttemptsExceeded) Error

func (e *MaxAttemptsExceeded) Error() string

Error is the message.

func (MaxAttemptsExceeded) ForJob

ForJob returns the error for a job that ran out of deliveries.

It is a method on the zero value rather than a package function, so `queue.MaxAttemptsExceeded{}.ForJob(j, 5)` names the error and the reason for it in one phrase.

func (*MaxAttemptsExceeded) Is

func (e *MaxAttemptsExceeded) Is(target error) bool

Is makes errors.Is(err, ErrMaxAttemptsExceeded) true for this error.

type ModelFinder

type ModelFinder func(ctx context.Context, g auth.Grant, id ModelIdentifier) (any, error)

ModelFinder loads a record back from its identifier.

It is what the application registers so a job can carry an id instead of a document. It takes the Grant the job runs under, which is what makes the reload obey the same policy the original read did.

type ModelIdentifier

type ModelIdentifier struct {
	// Class names the kind of record. It is a string a finder is registered
	// under -- "invoice", "user" -- because Go cannot reach a type from a name
	// in a string.
	Class string
	// ID is the primary key.
	ID string
	// Relations are the relations to load with it. Empty means none, which is
	// the right default: a job that needs a relation says so, and one that does
	// not should not pay for it.
	Relations []string
	// Connection names the database connection it came from, and empty means
	// the default.
	Connection string
	// TenantID is who it belongs to, and it is the thing that makes restoring
	// one safe: the finder is given a Grant built for this tenant, so a payload
	// that names another customer's row finds nothing.
	TenantID string
}

ModelIdentifier is a record reduced to what finds it again.

It is what SerializesModels puts on the wire in place of a record. Class is what the application registered a finder under, ID is the primary key, and Relations are the relations to load back with it.

type Module

type Module struct {
	// contains filtered or unexported fields
}

Module brings the jobs table, and reports on the queue.

It registers no routes and runs no worker: the worker is `aru work`, a separate process from the same image. What this module owns is the schema and the answer to "is anything draining this".

func NewModule

func NewModule(q Queue, queues ...string) *Module

NewModule returns the module.

Name the queues the application actually uses. A queue nobody watches is a queue that can stop draining without anyone noticing, which is the failure this module exists to make visible.

func (*Module) Diagnose

func (m *Module) Diagnose(ctx context.Context) []string

Diagnose reports a backlog and a dead letter queue, on the error page.

Next to the failure somebody is already looking at, which is the moment they are most likely to act on it.

func (*Module) Health

func (m *Module) Health(ctx context.Context) error

Health fails when a queue stops draining.

func (*Module) Migrations

func (m *Module) Migrations() []foundation.Migration

Migrations returns the schema the wired driver needs.

The jobs table belongs to DatabaseQueue and is declared there. An application wired to a driver that stores jobs elsewhere gets nothing here, rather than an empty table it will never read -- which is why this asks the driver instead of answering for it.

func (*Module) Name

func (*Module) Name() string

Name is the module identifier.

func (*Module) Routes

func (*Module) Routes(*routing.Router)

Routes registers nothing.

type NullQueue

type NullQueue struct{}

NullQueue accepts every job and keeps none of them.

It is what a test that does not care about the queue wires, and what the sync connection's registry-only worker is built over -- see SyncQueue.

It is a value and not a pointer, because it has no state and a nil pointer that silently swallows every job is a worse mistake than the one this type exists to make cheap.

func (NullQueue) Bulk

Bulk discards the jobs.

func (NullQueue) Clear

func (NullQueue) Clear(context.Context, string) (int, error)

Clear removes nothing.

func (NullQueue) CreationTimeOfOldestPendingJob

func (NullQueue) CreationTimeOfOldestPendingJob(context.Context, string) (time.Time, error)

CreationTimeOfOldestPendingJob is the zero time.

func (NullQueue) DelayedSize

func (NullQueue) DelayedSize(context.Context, string) (int, error)

DelayedSize is zero. It answers delayedSize().

func (NullQueue) DeleteJob

func (NullQueue) DeleteJob(context.Context, *jobs.Job) error

DeleteJob does nothing.

func (NullQueue) FailJob

func (NullQueue) FailJob(context.Context, *jobs.Job, error) error

FailJob does nothing.

func (NullQueue) Failed

func (NullQueue) Failed(context.Context, int) ([]jobs.Job, error)

Failed lists nothing.

func (NullQueue) GetConnectionName

func (NullQueue) GetConnectionName() string

GetConnectionName is "null".

It answers getConnectionName() without the embedded connection struct the other drivers use, because this one is a value with no state and giving it a settable name would be the one field on it.

func (NullQueue) Later

Later discards the job.

func (NullQueue) PendingSize

func (NullQueue) PendingSize(context.Context, string) (int, error)

PendingSize is zero.

func (NullQueue) Pop

Pop returns nothing.

func (NullQueue) Push

Push discards the job.

func (NullQueue) PushOn

PushOn discards the job.

func (NullQueue) PushRaw

PushRaw discards the job. It answers pushRaw().

func (NullQueue) ReleaseJob

ReleaseJob does nothing.

func (NullQueue) ReservedSize

func (NullQueue) ReservedSize(context.Context, string) (int, error)

ReservedSize is zero. It answers reservedSize().

func (NullQueue) Retry

Retry has nothing to retry.

func (NullQueue) Size

Size is zero.

type PayloadHook

type PayloadHook func(connection, queue string, j *jobs.Job)

PayloadHook may add to a job on its way onto a queue.

It is what CreatePayloadUsing registers, and it is how an observability package adds a trace id to every job it never wrote. It is given the connection, the queue and the job as it stands; what it changes, it changes on the job.

It takes a pointer rather than returning something to be merged over the record, because merging untyped maps is how a hook silently overwrites maxTries.

type PopCallback

type PopCallback func(ctx context.Context, pop func(queue string) ([]*jobs.Job, error)) ([]*jobs.Job, error)

PopCallback replaces how a worker chooses its next batch.

It is what PopUsing registers: pop is the driver's own pop, and what the callback returns is what the worker runs. It is the seam for weighting queues against each other.

type Queue

type Queue interface {
	// Push adds a job to the queue named on it.
	//
	// Inside database.Transaction the DatabaseQueue joins it, which is the
	// property that driver exists for: the job is committed by the same
	// transaction as the row it describes, so it cannot refer to a write that
	// rolled back.
	Push(ctx context.Context, g auth.Grant, j jobs.Job) error

	// PushOn adds a job to a named queue, ignoring the one on the job.
	PushOn(ctx context.Context, g auth.Grant, queue string, j jobs.Job) error

	// Later adds a job that becomes eligible after delay.
	Later(ctx context.Context, g auth.Grant, delay time.Duration, j jobs.Job) error

	// Bulk adds many jobs. A driver that can do it in one round trip does; the
	// contract only promises that either all of them arrive or an error says
	// otherwise.
	Bulk(ctx context.Context, g auth.Grant, js []jobs.Job) error

	// Pop takes up to n jobs off a queue and hides them for the lease.
	//
	// Jobs whose lease expires become visible again -- which is what makes a
	// worker crash recoverable and delivery at-least-once.
	//
	// The returned jobs carry Attempts INCLUDING this delivery: a job handed
	// over for the first time has Attempts == 1. A driver that returns the
	// count from before the delivery makes the worker park a job one attempt
	// early, and with a try limit of 2 it parks on the first failure and never
	// retries at all.
	//
	// Each job carries the queue it came off, so the worker settles it by
	// calling Release, Delete or Fail on the job rather than on the driver.
	Pop(ctx context.Context, queue string, n int, lease time.Duration) ([]*jobs.Job, error)

	// Size is how many jobs the queue holds, waiting or in flight.
	Size(ctx context.Context, queue string) (int, error)

	// PendingSize is how many jobs are waiting. It feeds the health check: a
	// queue that only grows is a worker that is not running.
	//
	// A job whose lease is still valid is running, not waiting.
	PendingSize(ctx context.Context, queue string) (int, error)

	// CreationTimeOfOldestPendingJob is when the oldest waiting job became
	// eligible, or the zero time when nothing is waiting.
	//
	// A stopped worker looks exactly like an idle one, and this is what tells
	// them apart. A job scheduled for the future returns a time in the future,
	// so a caller measuring lag with time.Since gets a negative number and can
	// treat it as no lag at all.
	CreationTimeOfOldestPendingJob(ctx context.Context, queue string) (time.Time, error)

	// Clear removes every job on a queue and returns how many went.
	Clear(ctx context.Context, queue string) (int, error)

	// Failed lists the jobs that gave up, most recent failure first.
	Failed(ctx context.Context, limit int) ([]jobs.Job, error)

	// Retry puts a failed job back in line with its attempts reset.
	//
	// Without it the only way out of a dead letter queue is SQL by hand, which
	// is how it becomes a table nobody touches.
	Retry(ctx context.Context, uuid string) error
}

Queue is what a driver implements.

Every push takes an auth.Grant, so the tenant comes from the Grant and not from an argument somebody can get wrong, and Pop takes a count and a lease rather than returning one job at a time -- a worker with a concurrency of four asks for four.

Pop rather than a channel of jobs: the job has to stay in the store, invisible to other workers, until it is settled. A worker that dies mid-job must not lose it, and that is not something a channel can offer.

type QueueManager

type QueueManager struct {
	// contains filtered or unexported fields
}

QueueManager holds the configured connections and hands them out by name.

There is no resolution from a config array: the application constructs its queues in bootstrap/app.go and passes them here. What is left is the part anybody calls -- `aru work --connection=redis` has to turn a string into a queue, the default connection has to have a name, and a queue has to be pausable while an incident is open.

m := queue.NewQueueManager().
	Extend("database", queue.NewDatabaseQueue(db)).
	Extend("sync", queue.NewSyncQueue(w))

func NewQueueManager

func NewQueueManager() *QueueManager

NewQueueManager returns an empty manager.

func (*QueueManager) AddConnector

func (m *QueueManager) AddConnector(name string, connect func() (Queue, error)) *QueueManager

AddConnector registers how to build a connection the first time it is asked for.

It is the one place lazy resolution earns its keep: a RESP connection opens a socket, and a binary that runs `aru migrate` should not open one on the way past.

The first one registered becomes the default, exactly as with QueueManager.Extend.

func (*QueueManager) After

func (m *QueueManager) After(listener any)

After registers a listener for the event after a job ran.

func (*QueueManager) Before

func (m *QueueManager) Before(listener any)

Before registers a listener for the event before a job runs.

This and the six methods under it are sugar over registering a listener for one event type, and they exist because `manager.Failing(...)` in bootstrap/app.go says what it does where the naked registration does not.

func (*QueueManager) Connected

func (m *QueueManager) Connected(name string) bool

Connected reports whether a connection has already been built.

It is false for a name registered with QueueManager.AddConnector and never asked for, and true once it has been.

func (*QueueManager) Connection

func (m *QueueManager) Connection(name string) (Queue, error)

Connection returns a registered queue. An empty name means the default.

A connector runs once and its queue is kept.

func (*QueueManager) Connections

func (m *QueueManager) Connections() []string

Connections lists the registered names, for `aru queue:monitor` to print.

Sorted, so two runs of a command print the same order -- a map's is random, and a listing that reshuffles is a listing nobody can diff.

func (*QueueManager) ExceptionOccurred

func (m *QueueManager) ExceptionOccurred(listener any)

ExceptionOccurred registers a listener for a job that returned an error.

func (*QueueManager) Extend

func (m *QueueManager) Extend(name string, q Queue) *QueueManager

Extend registers a connection under a name.

It takes the queue itself, because there is no config array to build it from and lazily resolving something the application already constructed is a delay with no payoff. QueueManager.AddConnector is the lazy form, for a driver whose connection is expensive to open.

The first one registered becomes the default, so an application with one queue never has to say which it means.

func (*QueueManager) Failing

func (m *QueueManager) Failing(listener any)

Failing registers a listener for a job that gave up.

func (*QueueManager) GetDefaultDriver

func (m *QueueManager) GetDefaultDriver() string

GetDefaultDriver is the name QueueManager.Connection uses when asked for "".

func (*QueueManager) GetName

func (m *QueueManager) GetName(name string) string

GetName is the connection name to use, resolving "" to the default.

It is what a command prints, so a log line says "database" rather than nothing.

func (*QueueManager) GetRoutes

func (m *QueueManager) GetRoutes() *QueueRoutes

GetRoutes is the routing table, for a dispatcher to consult. It is one object with one owner.

func (*QueueManager) IsPaused

func (m *QueueManager) IsPaused(ctx context.Context, connection, queue string) (bool, error)

IsPaused reports whether a queue is paused. It answers isPaused().

A manager with no cache answers false: nothing can have paused a queue when there is nowhere to record a pause.

func (*QueueManager) Looping

func (m *QueueManager) Looping(listener any)

Looping registers a listener for each pass of a worker's loop.

func (*QueueManager) Pause

func (m *QueueManager) Pause(ctx context.Context, connection, queue string) error

Pause stops workers taking new jobs off a queue, until it is resumed.

It answers pause(). Jobs already running finish; jobs already queued stay queued. It is the switch to reach for when a downstream system is failing and retrying into it is making things worse.

func (*QueueManager) PauseFor

func (m *QueueManager) PauseFor(ctx context.Context, connection, queue string, ttl time.Duration) error

PauseFor stops workers taking new jobs off a queue for a while, after which it resumes on its own.

It answers pauseFor(). It is the safer half of QueueManager.Pause: a pause with a deadline cannot be the thing somebody forgot to undo.

func (*QueueManager) Restart

func (m *QueueManager) Restart(ctx context.Context) error

Restart asks every running worker to stop after the job it is holding.

It writes the timestamp the workers poll, and it lives on the manager rather than in `queue:restart` so the command is three lines and the key is named once.

It is how a deploy replaces the workers: the new binary starts, the old ones notice the timestamp moved and exit cleanly, and whatever supervises them starts the new image. Nothing about it runs at boot.

func (*QueueManager) Resume

func (m *QueueManager) Resume(ctx context.Context, connection, queue string) error

Resume lets workers take jobs off a paused queue again. It answers resume().

func (*QueueManager) Route

func (m *QueueManager) Route(name, queue, connection string)

Route sends a job name to a connection and a queue.

See QueueRoutes for what it buys.

func (*QueueManager) SetCache

func (m *QueueManager) SetCache(c Cache) *QueueManager

SetCache gives the manager somewhere to keep the pause flags and the restart signal, and returns it so the call chains.

Without one, QueueManager.Pause and QueueManager.Restart have nowhere to write and say so, rather than silently doing nothing -- a queue somebody believes is paused and is not is worse than an error.

func (*QueueManager) SetDefaultDriver

func (m *QueueManager) SetDefaultDriver(name string)

SetDefaultDriver names the connection QueueManager.Connection uses when asked for "".

func (*QueueManager) SetEvents

func (m *QueueManager) SetEvents(d Dispatcher) *QueueManager

SetEvents gives the manager somewhere to send QueuePaused and QueueResumed, and returns it so the call chains.

func (*QueueManager) Starting

func (m *QueueManager) Starting(listener any)

Starting registers a listener for a worker starting.

func (*QueueManager) Stopping

func (m *QueueManager) Stopping(listener any)

Stopping registers a listener for a worker stopping.

func (*QueueManager) WithoutInterruptionPolling

func (m *QueueManager) WithoutInterruptionPolling()

WithoutInterruptionPolling stops workers reading the cache to find out whether they should pause or restart.

Every worker asking a shared cache once a second is a load somebody eventually notices, and an application whose deploy stops the processes itself does not need the signal.

It is process-wide: a worker built after this is called does not poll either.

type QueueRoutes

type QueueRoutes struct {
	// contains filtered or unexported fields
}

QueueRoutes says which connection and which queue a job name belongs on.

It exists so "reports go on the slow queue" is written once, at boot, instead of at every dispatch -- and so that moving them is one line rather than a search.

m.Route("report.monthly", "reports", "")
m.Route("invoice.send", "", "redis")

The key is the job name and there is no hierarchy to walk: a name is a name. What that costs is a route inherited by a family of jobs; what it buys is that the routing of a job is one lookup a person can predict.

func NewQueueRoutes

func NewQueueRoutes() *QueueRoutes

NewQueueRoutes returns an empty table.

func (*QueueRoutes) All

func (r *QueueRoutes) All() map[string]Route

All is every registered route. It answers all().

The copy is not a courtesy: the table is read by whatever dispatches jobs, concurrently, and handing out the map would be a data race with a nice name.

func (*QueueRoutes) GetConnection

func (r *QueueRoutes) GetConnection(name string) string

GetConnection is the connection a job name belongs on, or empty for the default. It answers getConnection().

func (*QueueRoutes) GetQueue

func (r *QueueRoutes) GetQueue(name string) string

GetQueue is the queue a job name belongs on, or empty for the default. It answers getQueue().

func (*QueueRoutes) GetRoute

func (r *QueueRoutes) GetRoute(name string) (Route, bool)

GetRoute is the route for a job name, and whether there is one.

func (*QueueRoutes) Names

func (r *QueueRoutes) Names() []string

Names is every routed job name, sorted, for `aru queue:monitor` to print.

func (*QueueRoutes) Set

func (r *QueueRoutes) Set(name, queue, connection string)

Set registers the route for a job name.

An empty queue or connection means "the default", and setting a name twice replaces the first -- the table is built at boot, in one place, and the last line wins the way the last assignment does.

type Route

type Route struct {
	// Connection is the queue connection, by the name it was registered under.
	// Empty means the default.
	Connection string
	// Queue is the queue on that connection. Empty means jobs.DefaultQueue.
	Queue string
}

Route is where a job goes when nothing at the call site said.

It is a struct rather than a pair, because two strings side by side are two strings somebody eventually puts in the wrong order.

type SerializesModels

type SerializesModels struct {
	// contains filtered or unexported fields
}

SerializesModels keeps records out of job payloads.

The rule it enforces is that a payload carries ids and facts, never a document: a record goes on the wire as a ModelIdentifier and is loaded back on the other side. This type is that rule with a name and a finder registry, and SerializesModels.RestoreModel is the reload.

var models queue.SerializesModels
models.FindModelsUsing("invoice", func(ctx context.Context, g auth.Grant, id queue.ModelIdentifier) (any, error) {
	return invoices.Find(ctx, g, id.ID)
})

Why it matters is not size. A job that carries a serialized record acts on the record as it was when the job was queued, and the queue exists precisely because time passes between those two moments: the invoice was voided, the address was corrected, the user was deleted. Reloading is what makes the job act on what is true now.

func (*SerializesModels) FindModelsUsing

func (s *SerializesModels) FindModelsUsing(class string, find ModelFinder) *SerializesModels

FindModelsUsing registers how to load a kind of record back.

The class on an identifier is a string, and this is what turns it into code -- the same trade the Worker makes for job names, because Go cannot reach a type from a name in a string.

func (*SerializesModels) GetRestoredPropertyValue

func (s *SerializesModels) GetRestoredPropertyValue(ctx context.Context, action auth.Action, value any) (any, error)

GetRestoredPropertyValue is the record an identifier names, or the value itself when it is not one.

func (*SerializesModels) GetSerializedPropertyValue

func (s *SerializesModels) GetSerializedPropertyValue(value any) any

GetSerializedPropertyValue is what goes on the wire in place of a record.

A value that can identify itself is reduced to its identifier; everything else is passed through, because a payload of plain facts is already what it should be.

func (*SerializesModels) RestoreModel

func (s *SerializesModels) RestoreModel(ctx context.Context, action auth.Action, id ModelIdentifier) (any, error)

RestoreModel loads the record an identifier names.

The Grant is rebuilt from the identifier's tenant, so the reload is scoped to the customer the job belongs to and a payload naming another customer's row finds nothing.

A record that is gone comes back as an error wrapping ErrMissingModel, not as a nil the handler has to remember to check.

type SyncQueue

type SyncQueue struct {
	// contains filtered or unexported fields
}

SyncQueue runs each job the moment it is pushed, in the caller's goroutine.

It exists for two reasons: a test that pushes a job wants to assert on what the job did without starting a worker, and a developer running the application locally wants the email to be sent rather than to sit in a table.

Nothing is stored, so nothing is retried: a handler that fails fails the push. That is the trade the sync connection makes, and it is why it is not a production driver.

It borrows the worker's handler registry rather than keeping one of its own, so a job name is registered once:

w := queue.NewWorker(queue.NullQueue{}, queue.WorkerOptions{})
w.HandleFunc("invoice.send", sendInvoice)
q := queue.NewSyncQueue(w)

The worker over a NullQueue is the registry with nothing to drain, which is exactly what the sync connection is.

func NewSyncQueue

func NewSyncQueue(h Handlers) *SyncQueue

NewSyncQueue returns the queue.

func (*SyncQueue) Bulk

func (q *SyncQueue) Bulk(ctx context.Context, g auth.Grant, js []jobs.Job) error

Bulk runs the jobs in order, stopping at the first failure.

func (*SyncQueue) Clear

func (q *SyncQueue) Clear(context.Context, string) (int, error)

Clear removes nothing.

func (*SyncQueue) CreationTimeOfOldestPendingJob

func (q *SyncQueue) CreationTimeOfOldestPendingJob(context.Context, string) (time.Time, error)

CreationTimeOfOldestPendingJob is the zero time: nothing waits.

func (*SyncQueue) DelayedSize

func (q *SyncQueue) DelayedSize(context.Context, string) (int, error)

DelayedSize is zero: a sync job has already run. It answers delayedSize().

func (*SyncQueue) DeleteJob

func (q *SyncQueue) DeleteJob(context.Context, *jobs.Job) error

DeleteJob does nothing.

func (*SyncQueue) FailJob

func (q *SyncQueue) FailJob(context.Context, *jobs.Job, error) error

FailJob does nothing: the error travels back out of Push.

func (*SyncQueue) Failed

func (q *SyncQueue) Failed(context.Context, int) ([]jobs.Job, error)

Failed lists nothing: a job that failed here failed the push, and the caller already has the error.

func (*SyncQueue) GetConnectionName

func (c *SyncQueue) GetConnectionName() string

GetConnectionName is the name this queue was registered under.

It is empty on a queue built directly and never handed to a QueueManager, which is what a test does.

func (*SyncQueue) Later

func (q *SyncQueue) Later(ctx context.Context, g auth.Grant, _ time.Duration, j jobs.Job) error

Later runs the job now.

The delay is dropped rather than slept through: the sync connection is what a test and a laptop use, and neither wants the process to stop for an hour because the job was scheduled for one.

func (*SyncQueue) PendingSize

func (q *SyncQueue) PendingSize(context.Context, string) (int, error)

PendingSize is zero.

func (*SyncQueue) Pop

Pop returns nothing: a sync job has already run by the time Push returns.

func (*SyncQueue) Push

func (q *SyncQueue) Push(ctx context.Context, g auth.Grant, j jobs.Job) error

Push runs the job.

func (*SyncQueue) PushOn

func (q *SyncQueue) PushOn(ctx context.Context, g auth.Grant, queue string, j jobs.Job) error

PushOn runs the job, recording the queue it was meant for.

func (*SyncQueue) PushRaw

func (q *SyncQueue) PushRaw(ctx context.Context, g auth.Grant, name string, payload []byte, queue string) error

PushRaw runs a job whose arguments are already encoded. It answers pushRaw().

func (*SyncQueue) ReleaseJob

func (q *SyncQueue) ReleaseJob(context.Context, *jobs.Job, time.Duration) error

ReleaseJob does nothing. There is no queue to put the job back on, and a handler that releases a sync job is told so by the error it gets back from its own Push.

func (*SyncQueue) ReservedSize

func (q *SyncQueue) ReservedSize(context.Context, string) (int, error)

ReservedSize is zero. It answers reservedSize().

func (*SyncQueue) Retry

func (q *SyncQueue) Retry(context.Context, string) error

Retry has nothing to retry.

func (*SyncQueue) SetConnectionName

func (c *SyncQueue) SetConnectionName(name string)

SetConnectionName names the connection this queue answers to.

It returns nothing rather than the queue, because an embedded struct cannot return the type that embeds it.

func (*SyncQueue) Size

func (q *SyncQueue) Size(context.Context, string) (int, error)

Size is zero. Nothing is ever stored.

type TestingT

type TestingT interface {
	// Helper marks the caller as a test helper, so a failure points at the line
	// in the test rather than at the line inside this package.
	Helper()
	// Errorf records the failure and marks the test as failed.
	Errorf(format string, args ...any)
}

TestingT is the two methods the assertions here need from *testing.T.

It is declared rather than imported so that nothing in this package depends on the testing package, which would put its flags on every binary that links the queue.

type TimeoutExceeded

type TimeoutExceeded struct {
	// Job is the job that ran too long.
	Job *jobs.Job
}

TimeoutExceeded is why a job was parked: the handler ran past its timeout.

It wraps ErrMaxAttemptsExceeded, so errors.Is against that still matches: a timed-out job is one that used up a delivery without finishing, and code that treats the two the same can.

func (*TimeoutExceeded) Error

func (e *TimeoutExceeded) Error() string

Error is the message.

func (TimeoutExceeded) ForJob

ForJob returns the error for a job whose handler ran too long.

func (*TimeoutExceeded) Unwrap

func (e *TimeoutExceeded) Unwrap() error

Unwrap makes a timeout match ErrMaxAttemptsExceeded.

type Worker

type Worker struct {

	// ShouldQuit ends the loop after the job in flight. It is what a signal
	// handler sets.
	ShouldQuit atomic.Bool
	// Paused stops the worker taking new jobs without ending the loop.
	Paused atomic.Bool
	// contains filtered or unexported fields
}

Worker runs jobs off a queue.

It is also the handler registry: something has to turn the name on the wire into code, and here it is a map somebody filled in, so the set of names a binary can run is visible at the call site.

It runs in the same binary as the application, started by `aru work`, which is the same image with a different argument. Not a second artifact: one image is one thing to build, ship and roll back.

func NewWorker

func NewWorker(q Queue, opts WorkerOptions) *Worker

NewWorker returns the worker.

func (*Worker) Daemon

func (w *Worker) Daemon(ctx context.Context) (int, error)

Daemon drains the queue until it is told to stop, and returns the exit status.

The connection, the queue and the options are on the worker rather than arguments, because a worker is constructed with them.

The status is one of ExitSuccess, ExitError and ExitMemoryLimit, and it is what `aru work` returns to the shell: a supervisor reads it to tell "asked to stop" from "ran out of memory".

func (*Worker) FlushState

func (w *Worker) FlushState()

FlushState forgets that the worker was ever asked to stop or pause.

A Worker is a long-lived object, and the two flags a signal sets outlive the loop that read them: without this, a worker started a second time in the same process -- which is what a test and `queue:work --once` in a loop both do -- stops immediately because something asked the first one to.

func (*Worker) GetManager

func (w *Worker) GetManager() *QueueManager

GetManager is the manager this worker resolves connections through, or nil.

func (*Worker) Handle

func (w *Worker) Handle(name string, h Handler) *Worker

Handle registers the handler for a job name.

Registering twice panics rather than replacing. Two handlers for one name is an import nobody meant to add, and finding out at boot beats finding out from work that silently went to the wrong place.

func (*Worker) HandleFunc

func (w *Worker) HandleFunc(name string, f HandlerFunc) *Worker

HandleFunc registers a function.

func (*Worker) Handler

func (w *Worker) Handler(name string) (Handler, bool)

Handler returns the handler registered for a job name.

It is what SyncQueue borrows, so the sync connection and the worker resolve a name the same way.

func (*Worker) Kill

func (w *Worker) Kill(status int)

Kill ends the process now.

It exists for the one case Worker.Stop cannot serve: a handler that has stopped responding to its context, where returning from the loop would wait forever on a goroutine that will never finish. Everything else stops by returning.

func (*Worker) MarkJobAsFailedIfAlreadyExceedsMaxAttempts

func (w *Worker) MarkJobAsFailedIfAlreadyExceedsMaxAttempts(ctx context.Context, j *jobs.Job) error

MarkJobAsFailedIfAlreadyExceedsMaxAttempts parks a job that used up its deliveries without ever reporting an error.

The case it exists for is a job that keeps timing out: the process dies with the job in flight, the lease expires, the job comes back, and nothing ever ran the code that counts a failure. Without this check that job is delivered forever.

It returns the reason it parked the job, so the caller stops, and nil when the job may run.

func (*Worker) MemoryExceeded

func (w *Worker) MemoryExceeded(limitMB int) bool

MemoryExceeded reports whether the process is holding more than limit megabytes.

The number read is MemStats.Sys, which is what the runtime has taken from the operating system rather than what is live. A limit of zero or less is no limit.

func (*Worker) Names

func (w *Worker) Names() []string

Names returns the registered job names, for `aru work` to print at start.

func (*Worker) Options

func (w *Worker) Options() WorkerOptions

Options is the configuration this worker is running under.

The worker holds its options, and [WorkCommand] reads them so its flags can override what the application built.

func (*Worker) Pause

func (w *Worker) Pause()

Pause stops this worker taking new jobs, without ending its loop.

It is a method rather than a signal handler because signal handling belongs to the program, not to a library: `aru work` installs the handlers and calls this.

func (*Worker) Process

func (w *Worker) Process(ctx context.Context, j *jobs.Job) (err error)

Process runs one job, instrumented when there is somewhere to send it.

Instrumenting it is the point: "the nightly job is slow" is the same investigation as "the page is slow", and it deserves the same page.

The error it returns is the handler's, after the job has been settled -- so a caller that wants to know what happened can, and the daemon that does not ignores it.

Every path out of it dispatches events.JobAttempted, which is what makes that event usable for counting deliveries.

func (*Worker) Restart

func (w *Worker) Restart()

Restart asks this worker to stop after the job it is holding.

The loop returns after the batch in flight, so a deploy does not lose work.

func (*Worker) Resume

func (w *Worker) Resume()

Resume lets this worker take jobs again.

func (*Worker) RunNextJob

func (w *Worker) RunNextJob(ctx context.Context) error

RunNextJob takes one batch off the queue and runs it.

It is what `queue:work --once` calls: one pass of the loop, with none of the bookkeeping that decides whether to make another. An empty queue sleeps once and returns.

func (*Worker) SetCache

func (w *Worker) SetCache(c Cache) *Worker

SetCache gives the worker somewhere to read the restart signal and the pause flag, and returns it so the call chains.

Without one the worker never restarts on `aru queue:restart` and never observes a paused queue -- which is what a test wants and what a single-process deployment can live with.

func (*Worker) SetEvents

func (w *Worker) SetEvents(d Dispatcher) *Worker

SetEvents gives the worker somewhere to send its events, and returns it so the call chains.

Nil means the events are not built at all, which is what production looks like when nobody is listening.

func (*Worker) SetManager

func (w *Worker) SetManager(m *QueueManager)

SetManager gives the worker a manager, so it can ask whether its queue is paused.

func (*Worker) SetName

func (w *Worker) SetName(name string) *Worker

SetName names this worker, and returns it so the call chains.

The name is what PopUsing keys on.

func (*Worker) SetOptions

func (w *Worker) SetOptions(o WorkerOptions) *Worker

SetOptions replaces the configuration, and returns the worker so the call chains.

It is the setter half of Worker.Options, and it is what `queue:work --queue=reports` uses. Calling it while the loop is running is a data race: the options are read on every pass, and a command sets them before it starts the worker.

func (*Worker) Sleep

func (w *Worker) Sleep(ctx context.Context, d time.Duration)

Sleep waits for d, or until the context is cancelled.

Taking the context is the point: a worker asleep for three seconds when SIGTERM arrives would hold up the deploy for three seconds, and this one does not.

func (*Worker) Stop

func (w *Worker) Stop(status int, reason WorkerStopReason) int

Stop ends the loop and returns the exit status.

The event goes out here rather than at each call site, which is what makes "the worker stopped and this is why" one line in a log instead of six places that might have logged it.

type WorkerOptions

type WorkerOptions struct {
	// Name identifies this worker, for [PopUsing] and for the log line. Empty
	// means "default".
	Name string
	// Queue is which queue to drain. Empty means jobs.DefaultQueue.
	Queue string
	// Concurrency is how many jobs run at once. Default 4.
	//
	// One process drains a batch, and the batch is this wide.
	Concurrency int
	// Lease is how long a popped job stays invisible to other workers. It has
	// to exceed the longest handler, or a second worker picks up work still in
	// progress. Default 5 minutes.
	Lease time.Duration
	// Timeout is how long one handler may run before its context is cancelled.
	// Zero means Lease.
	//
	// It defaults to Lease so there is one number until somebody needs two, and
	// a Timeout longer than the Lease is the misconfiguration that hands a
	// running job to a second worker.
	Timeout time.Duration
	// Sleep is how long to wait before asking again when the queue was empty.
	// Default 1 second.
	Sleep time.Duration
	// Rest is how long to wait after each job, whatever the queue holds. Zero
	// means no rest. It exists to stop a fast queue from saturating a database
	// that everything else shares.
	Rest time.Duration
	// MaxTries is how many deliveries a job gets before it is parked.
	// Default 5.
	//
	// A job's own jobs.Job.MaxTries overrides it, which is what makes one
	// handler able to retry more than the rest.
	MaxTries int
	// Force runs the worker even while the application is down for
	// maintenance.
	Force bool
	// StopWhenEmpty ends the loop the first time the queue has nothing. It is
	// what a pipeline step waits on.
	StopWhenEmpty bool
	// MaxJobs stops the worker after this many jobs. Zero means no limit.
	MaxJobs int
	// MaxTime stops the worker after this long. Zero means no limit.
	MaxTime time.Duration
	// Memory stops the worker when the process is holding more than this many
	// megabytes. Zero means no limit.
	Memory int
	// Middleware wraps every job this worker runs, outermost first.
	//
	// It is the worker's list rather than the job's, because a job here is a
	// name and a payload and not a class with attributes on it: there is
	// nowhere on the wire to put "this one is rate limited". A handler that
	// needs different treatment gets its own worker, which is also how it gets
	// its own queue.
	Middleware []middleware.Middleware
	// Recorder receives each finished job, so it shows on /_arandu/debug with
	// its queries and its timeline -- exactly like a request.
	//
	// Nil means no instrumentation, and that is what production looks like: no
	// Collector is built, log.FromContext returns nil, and every Record method
	// is a no-op on a nil receiver. Zero cost, not low cost.
	//
	// Building a Collector on every job unconditionally would make production
	// pay for recording every query with its bound arguments and its caller
	// frames, with nobody able to read any of it. Pass the application's
	// recorder to turn it on.
	Recorder *log.Recorder
	// Backoff returns how long to wait before attempt n. Default is
	// exponential, capped at an hour.
	Backoff func(attempt int) time.Duration
}

WorkerOptions configures the loop.

Two of the fields are there because of how this worker runs: a Concurrency, because it drains a batch at once rather than one job per process, and a Recorder, because a job is instrumented like a request.

type WorkerStopReason

type WorkerStopReason = events.WorkerStopReason

WorkerStopReason is why a worker's loop ended.

It is a string type, so its wire value is what a process supervisor greps for.

It exists so the thing that restarts the process can tell the two kinds of exit apart: a worker that stopped because it hit its job limit is healthy and should be started again, and one that stopped because it ran out of memory is a bug somebody should see.

It is an alias for the type declared in github.com/arandu-io/hesape/queue/events, because events.WorkerStopping carries it and this package is the one that dispatches that event.

Directories

Path Synopsis
Package attributes is the per-job settings a job carries: how many tries, how long to wait between them, how long the handler may run.
Package attributes is the per-job settings a job carries: how many tries, how long to wait between them, how long the handler may run.
Package capsule is empty, and nothing is coming.
Package capsule is empty, and nothing is coming.
Package connectors opens queue connections.
Package connectors opens queue connections.
redis module
Package console is the queue's commands.
Package console is the queue's commands.
concerns
Package concerns is what more than one queue command needs.
Package concerns is what more than one queue command needs.
Package events is everything the queue announces about itself.
Package events is everything the queue announces about itself.
Package failed is where a job goes when it gives up.
Package failed is where a job goes when it gives up.
Package jobs is the job itself: what Push takes, and what the worker holds while the work runs.
Package jobs is the job itself: what Push takes, and what the worker holds while the work runs.
Package middleware wraps the handling of a job.
Package middleware wraps the handling of a job.

Jump to

Keyboard shortcuts

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