jobs

package
v0.25.1 Latest Latest
Warning

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

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

Documentation

Overview

Package jobs is the job itself: what Push takes, and what the worker holds while the work runs.

There is one Job and one Driver interface, rather than a job type per store. The three things a running job does about itself -- release, delete and fail -- differ only in which store they write to, and a type per store would be a second way to say the same thing. The queue that popped the job supplies the Driver, and DatabaseJob and SyncJob are aliases of Job.

The split between this package and github.com/arandu-io/hesape/queue is the job on one side and the thing that holds jobs on the other. It is what keeps the import graph one-directional -- queue imports jobs, jobs imports nothing of queue -- so a driver in its own module can build a Job without depending on the drivers that ship in the collection.

A field is an accessor

The values a job carries are fields, not accessors: they are columns of the record, and nothing has to be decoded to read them. The methods that remain are the ones that compute something: Job.Backoff picks the entry for this attempt, Job.ResolveName prefers the display name, Job.IsDeletedOrReleased reads two flags.

Index

Constants

View Source
const DefaultQueue = "default"

DefaultQueue is where a job goes when nobody said otherwise.

View Source
const MaxPayload = 32 << 10

MaxPayload is the largest payload a job may carry, in bytes.

One number for every driver rather than one per store. A queue whose drivers accepted different jobs would be several queues: moving an application from one connection to another would be a migration instead of a line in its wiring, and the connection a developer runs locally would take work the connection in production refuses.

So it is the narrowest thing a payload has to fit through, and there are three. The payload column is TEXT, and the narrowest engine with a connector holds 65,535 bytes in one of those -- past that the insert is either refused or, on a server that is not in strict mode, silently shortened, and a shortened payload is a handler decoding arguments the pusher never wrote. The dead letter table writes the payload a second time into a column of the same kind, so a job that gave up has to still fit at the moment its record is the only thing left of it. And a driver that runs the job in another process hands it over on an argument list, where one entry is capped far below what any store would take.

Half the column rather than all of it, because the payload is counted in bytes here and the column is counted in bytes there, and nothing should depend on the two agreeing to the last one.

It is also what a payload is for. 32 KiB of JSON is thousands of ids; more than that is a document, and a document belongs in storage with the job carrying its key.

Variables

View Source
var ErrDetached = errors.New("queue: this job was built to be pushed and never came off a queue, so there is nothing to release, delete or fail")

ErrDetached is returned when a job that was never popped is asked to settle itself.

View Source
var ErrForged = errors.New("queue: the job does not match the Grant pushing it")

ErrForged is returned when a job claims an action or a tenant the Grant pushing it does not carry.

View Source
var ErrNoName = errors.New("queue: a job with no name cannot be routed to a handler")

ErrNoName is returned when a job has no name to route by.

View Source
var ErrNoTenant = errors.New("queue: the Grant carries no tenant, and a job without one cannot be scoped")

ErrNoTenant is returned when a Grant carries no tenant.

It is an error rather than a default: a job with no tenant cannot be scoped, and everything the handler touches would read across customers.

View Source
var ErrPayloadTooLarge = errors.New("queue: the job's payload is larger than a queue stores")

ErrPayloadTooLarge is returned when a job's payload is over MaxPayload.

A sentinel rather than a message, so a caller can tell a payload it can do something about -- write the document to storage, push its key -- from a store that is down. It is not the same refusal as arguments that cannot be encoded at all: those are a defect in the value, and this is a value that is merely too big to travel as one.

Functions

func Authorized

func Authorized(g auth.Grant, j Job) error

Authorized reports whether a job may be pushed under this Grant, and whether it is small enough for a queue to store.

Every driver calls it at the top of Push, and it closes an escalation the contract otherwise allows. New builds a job from the Grant, so what it produces always matches -- but Push takes a Job, and a Job is a struct anybody can fill in:

j := jobs.Job{UUID: id, Name: "invoice.send", Action: "invoice.delete", TenantID: other}
q.Push(ctx, viewGrant, j)

The worker rebuilds the Grant from the record -- GrantFor gives auth.SystemGrant(j.Action, j.TenantID) -- so the handler would run with an action nobody authorized, in a tenant nobody authorized, and every Policy downstream would say yes because the Grant looks legitimate. The queue would be the one way past the authorization the whole collection exists to enforce.

Checked here rather than in each driver, because a driver that forgets is a driver that reopens it.

MaxPayload is checked here as well, and last. It is not about the Grant, and it is in this function for the reason the name is: this is the one call every Push in every driver already makes and already returns the error of, so a ceiling that lives here cannot be the ceiling one driver has and the next one does not. Last, because a job that is both forged and oversized is a forged job first, whatever its size.

A job over the limit is refused before anything is written. The caller holding the payload is the only one who can do something about it -- put the document in storage and push its key -- and it is holding it now; refused at the pop instead, the job has already been stored and the news arrives in a worker's log with the code that built the payload long since returned.

func GrantFor

func GrantFor(j *Job) auth.Grant

GrantFor rebuilds the Grant a job runs under.

The action and the tenant come from the record, so the worker reissues exactly what the push authorized -- not more. A worker that invented its own Grant would be a way to reach the database with permissions nobody granted.

Types

type DatabaseJob

type DatabaseJob = Job

DatabaseJob is a job that came off a DatabaseQueue.

It is an alias rather than a type of its own: there is one Job, and which store it came off is the Driver it carries. A type per store would differ only in where release and delete write, which is what the Driver already says.

type DatabaseJobRecord

type DatabaseJobRecord struct {
	// Job is the record as the worker will see it.
	*Job

	// ReservedUntil is when this row becomes visible to other workers again.
	// Zero means it is not reserved.
	ReservedUntil time.Time
}

DatabaseJobRecord is the jobs table row while the driver is holding it.

It is the bookkeeping the database driver does to a row between reading it and handing the job over: counting the delivery and stamping the reservation.

It is a separate type from Job because it holds what only the table has. ReservedUntil is not on the job: a handler has no use for it, and a field a handler can read is a field a handler will one day write.

func (*DatabaseJobRecord) Increment

func (r *DatabaseJobRecord) Increment() int

Increment counts this delivery and returns the new count.

Attempts includes the current delivery, so the first call returns 1 and the worker comparing it against MaxTries is comparing deliveries to deliveries.

func (*DatabaseJobRecord) Touch

func (r *DatabaseJobRecord) Touch(lease time.Duration) time.Time

Touch stamps the reservation forward by lease and returns the new deadline.

It pushes the reservation out so a job that is still running is not handed to a second worker. A non-positive lease leaves the deadline alone rather than expiring the reservation, because a caller that asked for no lease meant "do not change it" and unreserving a running job is the one outcome nobody wants.

type Driver

type Driver interface {
	// ReleaseJob puts the job back on its queue, eligible again after delay.
	ReleaseJob(ctx context.Context, j *Job, delay time.Duration) error
	// DeleteJob removes a finished job.
	DeleteJob(ctx context.Context, j *Job) error
	// FailJob parks the job: it stops being delivered and starts being
	// something a person can list and retry.
	FailJob(ctx context.Context, j *Job, cause error) error
}

Driver is the three things a running job asks of the queue it came off.

A queue implements it and hands itself to Popped, so the job can settle itself without the worker having to know which store it lives in.

type FakeJob

type FakeJob struct {
	// Job is the job a handler is given.
	*Job

	// ReleaseDelay is how long the job asked to wait before its next delivery.
	// It is zero on a job that was not released.
	ReleaseDelay time.Duration

	// FailedWith is the cause the job was parked with. It is nil on a job that
	// was not parked.
	FailedWith error
}

FakeJob is a job that came off nothing, and records what was done to it.

It is what a test wires when it wants to call a handler directly and then ask whether the handler released, deleted or failed its job -- without a store, a worker or a transaction.

f := jobs.NewFakeJob("invoice.send", tenant)
err := sendInvoice(ctx, auth.SystemGrant("invoice.send", tenant), f.Job)
if !f.IsReleased() || f.ReleaseDelay != time.Minute { ... }

It is both the job and the Driver behind it, which is the whole trick: the job settles itself against the fake, and the fake is what the assertions read.

func NewFakeJob

func NewFakeJob(name, tenant string) *FakeJob

NewFakeJob returns a job for name, belonging to tenant, that came off a fake queue.

Attempts is 1: the job is being handled, and the handling counts.

func (*FakeJob) DeleteJob

func (f *FakeJob) DeleteJob(context.Context, *Job) error

DeleteJob records the deletion. It answers delete() on the fake.

func (*FakeJob) FailJob

func (f *FakeJob) FailJob(_ context.Context, _ *Job, cause error) error

FailJob records the cause. It answers fail() on the fake.

func (*FakeJob) ReleaseJob

func (f *FakeJob) ReleaseJob(_ context.Context, _ *Job, delay time.Duration) error

ReleaseJob records the release. It answers release() on the fake.

type Job

type Job struct {
	// UUID is the deduplication key. It is stable across retries, which is what
	// makes a handler able to recognize work it already did. It is also the
	// job's id in the store, because the id is minted by the application rather
	// than handed back by the store (see database.NewID).
	UUID string
	// Queue separates work by urgency: a password reset email and a monthly
	// report should not wait behind each other.
	Queue string
	// Name routes the job to its handler: "invoice.send", "report.monthly".
	Name string
	// DisplayName is what a person reads on a console, a log line or the failed
	// job list, and empty means Name. The only thing that sets it is a job that
	// wraps another -- a queued closure, a queued listener -- where the name
	// that routes and the name that explains are different strings.
	DisplayName string
	// TenantID is who the work belongs to. It comes from the Grant at Push, and
	// the worker rebuilds a Grant from it -- a job with no tenant cannot be
	// scoped, and everything downstream of it reads across customers.
	TenantID string
	// Payload is the arguments, as JSON. Keep it to facts and ids: a payload
	// that says "look it up" is a payload that reads a row which has already
	// changed.
	Payload []byte
	// AuthorizedBy and Action record the Grant that pushed it, which is the
	// audit trail and what the worker reissues the work under.
	AuthorizedBy string
	Action       string
	// RunAt is when it becomes eligible. Zero means now.
	RunAt time.Time
	// Attempts counts the deliveries INCLUDING the current one: a job being
	// handled for the first time has Attempts == 1. LastError is why the most
	// recent one failed -- stored rather than logged, because the thing anyone
	// needs at 3am is "this failed twelve times with this message".
	Attempts  int
	LastError string
	// Exceptions counts the deliveries that ended in an error, which is not the
	// same number as Attempts: a middleware that releases the job -- rate
	// limited, overlapping -- spends an attempt without the handler ever having
	// thrown. It is what Attributes.MaxExceptions is compared against.
	Exceptions int
	// Attributes are the per-job settings: tries, backoff, timeout, the queue
	// and the connection. They travel with the job, so the worker reads back
	// exactly what the push declared.
	Attributes attributes.Attributes
	// contains filtered or unexported fields
}

Job is one unit of work.

It is both the record a driver stores and the handle the worker settles. A job built by New carries no driver and can only be pushed; a job returned by a queue's Pop carries the queue it came off, and Job.Release, Job.Delete and Job.Fail are how it settles itself.

A job popped from a queue is used through a pointer, because releasing it and deleting it are decisions the worker and the middleware make about the same job and a copy would lose them.

func New

func New(g auth.Grant, queue, name string, payload any) (Job, error)

New builds a job from a Grant and a payload.

This is the only constructor, so every job in the system carries a tenant, an id and the Grant that authorized it -- there is no shape of Job that skipped any of the three.

func Popped

func Popped(d Driver, connection string, j Job) *Job

Popped attaches a job to the queue it came off, and to that connection's name.

Every driver's Pop ends with it.

func Prepare

func Prepare(g auth.Grant, j Job) Job

Prepare fills in the defaults a driver would otherwise each have to remember: the queue, the run time and the tenant.

It runs after Authorized and returns the job a driver stores. Three drivers deciding independently what an empty Queue means is three chances to disagree.

func (*Job) Backoff

func (j *Job) Backoff() time.Duration

Backoff is how long to wait before the next delivery, or zero when the worker's own schedule decides.

The attempt counts from one, and the last entry of the list repeats once the list runs out.

func (*Job) Decode

func (j *Job) Decode(v any) error

Decode unmarshals the payload into v.

func (*Job) Delete

func (j *Job) Delete(ctx context.Context) error

Delete removes the job.

func (*Job) Fail

func (j *Job) Fail(ctx context.Context, cause error) error

Fail parks the job.

Parked rather than retried: the job stops being delivered and starts being something a person can see, which is what a dead letter queue is for. It marks the job deleted too, so nothing downstream settles it a second time.

func (*Job) GetConnectionName

func (j *Job) GetConnectionName() string

GetConnectionName is the name of the queue connection this job came off.

func (*Job) GetJobID

func (j *Job) GetJobID() string

GetJobID is the identifier of the job in its store.

It is the same value as UUID, because the id is minted by the application rather than handed back by the store: see database.NewID.

func (*Job) GetJobRecord

func (j *Job) GetJobRecord() *DatabaseJobRecord

GetJobRecord is the table row behind a job.

The job is the record, so this wraps it in the type that carries the reservation and returns that.

func (*Job) GetName

func (j *Job) GetName() string

GetName is the name of the queued job.

It is the name a handler is registered under: "invoice.send".

func (*Job) GetQueue

func (j *Job) GetQueue() string

GetQueue is the name of the queue the job belongs to.

func (*Job) GetRawBody

func (j *Job) GetRawBody() []byte

GetRawBody is the job's payload as it was stored.

The envelope is the record's own columns, and this is the arguments, which is the half a handler decodes.

func (*Job) HasFailed

func (j *Job) HasFailed() bool

HasFailed reports whether the job was parked.

func (*Job) IsDeleted

func (j *Job) IsDeleted() bool

IsDeleted reports whether the job was removed.

func (*Job) IsDeletedOrReleased

func (j *Job) IsDeletedOrReleased() bool

IsDeletedOrReleased reports whether the job has already been settled.

It is what the worker asks before it acts on the outcome of a handler: a middleware may have released the job before the handler ever ran, and settling it twice is how a job runs twice.

func (*Job) IsReleased

func (j *Job) IsReleased() bool

IsReleased reports whether the job was put back on its queue.

func (*Job) MarkAsFailed

func (j *Job) MarkAsFailed()

MarkAsFailed records that the job failed without telling the store.

It is the half of Job.Fail that does not park the job: a handler that has already written its own failure somewhere calls this so the worker does not release the job for another attempt.

func (*Job) MaxExceptions

func (j *Job) MaxExceptions() int

MaxExceptions is how many failures this job gets before it is parked, or zero when only MaxTries decides.

func (*Job) MaxTries

func (j *Job) MaxTries() int

MaxTries is how many deliveries this job gets before it is parked, or zero when the worker's own limit decides.

func (*Job) Release

func (j *Job) Release(ctx context.Context, delay time.Duration) error

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

A middleware that decides the work should not run now calls it -- WithoutOverlapping when another worker holds the lock, RateLimited when the budget is spent -- and so does the worker after a failure that has attempts left.

func (*Job) ResolveName

func (j *Job) ResolveName() string

ResolveName is the name to show for this job.

func (*Job) ResolveQueuedJobClass

func (j *Job) ResolveQueuedJobClass() string

ResolveQueuedJobClass is the name of the work behind a wrapper job.

A queued closure and a queued listener arrive under the wrapper's name and record what they wrap; a job that wraps nothing answers with its own name.

func (*Job) RetryUntil

func (j *Job) RetryUntil() time.Time

RetryUntil is the deadline after which this job stops being retried, whatever MaxTries says, or the zero time when there is none.

func (*Job) SetConnectionName

func (j *Job) SetConnectionName(name string)

SetConnectionName names the connection this job belongs to.

The driver sets it when it pops.

func (*Job) ShouldFailOnTimeout

func (j *Job) ShouldFailOnTimeout() bool

ShouldFailOnTimeout reports whether running past the timeout parks the job rather than releasing it.

func (*Job) Timeout

func (j *Job) Timeout() time.Duration

Timeout is how long the handler may run, or zero when the worker's lease decides.

type JobName

type JobName struct{}

JobName reads the several names a job has.

It is a struct with no fields: the three methods belong together and none of them needs a receiver, so the type is a name for the reading and `jobs.JobName{}.Resolve(...)` is all there is to it.

A job has three names because a wrapper has two:

Name          what routes it to a handler
DisplayName   what a person reads
class name    what is really behind a wrapper

For an ordinary job all three are the same string, and these three methods are what say so in one place instead of at every call site.

func (JobName) Parse

func (JobName) Parse(name string) (string, string)

Parse splits a job name into the name and the method behind it.

A name here is "invoice.send" and there is no method to call, so a name with no "@" comes back with the default method "fire".

It is kept because a record written by an older release, or by another system pushing onto the same store, can still carry the "name@method" form, and routing "SendInvoice@handle" to a handler registered as "SendInvoice" is better than failing to route it at all.

func (JobName) Resolve

func (JobName) Resolve(name string, j *Job) string

Resolve is the name to show for a job.

It prefers the job's DisplayName. A nil job, or one with no DisplayName, answers with name.

func (JobName) ResolveClassName

func (JobName) ResolveClassName(name string, j *Job) string

ResolveClassName is the name of the work behind a wrapper.

The name that routes is the name that was registered, so for a job that wraps nothing this is that name.

type SyncJob

type SyncJob = Job

SyncJob is a job that ran the moment it was pushed.

It is an alias for the same reason DatabaseJob is.

Jump to

Keyboard shortcuts

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