Versions in this module Expand all Collapse all v0 v0.0.1 Aug 2, 2026 Changes in this version + const DefaultConsumerConcurrency + const DefaultConsumerGroup + const DefaultConsumerPollInterval + const DefaultConsumerReadCount + const DefaultDLQListCount + const DefaultMaxAttempts + const DefaultMaxPayloadBytes + const DefaultNamespace + const DefaultPendingMinIdle + const DefaultQueueName + const DefaultRetryInterval + const DefaultSchedulerBatchSize + const DefaultSchedulerLeaseTTL + const DefaultSchedulerShutdownTimeout + const DefaultSchedulerTickInterval + const DefaultShardCount + var ErrDLQMessageNotFound = errors.New("redq: dlq message not found") + var ErrMissingPayload = errors.New("redq: task payload missing") + var ErrStaleMessage = errors.New("redq: stale stream message") + func NewUUID() (string, error) + func ValidateName(field string, name string) error + func ValidateQueueName(name string) error + type CancelResult struct + Canceled bool + Queue string + Shard int + TaskID string + type CancelableProducer interface + type Canceler interface + Cancel func(ctx context.Context, taskID string) (*CancelResult, error) + type Config struct + Namespace string + Queues []QueueConfig + Redis RedisConfig + func DefaultConfig() Config + func (c Config) Validate() error + type Consumer interface + Start func(ctx context.Context) error + type ConsumerClient interface + Eval func(ctx context.Context, script string, keys []string, args ...interface{}) *redis.Cmd + HGet func(ctx context.Context, key, field string) *redis.StringCmd + XAutoClaim func(ctx context.Context, args *redis.XAutoClaimArgs) *redis.XAutoClaimCmd + XGroupCreateMkStream func(ctx context.Context, stream, group, start string) *redis.StatusCmd + XReadGroup func(ctx context.Context, args *redis.XReadGroupArgs) *redis.XStreamSliceCmd + type ConsumerOptions struct + Concurrency int + ConsumerName string + Group string + MissingPayload MissingPayloadPolicy + Namespace string + Now func() time.Time + PayloadCleanup PayloadCleanupPolicy + PendingMinIdle time.Duration + PollInterval time.Duration + Queue QueueConfig + ReadCount int + RecoveryCount int + RetryPolicy RetryPolicy + StreamRetention StreamRetentionPolicy + func DefaultConsumerOptions(queue string) ConsumerOptions + func (o ConsumerOptions) Validate() error + type DLQClient interface + Eval func(ctx context.Context, script string, keys []string, args ...interface{}) *redis.Cmd + HGet func(ctx context.Context, key, field string) *redis.StringCmd + XRangeN func(ctx context.Context, stream, start, stop string, count int64) *redis.XMessageSliceCmd + type DLQListOptions struct + Count int64 + Shard int + Start string + Stop string + func (o DLQListOptions) Validate() error + type DLQMessage struct + Error string + FailedAt time.Time + ID string + Shard int + Task Task + type DLQOptions struct + Namespace string + Now func() time.Time + Queue QueueConfig + func DefaultDLQOptions(queue string) DLQOptions + func (o DLQOptions) Validate() error + type DLQReplayBatchOptions struct + Count int64 + RunAt time.Time + Shard int + Start string + Stop string + func (o DLQReplayBatchOptions) Validate() error + type DLQReplayResult struct + Queue string + Replayed bool + RunAt time.Time + Shard int + TaskID string + type EnqueueResult struct + Duplicate bool + Queue string + Shard int + TaskID string + type ExponentialRetryPolicy struct + Initial time.Duration + Max time.Duration + Multiplier float64 + func (p ExponentialRetryPolicy) NextDelay(attempt int) time.Duration + type FixedRetryPolicy struct + Interval time.Duration + func (p FixedRetryPolicy) NextDelay(int) time.Duration + type Handler interface + Handle func(ctx context.Context, task *Task) error + type IDGenerator func() (string, error) + type MissingPayloadPolicy int + const MissingPayloadDefault + const MissingPayloadKeepPending + const MissingPayloadMoveToDLQ + type PayloadCleanupPolicy int + const PayloadCleanupDefault + const PayloadCleanupNever + const PayloadCleanupOnSuccess + type Producer interface + Enqueue func(ctx context.Context, task Task) (*EnqueueResult, error) + type ProducerClient interface + Eval func(ctx context.Context, script string, keys []string, args ...interface{}) *redis.Cmd + type ProducerOptions struct + IDGenerator IDGenerator + Namespace string + Now func() time.Time + Queue QueueConfig + func DefaultProducerOptions(queue string) ProducerOptions + func (o ProducerOptions) Validate() error + type QueueConfig struct + DefaultAttempts int + MaxPayloadBytes int + Name string + Shards int + func DefaultQueueConfig(name string) QueueConfig + func (c QueueConfig) Validate() error + type RedisClient = redis.UniversalClient + func NewRedisClient(config RedisConfig) (RedisClient, error) + type RedisConfig struct + Addrs []string + DB int + DialTimeout time.Duration + MinIdleConns int + Password string + PoolSize int + Protocol int + ReadTimeout time.Duration + Username string + WriteTimeout time.Duration + func DefaultRedisConfig() RedisConfig + func (c RedisConfig) UniversalOptions() *redis.UniversalOptions + func (c RedisConfig) Validate() error + type RedisConsumer struct + func NewConsumer(client ConsumerClient, handler Handler, options ConsumerOptions) (*RedisConsumer, error) + func (c *RedisConsumer) RecoverPendingOnce(ctx context.Context) (int, error) + func (c *RedisConsumer) RunOnce(ctx context.Context) (int, error) + func (c *RedisConsumer) Start(ctx context.Context) error + type RedisDLQ struct + func NewDLQ(client DLQClient, options DLQOptions) (*RedisDLQ, error) + func (d *RedisDLQ) ListMessages(ctx context.Context, options DLQListOptions) ([]DLQMessage, error) + func (d *RedisDLQ) ReplayBatch(ctx context.Context, options DLQReplayBatchOptions) ([]DLQReplayResult, error) + func (d *RedisDLQ) ReplayMessage(ctx context.Context, shard int, dlqMessageID string, runAt time.Time) (*DLQReplayResult, error) + type RedisProducer struct + func NewProducer(client ProducerClient, options ProducerOptions) (*RedisProducer, error) + func (p *RedisProducer) Cancel(ctx context.Context, taskID string) (*CancelResult, error) + func (p *RedisProducer) Enqueue(ctx context.Context, task Task) (*EnqueueResult, error) + type RedisScheduler struct + func NewScheduler(client SchedulerClient, options SchedulerOptions) (*RedisScheduler, error) + func (s *RedisScheduler) RunOnce(ctx context.Context) (int, error) + func (s *RedisScheduler) Start(ctx context.Context) error + type RetryPolicy interface + NextDelay func(attempt int) time.Duration + type Scheduler interface + Start func(ctx context.Context) error + type SchedulerClient interface + Eval func(ctx context.Context, script string, keys []string, args ...interface{}) *redis.Cmd + type SchedulerOptions struct + BatchSize int + InstanceID string + LeaseTTL time.Duration + Namespace string + Now func() time.Time + Queue QueueConfig + ShutdownTimeout time.Duration + TickInterval time.Duration + func DefaultSchedulerOptions(queue string) SchedulerOptions + func (o SchedulerOptions) Validate() error + type StreamRetentionPolicy int + const StreamRetentionDefault + const StreamRetentionDeleteOnAck + const StreamRetentionKeep + type Task struct + Attempt int + ID string + IdempotencyKey string + LastError string + MaxAttempts int + Metadata map[string]string + Payload []byte + Queue string + RunAt time.Time + func (t Task) EffectiveMaxAttempts() int + func (t Task) Validate() error + type ValidationError struct + Field string + Reason string + func (e ValidationError) Error() string