drivers

package
v1.0.0 Latest Latest
Warning

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

Go to latest
Published: Jul 3, 2026 License: MIT Imports: 8 Imported by: 0

Documentation

Index

Constants

View Source
const DefaultVisibility = 10 * time.Minute

DefaultVisibility is how long a popped payload may stay unacknowledged before the driver assumes the worker died and requeues it. It must exceed the manager's per-job timeout (5 minutes) so a slow-but-alive job is not delivered twice.

Variables

This section is empty.

Functions

This section is empty.

Types

type Memory

type Memory struct {

	// Visibility overrides DefaultVisibility when set (useful in tests).
	Visibility time.Duration
	// contains filtered or unexported fields
}

Memory is a non-durable in-process queue driver. Suitable for development.

Pop reserves rather than removes: a popped payload is tracked in-flight and requeued if neither Ack, Release, nor Dead is called before its visibility timeout expires (at-least-once delivery).

func NewMemory

func NewMemory() *Memory

NewMemory creates an in-memory queue driver.

func (*Memory) Ack

func (m *Memory) Ack(_ context.Context, p *queue.Payload) error

func (*Memory) Dead

func (m *Memory) Dead(_ context.Context, p *queue.Payload) error

func (*Memory) DeadLetters

func (m *Memory) DeadLetters() []*queue.Payload

DeadLetters returns all dead-lettered payloads (for testing/inspection).

func (*Memory) InFlight

func (m *Memory) InFlight() int

InFlight returns the number of reserved (popped, unacknowledged) payloads.

func (*Memory) Pop

func (m *Memory) Pop(ctx context.Context, q string) (*queue.Payload, error)

func (*Memory) Push

func (m *Memory) Push(_ context.Context, q string, p *queue.Payload) error

func (*Memory) Release

func (m *Memory) Release(_ context.Context, p *queue.Payload, delay time.Duration) error

func (*Memory) Size

func (m *Memory) Size(q string) int

Size returns the number of pending jobs in a queue (excludes in-flight).

type Redis

type Redis struct {

	// Visibility overrides DefaultVisibility when set.
	Visibility time.Duration
	// contains filtered or unexported fields
}

Redis is a Redis-backed queue driver: lists for ready jobs, a sorted set for delayed jobs, and a per-queue reservation set for in-flight jobs.

Pop reserves rather than removes: the popped payload is atomically moved into a processing set scored by its visibility deadline. Ack/Release/Dead clear the reservation; if a worker crashes first, the reaper (run on every Pop) pushes the payload back onto the ready list once the deadline passes. This gives at-least-once delivery instead of losing jobs on crash.

func NewRedis

func NewRedis(redisURL string) (*Redis, error)

NewRedis creates a Redis queue driver.

d := drivers.NewRedis("redis://localhost:6379")

func NewRedisFromClient

func NewRedisFromClient(c *redis.Client) *Redis

NewRedisFromClient creates a Redis queue driver from an existing client.

func (*Redis) Ack

func (r *Redis) Ack(ctx context.Context, p *queue.Payload) error

func (*Redis) Dead

func (r *Redis) Dead(ctx context.Context, p *queue.Payload) error

func (*Redis) DeadLetters

func (r *Redis) DeadLetters(ctx context.Context, n int64) ([]*queue.Payload, error)

DeadLetters returns up to n dead-lettered payloads.

func (*Redis) Pop

func (r *Redis) Pop(ctx context.Context, q string) (*queue.Payload, error)

func (*Redis) Push

func (r *Redis) Push(ctx context.Context, q string, p *queue.Payload) error

func (*Redis) Release

func (r *Redis) Release(ctx context.Context, p *queue.Payload, delay time.Duration) error

Jump to

Keyboard shortcuts

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