Versions in this module Expand all Collapse all v1 v1.0.1 Aug 28, 2026 v1.0.0 Aug 26, 2026 Changes in this version + var ErrInvalidManagementStatus = errors.New("redisstream: invalid management status") + var ErrManagementControlDisabled = errors.Join(errors.New("redisstream: management control disabled"), ...) + var ErrManagementStatusDisabled = errors.New("redisstream: management status disabled") + type Option func(*options) + func WithAddr(addr string) Option + func WithBlockTime(m time.Duration) Option + func WithCluster() Option + func WithCommandTimeout(timeout time.Duration) Option + func WithConnectTimeout(timeout time.Duration) Option + func WithConnectionString(connectionString string) Option + func WithConsumer(name string) Option + func WithDB(db int) Option + func WithDeadLetter(stream string, maxAttempts int64) Option + func WithFailureStream(stream string) Option + func WithGroup(name string) Option + func WithLogger(l queue.Logger) Option + func WithManagementStatus(metadata management.StatusMetadata) Option + func WithMaxLength(m int64) Option + func WithPassword(passwd string) Option + func WithReclaim(minIdle, interval time.Duration, batchSize int64) Option + func WithRecordRetention(maxRecords int64) Option + func WithReplayDestinations(destinations ...string) Option + func WithRequestTimeout(timeout time.Duration) Option + func WithRunFunc(fn func(context.Context, core.TaskMessage) error) Option + func WithSkipTLSVerify() Option + func WithStreamName(name string) Option + func WithTLS() Option + func WithUsername(username string) Option + type Stats struct + Depth int64 + Lag int64 + LagKnown bool + OldestJobAge time.Duration + Pending int64 + type Worker struct + func NewWorker(opts ...Option) *Worker + func NewWorkerE(opts ...Option) (*Worker, error) + func (*Worker) BackendName() string + func (w *Worker) Execute(ctx context.Context, command management.Command) (management.CommandResult, error) + func (w *Worker) Inspect(ctx context.Context, request management.InspectRequest) (management.JobRecord, error) + func (w *Worker) ListDeadLetters(ctx context.Context, request management.PageRequest) (management.RecordPage, error) + func (w *Worker) ListFailures(ctx context.Context, request management.PageRequest) (management.RecordPage, error) + func (w *Worker) ObserveQueue(ctx context.Context) (management.QueueStatus, error) + func (w *Worker) ObserveWorker(ctx context.Context) (management.WorkerStatus, error) + func (w *Worker) Queue(task core.TaskMessage) error + func (w *Worker) QueueName() string + func (w *Worker) Request() (core.TaskMessage, error) + func (w *Worker) Run(ctx context.Context, task core.TaskMessage) error + func (w *Worker) Shutdown() error + func (w *Worker) Stats(ctx context.Context) (Stats, error)