redis

package
v0.2.0 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 14 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrStreamsClosed reports an operation attempted after the streams
	// service started shutting down.
	ErrStreamsClosed = errors.New("redis streams service is closed")
	// ErrStreamsWorkerCapacity reports that Consume requested more
	// long-lived Redis group readers than the service-owned pool has available.
	ErrStreamsWorkerCapacity = errors.New("redis streams worker capacity is exhausted")
	// ErrStreamsDeliveryLost reports that a message was reclaimed by another
	// consumer before this delivery could acknowledge it.
	ErrStreamsDeliveryLost = errors.New("redis stream delivery ownership was lost")
)
View Source
var ConfigToken = gioc.NewToken("RedisConfig")
View Source
var Module = ForFeature()

Module preserves the base Redis module API for applications that do not use an optional feature.

View Source
var ServiceInjections = react.InjectFromBase(ConfigToken)
View Source
var ServiceToken = gioc.NewToken("RedisService")
View Source
var Streams = Feature{
	// contains filtered or unexported fields
}

Streams enables the lifecycle-owned Redis Streams worker hub.

View Source
var StreamsConfigToken = gioc.NewToken("StreamsConfig")

StreamsConfigToken is intentionally separate from ConfigToken. The application must provide it only when redis.ForFeature(redis.Streams) is selected.

View Source
var StreamsServiceToken = gioc.NewToken("StreamsService")

Functions

func ForFeature added in v0.2.0

func ForFeature(features ...Feature) *gioc.Module

ForFeature creates a fresh Redis module with the requested optional capabilities. Select redis.Streams to make StreamsServiceToken injectable. Its dedicated config remains the application's responsibility.

func ProvideConfig added in v0.2.0

func ProvideConfig(config *Config) gioc.IProvider

func ProvideStreamsConfig added in v0.2.0

func ProvideStreamsConfig(config *StreamsConfig) gioc.IProvider

ProvideStreamsConfig makes the application-owned configuration available through StreamsConfigToken.

Types

type Config

type Config struct {
	URL string `env:"REDIS_URL"`
}

type Feature added in v0.2.0

type Feature struct {
	// contains filtered or unexported fields
}

Feature is an immutable Redis capability selected by ForFeature. Features add only their own providers; the base Redis service is always available.

func (Feature) Name added in v0.2.0

func (feature Feature) Name() string

Name returns the stable feature name used in the generated module token.

type IStreamsService added in v0.2.0

type IStreamsService interface {
	Publish(ctx context.Context, stream string, key string, value any) error
	Consume(ctx context.Context, group string, stream string, configs ...StreamsConsumerConfig) (<-chan StreamsMessage, error)
	Ack(ctx context.Context, message StreamsMessage) error
}

IStreamsService is the lightweight application contract for Redis Streams. Consume uses manual acknowledgements and at-least-once delivery.

type Service

type Service struct {
	react.BaseConfigurableService[*Config]
	// contains filtered or unexported fields
}

func NewService

func NewService(injections gioc.Injections) (*Service, error)

func (*Service) Add

func (service *Service) Add(key string, members ...any) (int64, error)

func (*Service) Client added in v0.2.0

func (service *Service) Client() goredis.UniversalClient

Client exposes the application-owned, context-aware go-redis client for infrastructure adapters such as outbox.RedisStore. Callers must not close it; Service retains lifecycle ownership. Individual operations must pass their request or application context rather than reusing Service.Ctx.

func (*Service) Delete

func (service *Service) Delete(keys ...string) error

func (*Service) Drop

func (service *Service) Drop(key string, members ...any) (int64, error)

func (*Service) Get

func (service *Service) Get(key string) (string, bool, error)

func (*Service) GetBool

func (service *Service) GetBool(key string) (bool, bool, error)

func (*Service) IncludesMany

func (service *Service) IncludesMany(key string, members ...string) ([]bool, error)

func (*Service) IncludesOne

func (service *Service) IncludesOne(key string, member any) (bool, error)

func (*Service) Keys

func (service *Service) Keys(template string) ([]string, error)

func (*Service) Scan

func (service *Service) Scan(template string, count int64) ([]string, error)

func (*Service) Set

func (service *Service) Set(key string, value any, ttl time.Duration) error

func (*Service) SetIfNotExists

func (service *Service) SetIfNotExists(key string, value any, ttl time.Duration) (bool, error)

func (*Service) Size

func (service *Service) Size(key string) (int64, error)

type StreamsConfig added in v0.2.0

type StreamsConfig struct {
	WorkerCount          int
	ChannelSize          int
	DefaultConsumerCount int
	DefaultBatchSize     int64
	BlockTimeout         time.Duration
	ReclaimInterval      time.Duration
	ReclaimAfter         time.Duration
	MaximumDeliveries    int64
	DeadLetterSuffix     string
	RetryMinimumDelay    time.Duration
	RetryMaximumDelay    time.Duration
	MaximumMessageBytes  int
	MaximumStreamLength  int64
}

StreamsConfig controls the single inbound worker pool and the default bounded channel created by every Consume call. Zero-valued fields receive the values returned by DefaultStreamsConfig.

func DefaultStreamsConfig added in v0.2.0

func DefaultStreamsConfig() StreamsConfig

DefaultStreamsConfig returns bounded defaults suitable for the common one-group, short-message case. MaximumStreamLength is zero by design: stream retention is an application decision and is never enabled implicitly.

type StreamsConsumerConfig added in v0.2.0

type StreamsConsumerConfig struct {
	ConsumerCount int
	BatchSize     int64
	StartFrom     StreamsStart
}

StreamsConsumerConfig overrides per-subscription read concurrency and batch size. The output channel capacity always comes from StreamsConfig so memory use remains bounded at service level.

type StreamsMessage added in v0.2.0

type StreamsMessage struct {
	ID       string
	Stream   string
	Group    string
	Key      string
	Attempts int64
	Payload  json.RawMessage
	// contains filtered or unexported fields
}

StreamsMessage is one manual-acknowledgement group delivery. Payload is the JSON value supplied to Publish. Attempts starts at one and increases when an idle pending entry is reclaimed.

func (StreamsMessage) Decode added in v0.2.0

func (message StreamsMessage) Decode(target any) error

Decode unmarshals the message payload into target.

type StreamsService added in v0.2.0

type StreamsService struct {
	ApplicationService *react.ApplicationService
	Logger             react.ILogger
	// contains filtered or unexported fields
}

StreamsService owns one fixed workers.Pool. Each Consume call reserves one or more pool workers as long-lived XREADGROUP readers and receives a service-sized bounded channel. It never creates a worker pool per stream.

func NewStreamsService added in v0.2.0

func NewStreamsService(injections gioc.Injections) (*StreamsService, error)

NewStreamsService resolves every dependency from the feature factory, validates the application-provided configuration, and starts the one fixed inbound worker pool.

func (*StreamsService) Ack added in v0.2.0

func (service *StreamsService) Ack(ctx context.Context, message StreamsMessage) error

Ack removes a delivery from the group's pending list. Ownership and delivery count are checked atomically so an old delivery cannot acknowledge a message that another consumer has already reclaimed. Ack is idempotent after success.

func (*StreamsService) Consume added in v0.2.0

func (service *StreamsService) Consume(
	ctx context.Context,
	group string,
	stream string,
	configs ...StreamsConsumerConfig,
) (<-chan StreamsMessage, error)

Consume creates the group when needed, reserves readers from the one service pool, and returns a bounded manual-acknowledgement channel. At most one optional config is accepted. Cancelling ctx releases the reserved workers and closes the returned channel.

func (*StreamsService) Publish added in v0.2.0

func (service *StreamsService) Publish(ctx context.Context, stream string, key string, value any) error

Publish JSON-encodes value and appends it to stream. A successful return means Redis accepted XADD; Redis persistence and replication remain server policy. No retention is applied unless MaximumStreamLength is configured.

type StreamsStart added in v0.2.0

type StreamsStart string

StreamsStart is the initial group position used only when Consume has to create the group. Existing groups keep their persisted position.

const (
	StreamsStartBeginning StreamsStart = "0-0"
	StreamsStartLatest    StreamsStart = "$"
)

Jump to

Keyboard shortcuts

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