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
- Variables
- func Authorized(g auth.Grant, j Job) error
- func GrantFor(j *Job) auth.Grant
- type DatabaseJob
- type DatabaseJobRecord
- type Driver
- type FakeJob
- type Job
- func (j *Job) Backoff() time.Duration
- func (j *Job) Decode(v any) error
- func (j *Job) Delete(ctx context.Context) error
- func (j *Job) Fail(ctx context.Context, cause error) error
- func (j *Job) GetConnectionName() string
- func (j *Job) GetJobID() string
- func (j *Job) GetJobRecord() *DatabaseJobRecord
- func (j *Job) GetName() string
- func (j *Job) GetQueue() string
- func (j *Job) GetRawBody() []byte
- func (j *Job) HasFailed() bool
- func (j *Job) IsDeleted() bool
- func (j *Job) IsDeletedOrReleased() bool
- func (j *Job) IsReleased() bool
- func (j *Job) MarkAsFailed()
- func (j *Job) MaxExceptions() int
- func (j *Job) MaxTries() int
- func (j *Job) Release(ctx context.Context, delay time.Duration) error
- func (j *Job) ResolveName() string
- func (j *Job) ResolveQueuedJobClass() string
- func (j *Job) RetryUntil() time.Time
- func (j *Job) SetConnectionName(name string)
- func (j *Job) ShouldFailOnTimeout() bool
- func (j *Job) Timeout() time.Duration
- type JobName
- type SyncJob
Constants ¶
const DefaultQueue = "default"
DefaultQueue is where a job goes when nobody said otherwise.
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 ¶
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.
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.
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.
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.
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 ¶
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.
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 ¶
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.
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 ¶
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 ¶
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 ¶
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 ¶
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) Fail ¶
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 ¶
GetConnectionName is the name of the queue connection this job came off.
func (*Job) GetJobID ¶
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 ¶
GetName is the name of the queued job.
It is the name a handler is registered under: "invoice.send".
func (*Job) GetRawBody ¶
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) IsDeletedOrReleased ¶
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 ¶
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 ¶
MaxExceptions is how many failures this job gets before it is parked, or zero when only MaxTries decides.
func (*Job) MaxTries ¶
MaxTries is how many deliveries this job gets before it is parked, or zero when the worker's own limit decides.
func (*Job) Release ¶
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 ¶
ResolveName is the name to show for this job.
func (*Job) ResolveQueuedJobClass ¶
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 ¶
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 ¶
SetConnectionName names the connection this job belongs to.
The driver sets it when it pops.
func (*Job) ShouldFailOnTimeout ¶
ShouldFailOnTimeout reports whether running past the timeout parks the job rather than releasing it.
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 ¶
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.
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.