queue

package
v1.10.0 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 28 Imported by: 0

README

Queue Package for Nimbus

Background job processing (Laravel-inspired). Queue is core—no plugin needed. Call queue.Boot() in your app bootstrap.

Installation

Queue is initialized in bin/server.go when creating a new app. Ensure you have:

import "github.com/CodeSyncr/nimbus/queue"

queue.Boot(&queue.BootConfig{RegisterJobs: start.RegisterQueueJobs})

Configuration

Variable Description Default
QUEUE_DRIVER sync, redis, database, sqs, kafka sync
REDIS_URL Redis URL (for redis driver) redis://localhost:6379
QUEUE_REDIS_VISIBILITY_TIMEOUT_SECONDS How long a Redis job stays leased without a worker heartbeat 60
QUEUE_DB_LEASE_SECONDS How long a database job stays leased without a worker heartbeat 120
QUEUE_BOOT_STRICT Fail boot on unknown driver values false
SQS_QUEUE_URL AWS SQS queue URL —
KAFKA_BROKERS Kafka brokers (comma-separated) —
KAFKA_TOPIC Kafka topic nimbus-queue
KAFKA_GROUP_ID Consumer group ID nimbus-queue

Drivers:

  • sync — Runs jobs immediately (no worker). Useful for dev.
  • redis — Redis lists. Requires REDIS_URL.
  • database — GORM. Uses database.Get().
  • sqs — AWS SQS.
  • kafka — Apache Kafka.

Defining jobs

Implement the queue.Job interface:

package jobs

import (
    "context"
    "github.com/CodeSyncr/nimbus/queue"
)

type SendEmail struct {
    UserID  int
    Subject string
}

func (j *SendEmail) Handle(ctx context.Context) error {
    // Send email...
    return nil
}

Optional Failed for cleanup when job fails permanently:

func (j *SendEmail) Failed(ctx context.Context, err error) {
    log.Printf("SendEmail failed for user %d: %v", j.UserID, err)
}

Dispatching jobs

import "github.com/CodeSyncr/nimbus/queue"

queue.Dispatch(&jobs.SendEmail{UserID: 12, Subject: "Welcome"}).Dispatch(ctx)

// Delayed
queue.Dispatch(&jobs.SendEmail{...}).Delay(5 * time.Minute).Dispatch(ctx)

// Specific queue
queue.Dispatch(&jobs.Report{}).OnQueue("reports").Dispatch(ctx)

Registering jobs

In start/jobs.go (or equivalent):

package start

import (
    "github.com/CodeSyncr/nimbus/queue"
    "myapp/jobs"
)

func RegisterQueueJobs() {
    queue.Register(&jobs.SendEmail{})
    queue.Register(&jobs.ProcessVideo{})
}

Running the worker

nimbus queue:work

Or from your app:

queue.RunWorker(ctx, "default")

Rate limiting

Pass RateLimitPerSec and RateLimitBurst in BootConfig to throttle job processing. The limit is per queue. With the Redis driver it is kept in Redis, so it holds across all workers and instances; with other drivers each worker process gets the full limit.

Delivery guarantees

  • With redis and database, each job is handed to one worker at a time. While a job runs, its worker renews the job's lease every third of the lease period, so long jobs are not picked up twice.
  • If a worker dies, its job returns to the queue after one lease period and the lost run counts as an attempt. A job that keeps killing its worker fails for good once it uses up its retries, instead of looping forever.
  • Delivery is still at least once (a network partition can outlast a lease), so keep handlers idempotent.
  • Return queue.Release(delay) from Handle to put a job back without counting an attempt.
  • Database driver: finished jobs are deleted, and retries reuse their job's row.

Batches, chains, unique jobs, overlap

All four work across workers and instances with the redis and database drivers: queue.Boot keeps their locks and progress in the same backend.

// At boot, in every process that runs workers:
queue.RegisterBatch("import-users", queue.BatchCallbacks{
    Then:    func(ctx context.Context, b *queue.Batch) { /* all succeeded */ },
    Catch:   func(ctx context.Context, b *queue.Batch, err error) { /* one job failed */ },
    Finally: func(ctx context.Context, b *queue.Batch) { /* all done */ },
})

// Anywhere: queues every job and returns. Progress: queue.FindBatch(ctx, b.ID)
b := queue.NewBatch(jobs...).Named("import-users")
err := b.Dispatch(ctx)

// Each job is queued by the worker that finished the previous one.
queue.NewChain(&Fetch{}, &Transform{}, &Load{}).DispatchAsync(ctx)

// Skips the dispatch while an identical job is pending (lock released when it finishes).
queue.DispatchUnique(ctx, &RebuildIndex{TenantID: 7}) // implements UniqueID() / UniqueFor()

// Only one at a time per key; a blocked job is put back and retried shortly.
queue.Dispatch(queue.NewWithoutOverlapping(&SyncAccount{ID: 7}, "account-7")).Dispatch(ctx)

Jobs can also implement OverlapKey() string (queue.NonOverlapping) instead of using the wrapper. With the sync driver (or no manager), Batch.Dispatch runs the jobs concurrently in-process; Batch.Run always does.

Use this section as a starting baseline for reliable queue processing.

1) Use a durable driver in production
  • Prefer redis or database (avoid sync in production).
  • Run multiple workers (nimbus queue:work) behind a process supervisor.
  • Set QUEUE_BOOT_STRICT=true to fail fast on invalid queue driver config.
2) Tune lease / visibility timeouts
  • Redis (QUEUE_REDIS_VISIBILITY_TIMEOUT_SECONDS): start at 60.
  • Database (QUEUE_DB_LEASE_SECONDS): start at 120.
  • Workers heartbeat running jobs, so the timeout does not need to cover job runtime: it is how long a crashed worker's job waits before another worker takes it.
  • Too low: a worker stalled longer than the lease (long GC pause, network blip) can lose its job to another worker.
  • Too high: slow recovery when workers crash.

You can also set these in code:

queue.Boot(&queue.BootConfig{
    Driver:                 "redis",
    RedisURL:               "redis://localhost:6379",
    RedisVisibilityTimeout: 60 * time.Second,
    DatabaseLeaseDuration:  120 * time.Second,
    RegisterJobs:           start.RegisterQueueJobs,
})
3) Retry policy
  • Default retry backoff is exponential with jitter.
  • Keep retries bounded (Retries(n) per job).
  • Ensure job handlers are idempotent (safe to run more than once).
4) Observe the right signals

Nimbus now exports queue counters (via Horizon metrics):

  • nimbus_queue_jobs_dispatched_total
  • nimbus_queue_jobs_processed_total
  • nimbus_queue_jobs_failed_total
  • nimbus_queue_jobs_retried_total
  • nimbus_queue_jobs_reclaimed_total

Prometheus-formatted endpoint (when Horizon is enabled):

  • GET /horizon/api/metrics/prometheus
5) Suggested alert thresholds (starting point)
  • Retry spike: retried/processed ratio > 5% for 5-10 minutes.
  • Reclaim activity: reclaimed > 0 sustained for 10+ minutes.
  • Failure rate: failed/processed ratio > 1-2% for 5+ minutes.
  • Backlog growth: queue length rising continuously without recovery.

Tune thresholds per workload after collecting a week of baseline data.

Documentation

Index

Constants

This section is empty.

Variables

View Source
var ErrBatchNotFound = errors.New("queue: batch not found")

ErrBatchNotFound is returned when a batch ID is unknown (or expired).

View Source
var ErrLeaseLost = errors.New("queue: job lease lost")

ErrLeaseLost is returned by LeaseExtender.ExtendLease when the job's lease already expired and the job may have been handed to another worker.

Functions

func AcquireUniqueLock

func AcquireUniqueLock(job UniqueJob) bool

AcquireUniqueLock takes the unique lock for a job until UniqueFor passes or ReleaseUniqueLock is called.

func AddObserver

func AddObserver(o Observer)

AddObserver appends an observer (e.g. Telescope alongside Horizon).

func AfterBatch

func AfterBatch(fn func(context.Context, *Batch))

AfterBatch registers a callback invoked after a batch finishes all jobs (after Then / Finally hooks on the batch itself). Useful for observability (e.g. Telescope). Multiple subscribers are allowed.

func DispatchUnique

func DispatchUnique(ctx context.Context, job UniqueJob) error

DispatchUnique dispatches a job unless one with the same UniqueID is already pending. The lock lives in the configured Locker (Redis or the database with those drivers), so it holds across instances.

func IsUniqueLocked

func IsUniqueLocked(job UniqueJob) bool

IsUniqueLocked reports whether a job with the same UniqueID is pending.

func Register

func Register(job Job)

Register registers a job type with the global manager. Call at startup, typically from a central RegisterQueueJobs function in your application.

queue.Register(&jobs.SendEmail{})

func RegisterBatch added in v1.9.1

func RegisterBatch(name string, cb BatchCallbacks)

RegisterBatch registers callbacks for batches with this name. Call it at boot in every process that runs queue workers.

func RegisterChain added in v1.9.1

func RegisterChain(name string, onFailure func(ctx context.Context, failedJob Job, err error))

RegisterChain registers the failure callback for async chains with this name. Call it at boot in every worker process.

func Release added in v1.9.1

func Release(delay time.Duration) error

Release returns an error that puts the job back on the queue after delay without counting it as a failed attempt.

func ReleaseUniqueLock

func ReleaseUniqueLock(job UniqueJob)

ReleaseUniqueLock releases the unique lock for a job.

func RunWorker

func RunWorker(ctx context.Context, queueName string)

RunWorker runs the queue worker loop. Call from queue:work command.

func SetBatchStore added in v1.9.1

func SetBatchStore(s BatchStore)

SetBatchStore sets where durable batches keep progress. Boot calls it for the redis and database drivers.

func SetFailedJobStore

func SetFailedJobStore(s FailedJobStore)

SetFailedJobStore sets the global failed job store (e.g. from Horizon plugin).

func SetGlobal

func SetGlobal(m *Manager)

SetGlobal sets the global manager.

func SetLocker added in v1.9.1

func SetLocker(l Locker)

SetLocker sets the Locker used for unique jobs and WithoutOverlapping. Boot calls it for the redis and database drivers.

func SetObserver

func SetObserver(o Observer)

SetObserver replaces the observer list with a single observer (Horizon).

Types

type Adapter

type Adapter interface {
	// Push adds a job to the queue. delay is 0 for immediate.
	Push(ctx context.Context, payload *JobPayload) error

	// Pop blocks until a job is available or ctx is done. Returns nil when done.
	Pop(ctx context.Context, queue string) (*JobPayload, error)

	// Len returns approximate number of pending jobs (best-effort).
	Len(ctx context.Context, queue string) (int, error)
}

Adapter enqueues and dequeues jobs for processing.

type Batch

type Batch struct {
	ID string
	// contains filtered or unexported fields
}

Batch dispatches a group of jobs and runs callbacks when they finish.

With a real queue driver (redis, database, ...), Dispatch queues every job and returns; workers record progress in the BatchStore and whichever one finishes the last job runs the callbacks. Callbacks are found by batch name, so give the batch a name and register its callbacks at boot in every worker process with RegisterBatch. With the sync driver (or no manager), Dispatch runs the jobs concurrently in-process, like Run.

func NewBatch

func NewBatch(jobs ...Job) *Batch

NewBatch creates a new job batch.

func (*Batch) Catch

func (b *Batch) Catch(fn func(ctx context.Context, batch *Batch, err error)) *Batch

Catch sets a callback for when any job in the batch fails.

func (*Batch) Dispatch

func (b *Batch) Dispatch(ctx context.Context) error

Dispatch queues the batch's jobs (see Batch). It returns once every job is queued; use FindBatch or Refresh to follow progress.

func (*Batch) Errors

func (b *Batch) Errors() []error

Errors returns all errors from failed jobs.

func (*Batch) FailedJobs

func (b *Batch) FailedJobs() int32

FailedJobs returns the number of failed jobs.

func (*Batch) Finally

func (b *Batch) Finally(fn func(ctx context.Context, batch *Batch)) *Batch

Finally sets a callback that runs after all jobs complete (success or failure).

func (*Batch) Finished

func (b *Batch) Finished() bool

Finished returns true when all jobs have completed.

func (*Batch) HasFailures

func (b *Batch) HasFailures() bool

HasFailures returns true if any jobs failed.

func (*Batch) Name added in v1.9.1

func (b *Batch) Name() string

Name returns the batch name.

func (*Batch) Named added in v1.9.1

func (b *Batch) Named(name string) *Batch

Named names the batch. Workers look callbacks up by this name.

func (*Batch) OnQueue

func (b *Batch) OnQueue(name string) *Batch

OnQueue sets the queue for all jobs in the batch.

func (*Batch) PendingJobs

func (b *Batch) PendingJobs() int32

PendingJobs returns the number of jobs still pending.

func (*Batch) Refresh added in v1.9.1

func (b *Batch) Refresh(ctx context.Context) error

Refresh reloads a durable batch's counters from the BatchStore.

func (*Batch) Run added in v1.9.1

func (b *Batch) Run(ctx context.Context) error

Run runs all jobs in the batch concurrently in this process and blocks until they finish. Nothing is persisted.

func (*Batch) Then

func (b *Batch) Then(fn func(ctx context.Context, batch *Batch)) *Batch

Then sets a callback for when all non-failed jobs complete.

func (*Batch) TotalJobs

func (b *Batch) TotalJobs() int32

TotalJobs returns the total number of jobs in the batch.

type BatchCallbacks added in v1.9.1

type BatchCallbacks struct {
	// Then runs once when every job finished and none failed.
	Then func(ctx context.Context, batch *Batch)
	// Catch runs for each job that fails for good.
	Catch func(ctx context.Context, batch *Batch, err error)
	// Finally runs once when every job finished, failed or not.
	Finally func(ctx context.Context, batch *Batch)
}

BatchCallbacks are the callbacks of a named batch.

type BatchState added in v1.9.1

type BatchState struct {
	ID         string    `json:"id"`
	Name       string    `json:"name"`
	Total      int       `json:"total"`
	Pending    int       `json:"pending"`
	Failed     int       `json:"failed"`
	Errors     []string  `json:"errors,omitempty"`
	CreatedAt  time.Time `json:"created_at"`
	FinishedAt time.Time `json:"finished_at,omitzero"`
}

BatchState is a batch's stored progress.

func FindBatch added in v1.9.1

func FindBatch(ctx context.Context, id string) (*BatchState, error)

FindBatch looks up a dispatched batch's progress by ID.

type BatchStore added in v1.9.1

type BatchStore interface {
	Create(ctx context.Context, state BatchState) error
	// Record counts jobID as done (failed when jobErr is non-nil). Recording
	// the same job twice is a no-op.
	Record(ctx context.Context, batchID, jobID string, jobErr error) (BatchUpdate, error)
	Find(ctx context.Context, batchID string) (*BatchState, error)
}

BatchStore persists batch progress.

func GetBatchStore added in v1.9.1

func GetBatchStore() BatchStore

GetBatchStore returns the store durable batches use.

type BatchUpdate added in v1.9.1

type BatchUpdate struct {
	State BatchState
	// Counted is false when this job's outcome was already recorded (the
	// job was delivered twice), so callbacks must not run again.
	Counted bool
	// Finished is true for exactly one caller: the one that recorded the
	// batch's last outstanding job.
	Finished bool
}

BatchUpdate is the result of recording one job's outcome.

type BootConfig

type BootConfig struct {
	Driver   string // sync, redis, database, sqs, kafka
	RedisURL string
	// RedisVisibilityTimeout controls Redis in-flight lease timeout.
	RedisVisibilityTimeout time.Duration
	// DatabaseLeaseDuration controls how long processing DB jobs are leased before reclaim.
	DatabaseLeaseDuration time.Duration
	SQSQueueURL           string
	KafkaBrokers          string
	KafkaTopic            string
	KafkaGroupID          string
	RateLimitPerSec       float64
	RateLimitBurst        int
	Strict                bool
	RegisterJobs          func()
}

BootConfig configures queue boot. Pass nil for env-based config.

type Chain

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

Chain dispatches a sequence of jobs one after another. If any job in the chain fails (exhausts retries), the rest are skipped and the optional OnFailure callback is called.

func NewChain

func NewChain(jobs ...Job) *Chain

NewChain creates a new job chain.

func (*Chain) Dispatch

func (c *Chain) Dispatch(ctx context.Context) error

Dispatch runs the chain sequentially in the calling goroutine.

func (*Chain) DispatchAsync

func (c *Chain) DispatchAsync(ctx context.Context) error

DispatchAsync queues the first job; each job's worker queues the next one when it succeeds. Every job must be registered with queue.Register.

OnFailure runs on the worker that saw the failure, so register it there with RegisterChain under the chain's name.

func (*Chain) Named added in v1.9.1

func (c *Chain) Named(name string) *Chain

Named names the chain. DispatchAsync needs a name when OnFailure is set, so workers can find the callback (see RegisterChain).

func (*Chain) OnFailure

func (c *Chain) OnFailure(fn func(ctx context.Context, failedJob Job, err error)) *Chain

OnFailure sets a callback when a chain job fails.

func (*Chain) OnQueue

func (c *Chain) OnQueue(name string) *Chain

OnQueue sets the queue for all jobs in the chain.

type CompletableAdapter

type CompletableAdapter interface {
	Adapter
	Complete(ctx context.Context, payload *JobPayload) error
}

CompletableAdapter optionally deletes/acks a message after successful processing (e.g. SQS).

type DatabaseAdapter

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

DatabaseAdapter uses a SQL database for job storage.

func NewDatabaseAdapter

func NewDatabaseAdapter(db *lucid.DB) *DatabaseAdapter

NewDatabaseAdapter creates a database adapter.

func (*DatabaseAdapter) Complete

func (d *DatabaseAdapter) Complete(ctx context.Context, payload *JobPayload) error

Complete deletes a finished job's row, unless the job was meanwhile requeued (retry/release) or handed to another worker.

func (*DatabaseAdapter) EnsureTable

func (d *DatabaseAdapter) EnsureTable(ctx context.Context) error

EnsureTable creates or migrates the queue_jobs table and drops rows left in the "done" state by older versions.

func (*DatabaseAdapter) ExtendLease added in v1.9.1

func (d *DatabaseAdapter) ExtendLease(ctx context.Context, payload *JobPayload) error

ExtendLease implements LeaseExtender.

func (*DatabaseAdapter) LeaseDuration added in v1.9.1

func (d *DatabaseAdapter) LeaseDuration() time.Duration

LeaseDuration implements LeaseExtender.

func (*DatabaseAdapter) Len

func (d *DatabaseAdapter) Len(ctx context.Context, queue string) (int, error)

Len returns the number of pending jobs.

func (*DatabaseAdapter) Pop

func (d *DatabaseAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)

Pop blocks until a job is available. Uses SELECT ... FOR UPDATE SKIP LOCKED on Postgres/MySQL; elsewhere a guarded UPDATE decides who gets the job.

func (*DatabaseAdapter) Push

func (d *DatabaseAdapter) Push(ctx context.Context, payload *JobPayload) error

Push adds a job to the queue. Pushing a job ID that already has a row (a retry or release) resets that row to pending.

func (*DatabaseAdapter) SetLeaseDuration

func (d *DatabaseAdapter) SetLeaseDuration(v time.Duration)

SetLeaseDuration sets how long a processing job may go without a heartbeat before it is handed to another worker. Workers run by the Manager extend the lease while the job runs.

func (*DatabaseAdapter) SetPollInterval added in v1.9.1

func (d *DatabaseAdapter) SetPollInterval(v time.Duration)

SetPollInterval sets how often an idle worker checks for new jobs.

type DatabaseBatchStore added in v1.9.1

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

DatabaseBatchStore keeps batch progress in SQL tables.

func NewDatabaseBatchStore added in v1.9.1

func NewDatabaseBatchStore(db *lucid.DB) *DatabaseBatchStore

NewDatabaseBatchStore returns a database-backed BatchStore. Call EnsureTable once.

func (*DatabaseBatchStore) Create added in v1.9.1

func (s *DatabaseBatchStore) Create(ctx context.Context, state BatchState) error

func (*DatabaseBatchStore) EnsureTable added in v1.9.1

func (s *DatabaseBatchStore) EnsureTable(ctx context.Context) error

EnsureTable creates the queue_batches and queue_batch_jobs tables.

func (*DatabaseBatchStore) Find added in v1.9.1

func (s *DatabaseBatchStore) Find(ctx context.Context, batchID string) (*BatchState, error)

func (*DatabaseBatchStore) Record added in v1.9.1

func (s *DatabaseBatchStore) Record(ctx context.Context, batchID, jobID string, jobErr error) (BatchUpdate, error)

type DatabaseLocker added in v1.9.1

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

DatabaseLocker stores locks as rows keyed by lock name, so the primary key decides who wins.

func NewDatabaseLocker added in v1.9.1

func NewDatabaseLocker(db *lucid.DB) *DatabaseLocker

NewDatabaseLocker returns a database-backed Locker. Call EnsureTable once.

func (*DatabaseLocker) Acquire added in v1.9.1

func (d *DatabaseLocker) Acquire(ctx context.Context, key string, ttl time.Duration) (string, bool, error)

func (*DatabaseLocker) EnsureTable added in v1.9.1

func (d *DatabaseLocker) EnsureTable(ctx context.Context) error

EnsureTable creates the queue_locks table if it does not exist.

func (*DatabaseLocker) Release added in v1.9.1

func (d *DatabaseLocker) Release(ctx context.Context, key, token string) error

type DispatchBuilder

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

DispatchBuilder allows chaining dispatch options.

func Dispatch

func Dispatch(job Job) *DispatchBuilder

Dispatch enqueues a job using the global manager. It is the primary entry point for application code to push work onto a queue.

The returned DispatchBuilder can be used to set options before the job is actually queued:

queue.Dispatch(&jobs.SendEmail{UserID: 12}).
    Delay(5 * time.Minute).
    Dispatch(ctx)

queue.Dispatch(&jobs.Report{}).
    OnQueue("reports").
    Retries(5).
    Dispatch(ctx)

If no global manager has been configured via queue.SetGlobal (usually done by queue.Boot), Dispatch returns a no-op builder so that calls are safe in tests or environments where the queue subsystem is disabled.

func (*DispatchBuilder) Delay

Delay sets the delay before the job runs.

func (*DispatchBuilder) Dispatch

func (b *DispatchBuilder) Dispatch(ctx context.Context) error

Dispatch executes the dispatch.

func (*DispatchBuilder) OnQueue

func (b *DispatchBuilder) OnQueue(name string) *DispatchBuilder

OnQueue sets the queue name.

func (*DispatchBuilder) Priority

func (b *DispatchBuilder) Priority(n int) *DispatchBuilder

Priority sets job priority (1=highest, 10=lowest).

func (*DispatchBuilder) Retries

func (b *DispatchBuilder) Retries(n int) *DispatchBuilder

Retries sets max retry attempts.

type FailedJob

type FailedJob interface {
	Job
	Failed(ctx context.Context, err error)
}

FailedJob is optionally implemented for cleanup when job fails permanently.

type FailedJobRecord

type FailedJobRecord struct {
	ID         string    `json:"id"`
	UUID       string    `json:"uuid"`
	Queue      string    `json:"queue"`
	JobName    string    `json:"job_name"`
	Payload    []byte    `json:"payload"`
	Exception  string    `json:"exception"`
	FailedAt   time.Time `json:"failed_at"`
	Attempts   int       `json:"attempts"`
	MaxRetries int       `json:"max_retries"`
}

FailedJobRecord is a single failed job entry for listing and retry.

type FailedJobStore

type FailedJobStore interface {
	// Push adds a failed job to the store.
	Push(ctx context.Context, payload *JobPayload, errMsg string) error
	// List returns all failed job records (e.g. for dashboard).
	List(ctx context.Context) ([]FailedJobRecord, error)
	// Get returns one record by ID/UUID.
	Get(ctx context.Context, id string) (*FailedJobRecord, error)
	// Forget removes a single failed job.
	Forget(ctx context.Context, id string) error
	// ForgetAll removes all failed jobs.
	ForgetAll(ctx context.Context) error
	// Retry re-enqueues the job and removes it from failed store.
	Retry(ctx context.Context, id string, enqueue func(ctx context.Context, payload *JobPayload) error) error
}

FailedJobStore persists failed jobs for Horizon dashboard (list, forget, retry).

func GetFailedJobStore

func GetFailedJobStore() FailedJobStore

GetFailedJobStore returns the global failed job store.

type Job

type Job interface {
	Handle(ctx context.Context) error
}

Job is the interface for queue jobs.

type JobFunc

type JobFunc func(ctx context.Context) error

JobFunc adapts a function to Job.

func (JobFunc) Handle

func (f JobFunc) Handle(ctx context.Context) error

type JobLimiter added in v1.9.1

type JobLimiter interface {
	Wait(ctx context.Context, queue string) error
}

JobLimiter blocks until a job may be taken from queue.

type JobPayload

type JobPayload struct {
	ID         string                 `json:"id"`
	JobName    string                 `json:"job"`
	Queue      string                 `json:"queue"`
	Payload    []byte                 `json:"payload"`
	Attempts   int                    `json:"attempts"`
	MaxRetries int                    `json:"max_retries"`
	Delay      time.Duration          `json:"delay"`
	RunAt      time.Time              `json:"run_at"`
	Meta       map[string]interface{} `json:"meta,omitempty"`
}

JobPayload is the serialized form of a job for storage.

type KafkaAdapter

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

KafkaAdapter uses Kafka for job storage.

func NewKafkaAdapter

func NewKafkaAdapter(cfg KafkaConfig) *KafkaAdapter

NewKafkaAdapter creates a Kafka adapter.

func (*KafkaAdapter) Close

func (k *KafkaAdapter) Close() error

Close closes the writer and reader.

func (*KafkaAdapter) Len

func (k *KafkaAdapter) Len(ctx context.Context, queue string) (int, error)

Len returns 0 (Kafka doesn't provide simple count).

func (*KafkaAdapter) Pop

func (k *KafkaAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)

Pop blocks until a job is available.

func (*KafkaAdapter) Push

func (k *KafkaAdapter) Push(ctx context.Context, payload *JobPayload) error

Push adds a job to the queue.

type KafkaConfig

type KafkaConfig struct {
	Brokers  []string // e.g. []string{"localhost:9092"}
	Topic    string   // queue topic
	GroupID  string   // consumer group for workers
	MinBytes int      // min bytes to fetch (default 1)
	MaxBytes int      // max bytes to fetch (default 1e6)
}

KafkaConfig holds Kafka adapter configuration.

type LeaseExtender added in v1.9.1

type LeaseExtender interface {
	LeaseDuration() time.Duration
	ExtendLease(ctx context.Context, payload *JobPayload) error
}

LeaseExtender is implemented by adapters whose deliveries are leased for a limited time (Redis visibility timeout, database lease). While a job runs, the Manager calls ExtendLease every LeaseDuration()/3 so long jobs are not handed to a second worker, while a crashed worker's jobs still come back after one lease period.

type LocalJobLimiter added in v1.9.1

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

LocalJobLimiter rate-limits each queue within this process.

func NewLocalJobLimiter added in v1.9.1

func NewLocalJobLimiter(limitPerSec float64, burst int) *LocalJobLimiter

NewLocalJobLimiter returns an in-process JobLimiter.

func (*LocalJobLimiter) Wait added in v1.9.1

func (l *LocalJobLimiter) Wait(ctx context.Context, queue string) error

type Locker added in v1.9.1

type Locker interface {
	// Acquire takes key for ttl. ok is false when someone else holds it.
	// The returned token releases the lock.
	Acquire(ctx context.Context, key string, ttl time.Duration) (token string, ok bool, err error)
	// Release drops key if token still holds it. An empty token drops the
	// lock whoever holds it.
	Release(ctx context.Context, key, token string) error
}

Locker hands out expiring exclusive locks.

func GetLocker added in v1.9.1

func GetLocker() Locker

GetLocker returns the Locker used for unique jobs and WithoutOverlapping.

type Manager

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

Manager manages adapters and job dispatch.

func Boot

func Boot(cfg *BootConfig) *Manager

Boot initializes the queue manager from config/env and sets it globally.

func BootWithError

func BootWithError(cfg *BootConfig) (*Manager, error)

BootWithError initializes the queue manager from config/env and returns an explicit error when configuration/adapter setup fails.

func GetGlobal

func GetGlobal() *Manager

GetGlobal returns the global manager.

func NewManager

func NewManager(adapter Adapter) *Manager

NewManager creates a manager with the given adapter. Pass nil to use SyncAdapter.

func (*Manager) Adapter

func (m *Manager) Adapter() Adapter

Adapter returns the underlying queue adapter (for Horizon retry, etc.).

func (*Manager) Dispatch

func (m *Manager) Dispatch(job Job) *DispatchBuilder

Dispatch enqueues a job. Returns a DispatchBuilder for options.

func (*Manager) Process

func (m *Manager) Process(ctx context.Context, queue string) error

Process pops a job from the adapter, deserializes, and runs it.

func (*Manager) Register

func (m *Manager) Register(job Job)

Register registers a job type for deserialization. Call with a zero-value instance.

queue.Register(&jobs.SendEmail{})

func (*Manager) RegisterFunc

func (m *Manager) RegisterFunc(name string, fn func() Job)

RegisterFunc registers a job by name with a constructor.

type MemoryBatchStore added in v1.9.1

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

MemoryBatchStore keeps batches in this process. Batches finished more than an hour ago are dropped.

func NewMemoryBatchStore added in v1.9.1

func NewMemoryBatchStore() *MemoryBatchStore

NewMemoryBatchStore returns a process-local BatchStore.

func (*MemoryBatchStore) Create added in v1.9.1

func (s *MemoryBatchStore) Create(_ context.Context, state BatchState) error

func (*MemoryBatchStore) Find added in v1.9.1

func (s *MemoryBatchStore) Find(_ context.Context, batchID string) (*BatchState, error)

func (*MemoryBatchStore) Record added in v1.9.1

func (s *MemoryBatchStore) Record(_ context.Context, batchID, jobID string, jobErr error) (BatchUpdate, error)

type MemoryLocker added in v1.9.1

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

MemoryLocker keeps locks in this process only.

func NewMemoryLocker added in v1.9.1

func NewMemoryLocker() *MemoryLocker

NewMemoryLocker returns a process-local Locker.

func (*MemoryLocker) Acquire added in v1.9.1

func (m *MemoryLocker) Acquire(_ context.Context, key string, ttl time.Duration) (string, bool, error)

func (*MemoryLocker) Release added in v1.9.1

func (m *MemoryLocker) Release(_ context.Context, key, token string) error

type NonOverlapping added in v1.9.1

type NonOverlapping interface {
	Job
	OverlapKey() string
}

NonOverlapping is implemented by jobs that must not run at the same time as another job with the same OverlapKey. When a worker picks one up while another is running, it puts it back on the queue (without counting an attempt) and tries again later. The lock lives in the configured Locker, so it holds across workers and instances.

type Observer

type Observer interface {
	JobDispatched(payload *JobPayload)
	JobProcessed(payload *JobPayload, err error)
}

Observer can be used to observe queue lifecycle events (for dashboards like Horizon). It is optional and only called when set.

type ObserverV2

type ObserverV2 interface {
	JobProcessedV2(payload *JobPayload, duration time.Duration, err error)
}

ObserverV2 is an optional extension interface for richer job lifecycle metadata. Observers can implement this in addition to Observer.

type OverlapOptions added in v1.9.1

type OverlapOptions interface {
	// OverlapExpiresAfter bounds how long the lock is held if the worker
	// dies mid-job. Default 10 minutes.
	OverlapExpiresAfter() time.Duration
	// OverlapReleaseAfter is how long a blocked job waits before it is
	// tried again. Default 3 seconds.
	OverlapReleaseAfter() time.Duration
}

OverlapOptions optionally tunes a NonOverlapping job.

type Queue

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

Queue is a legacy in-memory worker pool. Prefer queue.Dispatch() with a configured Manager (Redis/Database) for production.

func NewQueue

func NewQueue(workers int) *Queue

NewQueue creates a queue with n workers. Start with Run().

func (*Queue) Push

func (q *Queue) Push(job Job)

Push enqueues a job (non-blocking if buffer full; can backpressure).

func (*Queue) Run

func (q *Queue) Run(ctx context.Context)

Run starts the worker pool. Call in a goroutine or block.

func (*Queue) Stop

func (q *Queue) Stop()

Stop stops the queue.

type QueueBatch added in v1.9.1

type QueueBatch struct {
	ID         string `gorm:"primaryKey;size:36"`
	Name       string `gorm:"size:191"`
	Total      int    `gorm:"not null"`
	Pending    int    `gorm:"not null"`
	Failed     int    `gorm:"not null;default:0"`
	Errors     string `gorm:"type:text"` // JSON array
	CreatedAt  time.Time
	FinishedAt *time.Time `gorm:"index"`
}

QueueBatch is the database model for DatabaseBatchStore.

func (QueueBatch) TableName added in v1.9.1

func (QueueBatch) TableName() string

type QueueBatchJob added in v1.9.1

type QueueBatchJob struct {
	BatchID string `gorm:"primaryKey;size:36"`
	JobID   string `gorm:"primaryKey;size:64"`
}

QueueBatchJob records which jobs a batch has already counted.

func (QueueBatchJob) TableName added in v1.9.1

func (QueueBatchJob) TableName() string

type QueueJob

type QueueJob struct {
	ID         string    `gorm:"primaryKey;size:36"`
	Queue      string    `gorm:"index;index:idx_queue_jobs_ready,priority:1;index:idx_queue_jobs_lease,priority:1;size:64;not null"`
	Payload    []byte    `gorm:"type:text;not null"` // JSON of JobPayload
	RunAt      time.Time `gorm:"index;index:idx_queue_jobs_ready,priority:3;not null"`
	Status     string    `gorm:"size:16;default:pending;index:idx_queue_jobs_ready,priority:2;index:idx_queue_jobs_lease,priority:2"` // pending, processing
	ClaimToken string    `gorm:"size:32;not null;default:''"`
	Reclaims   int       `gorm:"not null;default:0"` // deliveries lost to dead workers
	CreatedAt  time.Time
	UpdatedAt  time.Time `gorm:"index:idx_queue_jobs_lease,priority:3"`
}

QueueJob is the database model for jobs. Rows are deleted once a job finishes; a retry reuses its job's row.

func (QueueJob) TableName

func (QueueJob) TableName() string

type QueueLock added in v1.9.1

type QueueLock struct {
	Key       string    `gorm:"column:lock_key;primaryKey;size:191"`
	Token     string    `gorm:"size:32;not null"`
	ExpiresAt time.Time `gorm:"index;not null"`
}

QueueLock is the database model for DatabaseLocker.

func (QueueLock) TableName added in v1.9.1

func (QueueLock) TableName() string

type RateLimitAdapter

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

RateLimitAdapter wraps an adapter with rate limiting.

func NewRateLimitAdapter

func NewRateLimitAdapter(inner Adapter, limitPerSec float64, burst int) *RateLimitAdapter

NewRateLimitAdapter wraps an adapter with an in-process rate limiter. limitPerSec: max jobs per second per queue (e.g. 10) burst: max burst size (e.g. 20)

func NewRateLimitAdapterWith added in v1.9.1

func NewRateLimitAdapterWith(inner Adapter, limiter JobLimiter, limitPerSec float64) *RateLimitAdapter

NewRateLimitAdapterWith wraps an adapter with the given limiter, e.g. a RedisJobLimiter shared by every worker.

func (*RateLimitAdapter) Complete

func (r *RateLimitAdapter) Complete(ctx context.Context, payload *JobPayload) error

Complete delegates if inner supports it.

func (*RateLimitAdapter) ExtendLease added in v1.9.1

func (r *RateLimitAdapter) ExtendLease(ctx context.Context, payload *JobPayload) error

ExtendLease delegates if inner leases jobs.

func (*RateLimitAdapter) Inner added in v1.9.1

func (r *RateLimitAdapter) Inner() Adapter

Inner returns the wrapped adapter.

func (*RateLimitAdapter) LeaseDuration added in v1.9.1

func (r *RateLimitAdapter) LeaseDuration() time.Duration

LeaseDuration delegates if inner leases jobs; zero disables heartbeats.

func (*RateLimitAdapter) Len

func (r *RateLimitAdapter) Len(ctx context.Context, queue string) (int, error)

Len delegates to inner.

func (*RateLimitAdapter) Pop

func (r *RateLimitAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)

Pop waits for rate limit before returning a job.

func (*RateLimitAdapter) Push

func (r *RateLimitAdapter) Push(ctx context.Context, payload *JobPayload) error

Push delegates to inner (no rate limit on push).

type RateLimitConfig

type RateLimitConfig struct {
	// Limit is the max number of jobs per second (e.g. 10 = 10 jobs/sec).
	Limit float64
	// Burst allows short bursts above the limit.
	Burst int
}

RateLimitConfig configures rate limiting per queue.

type ReclaimObserver

type ReclaimObserver interface {
	JobsReclaimed(queue string, count int)
}

ReclaimObserver can observe reclaimed in-flight jobs.

type RedisAdapter

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

RedisAdapter stores jobs in Redis: a pending list per queue, a sorted set of delayed jobs, and a sorted set of in-flight leases that are returned to the queue when a worker stops heartbeating (see LeaseExtender).

func NewRedisAdapter

func NewRedisAdapter(client *redis.Client) *RedisAdapter

NewRedisAdapter creates a Redis adapter. Pass a configured redis.Client.

func NewRedisAdapterFromURL

func NewRedisAdapterFromURL(url string) (*RedisAdapter, error)

NewRedisAdapterFromURL creates adapter from REDIS_URL (e.g. redis://localhost:6379).

func (*RedisAdapter) Client added in v1.9.1

func (r *RedisAdapter) Client() *redis.Client

Client returns the underlying Redis client.

func (*RedisAdapter) Complete

func (r *RedisAdapter) Complete(ctx context.Context, payload *JobPayload) error

Complete acknowledges and removes a processed in-flight Redis message.

func (*RedisAdapter) ExtendLease added in v1.9.1

func (r *RedisAdapter) ExtendLease(ctx context.Context, payload *JobPayload) error

ExtendLease implements LeaseExtender. It returns ErrLeaseLost when the lease already expired and the job was handed to another worker.

func (*RedisAdapter) LeaseDuration added in v1.9.1

func (r *RedisAdapter) LeaseDuration() time.Duration

LeaseDuration implements LeaseExtender.

func (*RedisAdapter) Len

func (r *RedisAdapter) Len(ctx context.Context, queue string) (int, error)

Len returns the number of pending jobs.

func (*RedisAdapter) Pop

func (r *RedisAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)

Pop blocks until a job is available or ctx is done.

func (*RedisAdapter) Push

func (r *RedisAdapter) Push(ctx context.Context, payload *JobPayload) error

Push adds a job to the queue.

func (*RedisAdapter) SetVisibilityTimeout

func (r *RedisAdapter) SetVisibilityTimeout(v time.Duration)

SetVisibilityTimeout sets how long a job stays leased to a worker without a heartbeat before it is handed to another worker. Workers run by the Manager extend the lease while the job runs, so this bounds how long a crashed worker's job waits, not how long a job may run.

type RedisBatchStore added in v1.9.1

type RedisBatchStore struct {
	TTL time.Duration
	// contains filtered or unexported fields
}

RedisBatchStore keeps batch progress in Redis. Batches expire TTL after their last update.

func NewRedisBatchStore added in v1.9.1

func NewRedisBatchStore(client *redis.Client) *RedisBatchStore

NewRedisBatchStore returns a Redis-backed BatchStore that keeps batches for 7 days after their last update.

func (*RedisBatchStore) Create added in v1.9.1

func (s *RedisBatchStore) Create(ctx context.Context, state BatchState) error

func (*RedisBatchStore) Find added in v1.9.1

func (s *RedisBatchStore) Find(ctx context.Context, batchID string) (*BatchState, error)

func (*RedisBatchStore) Record added in v1.9.1

func (s *RedisBatchStore) Record(ctx context.Context, batchID, jobID string, jobErr error) (BatchUpdate, error)

type RedisFailedStore

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

RedisFailedStore stores failed jobs in Redis for Horizon dashboard.

func NewRedisFailedStore

func NewRedisFailedStore(client *redis.Client) *RedisFailedStore

NewRedisFailedStore creates a failed job store using the given Redis client.

func (*RedisFailedStore) Forget

func (r *RedisFailedStore) Forget(ctx context.Context, id string) error

Forget removes a single failed job.

func (*RedisFailedStore) ForgetAll

func (r *RedisFailedStore) ForgetAll(ctx context.Context) error

ForgetAll removes all failed jobs.

func (*RedisFailedStore) Get

Get returns one record by ID.

func (*RedisFailedStore) List

List returns all failed job records (newest last).

func (*RedisFailedStore) Push

func (r *RedisFailedStore) Push(ctx context.Context, payload *JobPayload, errMsg string) error

Push adds a failed job to the store.

func (*RedisFailedStore) Retry

func (r *RedisFailedStore) Retry(ctx context.Context, id string, enqueue func(ctx context.Context, payload *JobPayload) error) error

Retry re-enqueues the job and removes it from the failed store.

type RedisJobLimiter added in v1.9.1

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

RedisJobLimiter rate-limits each queue across every worker that shares the Redis server.

func NewRedisJobLimiter added in v1.9.1

func NewRedisJobLimiter(client *redis.Client, limitPerSec float64, burst int) *RedisJobLimiter

NewRedisJobLimiter returns a JobLimiter kept in Redis.

func (*RedisJobLimiter) Wait added in v1.9.1

func (l *RedisJobLimiter) Wait(ctx context.Context, queue string) error

type RedisLocker added in v1.9.1

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

RedisLocker uses SET NX PX, so locks are shared by every process using the same Redis.

func NewRedisLocker added in v1.9.1

func NewRedisLocker(client *redis.Client) *RedisLocker

NewRedisLocker returns a Redis-backed Locker.

func (*RedisLocker) Acquire added in v1.9.1

func (r *RedisLocker) Acquire(ctx context.Context, key string, ttl time.Duration) (string, bool, error)

func (*RedisLocker) Release added in v1.9.1

func (r *RedisLocker) Release(ctx context.Context, key, token string) error

type RedisQueueWorkload

type RedisQueueWorkload struct {
	Name       string `json:"name"`
	Pending    int64  `json:"pending"`    // jobs waiting in the main list
	Delayed    int64  `json:"delayed"`    // jobs in delayed sorted set
	Processing int64  `json:"processing"` // jobs leased to workers
	InFlight   int64  `json:"in_flight"`  // same as Processing; kept for older dashboards
}

RedisQueueWorkload holds Redis list/set sizes for one logical queue (Redis driver key layout).

func RedisQueueWorkloads

func RedisQueueWorkloads(ctx context.Context, c *redis.Client, names []string) ([]RedisQueueWorkload, error)

RedisQueueWorkloads returns live depth metrics from Redis for the given queue names. Uses the same key prefixes as RedisAdapter. Pass nil or empty names to probe only "default".

type ReleaseError added in v1.9.1

type ReleaseError struct {
	Delay time.Duration
}

ReleaseError tells the worker to put the job back on the queue after Delay without counting an attempt. Return it from Handle with Release.

func (*ReleaseError) Error added in v1.9.1

func (e *ReleaseError) Error() string

type RetryObserver

type RetryObserver interface {
	JobRetried(payload *JobPayload, nextDelay time.Duration)
}

RetryObserver can observe retry scheduling events.

type SQSAdapter

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

SQSAdapter uses AWS SQS for job storage.

func NewSQSAdapter

func NewSQSAdapter(ctx context.Context, queueURL string) (*SQSAdapter, error)

NewSQSAdapter creates an SQS adapter. queueURL is the full SQS queue URL.

func NewSQSAdapterFromConfig

func NewSQSAdapterFromConfig(client *sqs.Client, queueURL string) *SQSAdapter

NewSQSAdapterFromConfig creates adapter with custom config.

func (*SQSAdapter) Complete

func (s *SQSAdapter) Complete(ctx context.Context, payload *JobPayload) error

Complete deletes the SQS message after successful processing.

func (*SQSAdapter) Len

func (s *SQSAdapter) Len(ctx context.Context, queue string) (int, error)

Len returns approximate message count (SQS ApproximateNumberOfMessages).

func (*SQSAdapter) Pop

func (s *SQSAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)

Pop blocks until a job is available (long polling).

func (*SQSAdapter) Push

func (s *SQSAdapter) Push(ctx context.Context, payload *JobPayload) error

Push adds a job to the queue.

type ScheduleLocker added in v1.9.1

type ScheduleLocker struct{ Locker Locker }

ScheduleLocker adapts a Locker to schedule.Locker, so scheduled tasks lock through the same backend as the queue.

func (ScheduleLocker) TryLock added in v1.9.1

func (s ScheduleLocker) TryLock(ctx context.Context, key string, ttl time.Duration) (func(), bool, error)

TryLock implements schedule.Locker.

type Scheduler

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

Scheduler runs jobs on a schedule.

func NewScheduler

func NewScheduler(m *Manager) *Scheduler

NewScheduler creates a scheduler that dispatches to the given manager.

func (*Scheduler) Cron

func (s *Scheduler) Cron(expr string, job Job) error

Cron adds a job with a cron expression. Examples: "0 0 * * *" = midnight daily, "*/5 * * * *" = every 5 min.

func (*Scheduler) Every

func (s *Scheduler) Every(d time.Duration, job Job) error

Every adds a job that runs at a fixed interval.

func (*Scheduler) Start

func (s *Scheduler) Start(ctx context.Context)

Start runs the scheduler. Blocks until ctx is done.

func (*Scheduler) Stop

func (s *Scheduler) Stop()

Stop stops the scheduler.

type Silenced

type Silenced interface {
	Job
	// Silenced returns true to hide this job from the completed list.
	Silenced() bool
}

Silenced is optionally implemented to hide the job from Horizon's completed jobs list.

type SyncAdapter

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

SyncAdapter runs jobs immediately in the same process. No persistence. Use for development/testing when you don't need a separate worker.

func NewSyncAdapter

func NewSyncAdapter(m *Manager) *SyncAdapter

NewSyncAdapter creates a sync adapter. Requires manager for job deserialization.

func (*SyncAdapter) Len

func (s *SyncAdapter) Len(ctx context.Context, queue string) (int, error)

Len always returns 0.

func (*SyncAdapter) Pop

func (s *SyncAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)

Pop always blocks (sync has no persisted jobs). Use Redis/Database for workers.

func (*SyncAdapter) Push

func (s *SyncAdapter) Push(ctx context.Context, payload *JobPayload) error

Push runs the job immediately in-process. It mirrors the behavior of Manager.Process enough to keep observers (e.g. Horizon) informed.

type Tagger

type Tagger interface {
	Job
	Tags() []string
}

Tagger is optionally implemented to assign tags for Horizon dashboard (Laravel-style).

type UniqueJob

type UniqueJob interface {
	Job

	// UniqueID returns a unique identifier for this job instance.
	// Jobs with the same UniqueID will not be dispatched while one is pending.
	UniqueID() string

	// UniqueFor returns the longest the unique lock is held. The lock is
	// released earlier when the job finishes (succeeds or fails for good).
	// Zero or less means one hour.
	UniqueFor() time.Duration
}

UniqueJob interface — jobs implementing this will be deduplicated.

type WithoutOverlapping

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

WithoutOverlapping wraps a Job so that only one job with the same uniqueID runs at a time, across all workers. Dispatch it like any job; the inner job must be registered with queue.Register.

func NewWithoutOverlapping

func NewWithoutOverlapping(job Job, uniqueID string) *WithoutOverlapping

NewWithoutOverlapping wraps a job to prevent overlapping execution.

func (*WithoutOverlapping) ExpireAfter added in v1.9.1

ExpireAfter bounds how long the lock is held if a worker dies mid-job.

func (*WithoutOverlapping) Handle

func (w *WithoutOverlapping) Handle(ctx context.Context) error

Handle runs the inner job. Workers take the lock before calling Handle; called directly, Handle waits for the lock itself.

func (*WithoutOverlapping) Inner added in v1.9.1

func (w *WithoutOverlapping) Inner() Job

Inner returns the wrapped job.

func (*WithoutOverlapping) MarshalJSON added in v1.9.1

func (w *WithoutOverlapping) MarshalJSON() ([]byte, error)

MarshalJSON lets the wrapper travel through a queue.

func (*WithoutOverlapping) OverlapExpiresAfter added in v1.9.1

func (w *WithoutOverlapping) OverlapExpiresAfter() time.Duration

OverlapExpiresAfter implements OverlapOptions.

func (*WithoutOverlapping) OverlapKey added in v1.9.1

func (w *WithoutOverlapping) OverlapKey() string

OverlapKey implements NonOverlapping.

func (*WithoutOverlapping) OverlapReleaseAfter added in v1.9.1

func (w *WithoutOverlapping) OverlapReleaseAfter() time.Duration

OverlapReleaseAfter implements OverlapOptions.

func (*WithoutOverlapping) ReleaseAfter added in v1.9.1

func (w *WithoutOverlapping) ReleaseAfter(d time.Duration) *WithoutOverlapping

ReleaseAfter sets how long a blocked job waits before it is tried again.

func (*WithoutOverlapping) UnmarshalJSON added in v1.9.1

func (w *WithoutOverlapping) UnmarshalJSON(data []byte) error

UnmarshalJSON rebuilds the inner job from the manager's job registry.

Jump to

Keyboard shortcuts

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