Documentation
¶
Index ¶
- Variables
- func ForFeature(features ...Feature) *gioc.Module
- func ProvideConfig(config *Config) gioc.IProvider
- func ProvideStreamsConfig(config *StreamsConfig) gioc.IProvider
- type Config
- type Feature
- type IStreamsService
- type Service
- func (service *Service) Add(key string, members ...any) (int64, error)
- func (service *Service) Client() goredis.UniversalClient
- func (service *Service) Delete(keys ...string) error
- func (service *Service) Drop(key string, members ...any) (int64, error)
- func (service *Service) Get(key string) (string, bool, error)
- func (service *Service) GetBool(key string) (bool, bool, error)
- func (service *Service) IncludesMany(key string, members ...string) ([]bool, error)
- func (service *Service) IncludesOne(key string, member any) (bool, error)
- func (service *Service) Keys(template string) ([]string, error)
- func (service *Service) Scan(template string, count int64) ([]string, error)
- func (service *Service) Set(key string, value any, ttl time.Duration) error
- func (service *Service) SetIfNotExists(key string, value any, ttl time.Duration) (bool, error)
- func (service *Service) Size(key string) (int64, error)
- type StreamsConfig
- type StreamsConsumerConfig
- type StreamsMessage
- type StreamsService
- func (service *StreamsService) Ack(ctx context.Context, message StreamsMessage) error
- func (service *StreamsService) Consume(ctx context.Context, group string, stream string, ...) (<-chan StreamsMessage, error)
- func (service *StreamsService) Publish(ctx context.Context, stream string, key string, value any) error
- type StreamsStart
Constants ¶
This section is empty.
Variables ¶
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") )
var ConfigToken = gioc.NewToken("RedisConfig")
var Module = ForFeature()
Module preserves the base Redis module API for applications that do not use an optional feature.
var ServiceInjections = react.InjectFromBase(ConfigToken)
var ServiceToken = gioc.NewToken("RedisService")
var Streams = Feature{ // contains filtered or unexported fields }
Streams enables the lifecycle-owned Redis Streams worker hub.
var StreamsConfigToken = gioc.NewToken("StreamsConfig")
StreamsConfigToken is intentionally separate from ConfigToken. The application must provide it only when redis.ForFeature(redis.Streams) is selected.
var StreamsServiceInjections = []gioc.Token{ ServiceToken, StreamsConfigToken, react.ApplicationContextServiceToken, react.LoggerToken, }
var StreamsServiceToken = gioc.NewToken("StreamsService")
Functions ¶
func ForFeature ¶ added in v0.2.0
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 ProvideStreamsConfig ¶ added in v0.2.0
func ProvideStreamsConfig(config *StreamsConfig) gioc.IProvider
ProvideStreamsConfig makes the application-owned configuration available through StreamsConfigToken.
Types ¶
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.
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) 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) IncludesMany ¶
func (*Service) IncludesOne ¶
func (*Service) SetIfNotExists ¶
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 = "$" )