Documentation
¶
Index ¶
- Constants
- type Memory
- func (m *Memory) Ack(_ context.Context, p *queue.Payload) error
- func (m *Memory) Dead(_ context.Context, p *queue.Payload) error
- func (m *Memory) DeadLetters() []*queue.Payload
- func (m *Memory) InFlight() int
- func (m *Memory) Pop(ctx context.Context, q string) (*queue.Payload, error)
- func (m *Memory) Push(_ context.Context, q string, p *queue.Payload) error
- func (m *Memory) Release(_ context.Context, p *queue.Payload, delay time.Duration) error
- func (m *Memory) Size(q string) int
- type Redis
- func (r *Redis) Ack(ctx context.Context, p *queue.Payload) error
- func (r *Redis) Dead(ctx context.Context, p *queue.Payload) error
- func (r *Redis) DeadLetters(ctx context.Context, n int64) ([]*queue.Payload, error)
- func (r *Redis) Pop(ctx context.Context, q string) (*queue.Payload, error)
- func (r *Redis) Push(ctx context.Context, q string, p *queue.Payload) error
- func (r *Redis) Release(ctx context.Context, p *queue.Payload, delay time.Duration) error
Constants ¶
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 (*Memory) DeadLetters ¶
DeadLetters returns all dead-lettered payloads (for testing/inspection).
func (*Memory) InFlight ¶
InFlight returns the number of reserved (popped, unacknowledged) payloads.
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 ¶
NewRedis creates a Redis queue driver.
d := drivers.NewRedis("redis://localhost:6379")
func NewRedisFromClient ¶
NewRedisFromClient creates a Redis queue driver from an existing client.
func (*Redis) DeadLetters ¶
DeadLetters returns up to n dead-lettered payloads.