Documentation
¶
Index ¶
- Variables
- func AcquireUniqueLock(job UniqueJob) bool
- func AddObserver(o Observer)
- func AfterBatch(fn func(context.Context, *Batch))
- func DispatchUnique(ctx context.Context, job UniqueJob) error
- func IsUniqueLocked(job UniqueJob) bool
- func Register(job Job)
- func RegisterBatch(name string, cb BatchCallbacks)
- func RegisterChain(name string, onFailure func(ctx context.Context, failedJob Job, err error))
- func Release(delay time.Duration) error
- func ReleaseUniqueLock(job UniqueJob)
- func RunWorker(ctx context.Context, queueName string)
- func SetBatchStore(s BatchStore)
- func SetFailedJobStore(s FailedJobStore)
- func SetGlobal(m *Manager)
- func SetLocker(l Locker)
- func SetObserver(o Observer)
- type Adapter
- type Batch
- func (b *Batch) Catch(fn func(ctx context.Context, batch *Batch, err error)) *Batch
- func (b *Batch) Dispatch(ctx context.Context) error
- func (b *Batch) Errors() []error
- func (b *Batch) FailedJobs() int32
- func (b *Batch) Finally(fn func(ctx context.Context, batch *Batch)) *Batch
- func (b *Batch) Finished() bool
- func (b *Batch) HasFailures() bool
- func (b *Batch) Name() string
- func (b *Batch) Named(name string) *Batch
- func (b *Batch) OnQueue(name string) *Batch
- func (b *Batch) PendingJobs() int32
- func (b *Batch) Refresh(ctx context.Context) error
- func (b *Batch) Run(ctx context.Context) error
- func (b *Batch) Then(fn func(ctx context.Context, batch *Batch)) *Batch
- func (b *Batch) TotalJobs() int32
- type BatchCallbacks
- type BatchState
- type BatchStore
- type BatchUpdate
- type BootConfig
- type Chain
- type CompletableAdapter
- type DatabaseAdapter
- func (d *DatabaseAdapter) Complete(ctx context.Context, payload *JobPayload) error
- func (d *DatabaseAdapter) EnsureTable(ctx context.Context) error
- func (d *DatabaseAdapter) ExtendLease(ctx context.Context, payload *JobPayload) error
- func (d *DatabaseAdapter) LeaseDuration() time.Duration
- func (d *DatabaseAdapter) Len(ctx context.Context, queue string) (int, error)
- func (d *DatabaseAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)
- func (d *DatabaseAdapter) Push(ctx context.Context, payload *JobPayload) error
- func (d *DatabaseAdapter) SetLeaseDuration(v time.Duration)
- func (d *DatabaseAdapter) SetPollInterval(v time.Duration)
- type DatabaseBatchStore
- func (s *DatabaseBatchStore) Create(ctx context.Context, state BatchState) error
- func (s *DatabaseBatchStore) EnsureTable(ctx context.Context) error
- func (s *DatabaseBatchStore) Find(ctx context.Context, batchID string) (*BatchState, error)
- func (s *DatabaseBatchStore) Record(ctx context.Context, batchID, jobID string, jobErr error) (BatchUpdate, error)
- type DatabaseLocker
- type DispatchBuilder
- func (b *DispatchBuilder) Delay(d time.Duration) *DispatchBuilder
- func (b *DispatchBuilder) Dispatch(ctx context.Context) error
- func (b *DispatchBuilder) OnQueue(name string) *DispatchBuilder
- func (b *DispatchBuilder) Priority(n int) *DispatchBuilder
- func (b *DispatchBuilder) Retries(n int) *DispatchBuilder
- type FailedJob
- type FailedJobRecord
- type FailedJobStore
- type Job
- type JobFunc
- type JobLimiter
- type JobPayload
- type KafkaAdapter
- type KafkaConfig
- type LeaseExtender
- type LocalJobLimiter
- type Locker
- type Manager
- type MemoryBatchStore
- type MemoryLocker
- type NonOverlapping
- type Observer
- type ObserverV2
- type OverlapOptions
- type Queue
- type QueueBatch
- type QueueBatchJob
- type QueueJob
- type QueueLock
- type RateLimitAdapter
- func (r *RateLimitAdapter) Complete(ctx context.Context, payload *JobPayload) error
- func (r *RateLimitAdapter) ExtendLease(ctx context.Context, payload *JobPayload) error
- func (r *RateLimitAdapter) Inner() Adapter
- func (r *RateLimitAdapter) LeaseDuration() time.Duration
- func (r *RateLimitAdapter) Len(ctx context.Context, queue string) (int, error)
- func (r *RateLimitAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)
- func (r *RateLimitAdapter) Push(ctx context.Context, payload *JobPayload) error
- type RateLimitConfig
- type ReclaimObserver
- type RedisAdapter
- func (r *RedisAdapter) Client() *redis.Client
- func (r *RedisAdapter) Complete(ctx context.Context, payload *JobPayload) error
- func (r *RedisAdapter) ExtendLease(ctx context.Context, payload *JobPayload) error
- func (r *RedisAdapter) LeaseDuration() time.Duration
- func (r *RedisAdapter) Len(ctx context.Context, queue string) (int, error)
- func (r *RedisAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)
- func (r *RedisAdapter) Push(ctx context.Context, payload *JobPayload) error
- func (r *RedisAdapter) SetVisibilityTimeout(v time.Duration)
- type RedisBatchStore
- type RedisFailedStore
- func (r *RedisFailedStore) Forget(ctx context.Context, id string) error
- func (r *RedisFailedStore) ForgetAll(ctx context.Context) error
- func (r *RedisFailedStore) Get(ctx context.Context, id string) (*FailedJobRecord, error)
- func (r *RedisFailedStore) List(ctx context.Context) ([]FailedJobRecord, error)
- func (r *RedisFailedStore) Push(ctx context.Context, payload *JobPayload, errMsg string) error
- func (r *RedisFailedStore) Retry(ctx context.Context, id string, ...) error
- type RedisJobLimiter
- type RedisLocker
- type RedisQueueWorkload
- type ReleaseError
- type RetryObserver
- type SQSAdapter
- func (s *SQSAdapter) Complete(ctx context.Context, payload *JobPayload) error
- func (s *SQSAdapter) Len(ctx context.Context, queue string) (int, error)
- func (s *SQSAdapter) Pop(ctx context.Context, queue string) (*JobPayload, error)
- func (s *SQSAdapter) Push(ctx context.Context, payload *JobPayload) error
- type ScheduleLocker
- type Scheduler
- type Silenced
- type SyncAdapter
- type Tagger
- type UniqueJob
- type WithoutOverlapping
- func (w *WithoutOverlapping) ExpireAfter(d time.Duration) *WithoutOverlapping
- func (w *WithoutOverlapping) Handle(ctx context.Context) error
- func (w *WithoutOverlapping) Inner() Job
- func (w *WithoutOverlapping) MarshalJSON() ([]byte, error)
- func (w *WithoutOverlapping) OverlapExpiresAfter() time.Duration
- func (w *WithoutOverlapping) OverlapKey() string
- func (w *WithoutOverlapping) OverlapReleaseAfter() time.Duration
- func (w *WithoutOverlapping) ReleaseAfter(d time.Duration) *WithoutOverlapping
- func (w *WithoutOverlapping) UnmarshalJSON(data []byte) error
Constants ¶
This section is empty.
Variables ¶
var ErrBatchNotFound = errors.New("queue: batch not found")
ErrBatchNotFound is returned when a batch ID is unknown (or expired).
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 ¶
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 ¶
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 ¶
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 ¶
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
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
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 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 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 (*Batch) Dispatch ¶
Dispatch queues the batch's jobs (see Batch). It returns once every job is queued; use FindBatch or Refresh to follow progress.
func (*Batch) FailedJobs ¶
FailedJobs returns the number of failed jobs.
func (*Batch) Finally ¶
Finally sets a callback that runs after all jobs complete (success or failure).
func (*Batch) HasFailures ¶
HasFailures returns true if any jobs failed.
func (*Batch) Named ¶ added in v1.9.1
Named names the batch. Workers look callbacks up by this name.
func (*Batch) PendingJobs ¶
PendingJobs returns the number of jobs still pending.
func (*Batch) Refresh ¶ added in v1.9.1
Refresh reloads a durable batch's counters from the BatchStore.
func (*Batch) Run ¶ added in v1.9.1
Run runs all jobs in the batch concurrently in this process and blocks until they finish. Nothing is persisted.
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.
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 (*Chain) DispatchAsync ¶
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
Named names the chain. DispatchAsync needs a name when OnFailure is set, so workers can find the callback (see RegisterChain).
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) 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) 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.
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 ¶
func (b *DispatchBuilder) Delay(d time.Duration) *DispatchBuilder
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 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 JobLimiter ¶ added in v1.9.1
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) 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.
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.
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 NewManager ¶
NewManager creates a manager with the given adapter. Pass nil to use SyncAdapter.
func (*Manager) Dispatch ¶
func (m *Manager) Dispatch(job Job) *DispatchBuilder
Dispatch enqueues a job. Returns a DispatchBuilder for options.
func (*Manager) Register ¶
Register registers a job type for deserialization. Call with a zero-value instance.
queue.Register(&jobs.SendEmail{})
func (*Manager) RegisterFunc ¶
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.
type NonOverlapping ¶ added in v1.9.1
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.
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.
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.
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) 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 ¶
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) 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
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 ¶
func (r *RedisFailedStore) Get(ctx context.Context, id string) (*FailedJobRecord, error)
Get returns one record by ID.
func (*RedisFailedStore) List ¶
func (r *RedisFailedStore) List(ctx context.Context) ([]FailedJobRecord, error)
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.
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.
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.
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
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) 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.
type Scheduler ¶
type Scheduler struct {
// contains filtered or unexported fields
}
Scheduler runs jobs on a schedule.
func NewScheduler ¶
NewScheduler creates a scheduler that dispatches to the given manager.
func (*Scheduler) Cron ¶
Cron adds a job with a cron expression. Examples: "0 0 * * *" = midnight daily, "*/5 * * * *" = every 5 min.
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) 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 ¶
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
func (w *WithoutOverlapping) ExpireAfter(d time.Duration) *WithoutOverlapping
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.