Versions in this module Expand all Collapse all v0 v0.2.0 Aug 23, 2026 Changes in this version + const MaxClaimDestinations + const MaximumTimestampUnixMicro + const OutboxAppenderToken + const OutboxConfigToken + const OutboxDeliveryStoreToken + const OutboxMaintenanceStoreToken + const OutboxModuleToken + const OutboxPostgresConfigToken + const OutboxReaderToken + const OutboxRedisConfigToken + const OutboxServiceToken + const OutboxStoreToken + var ErrClosed = errors.New("outbox: store closed") + var ErrConflict = errors.New("outbox: conflict") + var ErrDuplicateID = errors.New("outbox: duplicate ID or idempotency key") + var ErrInvalidArgument = errors.New("outbox: invalid argument") + var ErrInvalidTransition = errors.New("outbox: invalid state transition") + var ErrLeaseLost = errors.New("outbox: lease lost") + var ErrNotFound = errors.New("outbox: record not found") + var ErrRedisUnsafeDurability = errors.New("outbox redis: unsafe durability configuration") + var ErrUnsupportedCriteria = errors.New("outbox: unsupported query criteria") + var Postgres = DefineStoreFeature("postgres", PostgresStoreServiceProvider) + var PostgresStoreServiceInjections = []gioc.Token + var Redis = DefineStoreFeature("redis", RedisStoreServiceProvider) + var RedisStoreServiceInjections = []gioc.Token + var ServiceInjections = []gioc.Token + func CanonicalTime(value time.Time) time.Time + func CursorForRecord(record Record, field SortField, direction SortDirection) (string, error) + func EncodeCursor(cursor Cursor) (string, error) + func ForFeature(storage StoreFeature, configs ...Config) *gioc.Module + func ImmutableDigest(record NewRecord) string + func NormalizeQuery(query Query, limits Limits) (Query, Cursor, error) + func PostgresStoreProvider(token gioc.Token, config PostgresConfig) gioc.IProvider + func PostgresStoreServiceProvider(token gioc.Token) gioc.IProvider + func ProvidePostgresConfig(config PostgresConfig) gioc.IProvider + func ProvideRedisConfig(config RedisConfig) gioc.IProvider + func RecordAfterCursor(record Record, cursor Cursor) bool + func RecordSortTime(record Record, field SortField) time.Time + func RedisClusterSlot(key string) uint16 + func RedisStoreProvider(token gioc.Token, config RedisConfig) gioc.IProvider + func RedisStoreServiceProvider(token gioc.Token) gioc.IProvider + func ServiceCapabilityProviders(tokens Tokens) []gioc.IProvider + func SortRecords(records []Record, field SortField, direction SortDirection) + func ValidateDestination(destination string, limits Limits) error + func ValidateID(id ID, limits Limits) error + func ValidateLeaseDuration(field string, duration, maximum time.Duration) error + func ValidateLeaseOwner(owner string, limits Limits) error + func ValidateLeaseToken(token string, limits Limits) error + func ValidateTimestamp(field string, value time.Time) error + type AppendRequest struct + DuplicateMode DuplicateMode + Records []NewRecord + type Backlog struct + Dead int64 + Leased int64 + OldestDueAt *time.Time + Pending int64 + type ClaimRequest struct + Destinations []string + LeaseDuration time.Duration + Limit int + Owner string + RecoveryLimit int + type Config struct + ClaimBatchSize int + DeliveryTimeout time.Duration + ErrorClassifier IErrorClassifier + LeaseDuration time.Duration + LeaseRenewalThreshold time.Duration + Limits Limits + MaximumAttempts int + Owner string + PollMaximumInterval time.Duration + PollMinimumInterval time.Duration + RenewalEnabled bool + Retry ExponentialBackoffConfig + RetryPolicy IRetryPolicy + ShutdownTimeout time.Duration + WorkerCount int + func DefaultConfig() Config + func (config Config) Validate() error + type Cursor struct + Direction SortDirection + ID ID + Micros int64 + Sort SortField + Version int + func DecodeCursor(value string) (Cursor, error) + type DeliveryOutcome string + const OutcomeAmbiguous + const OutcomeRetryable + const OutcomeSuccess + const OutcomeTerminal + func (outcome DeliveryOutcome) Valid() bool + type DeliveryResult struct + type DestinationsConfig struct + Concurrency int + Destinations []string + type DuplicateMode int + const AcceptIdentical + const RejectDuplicate + func (m DuplicateMode) Valid() bool + type ErrorClassifierFunc func(error) DeliveryOutcome + func (f ErrorClassifierFunc) Classify(err error) DeliveryOutcome + type ExponentialBackoff struct + func NewExponentialBackoff(config ExponentialBackoffConfig) (*ExponentialBackoff, error) + func (policy *ExponentialBackoff) Next(attempt int, _ error) (time.Duration, bool) + type ExponentialBackoffConfig struct + Jitter float64 + MaxAttempts int + Maximum time.Duration + Minimum time.Duration + Multiplier float64 + type Failure struct + Code string + Message string + func BoundFailure(failure Failure, limits Limits) Failure + type FieldError struct + Field string + Message string + func (e *FieldError) Error() string + func (e *FieldError) Unwrap() error + type Health struct + Backlog Backlog + DurabilitySafe bool + Message string + Ready bool + StorageAvailable bool + type IAppender interface + Append func(ctx context.Context, records ...NewRecord) ([]Record, error) + type IBacklogReader interface + Backlog func(ctx context.Context) (Backlog, error) + type IBatchAppender interface + AppendBatch func(ctx context.Context, request AppendRequest) ([]Record, error) + type ID string + type IDeliveryStore interface + Acknowledge func(ctx context.Context, lease LeaseRef, result DeliveryResult) error + Claim func(ctx context.Context, request ClaimRequest) ([]Record, error) + DeadLetter func(ctx context.Context, lease LeaseRef, failure Failure) error + Release func(ctx context.Context, lease LeaseRef, availableAt time.Time) error + Renew func(ctx context.Context, lease LeaseRef, until time.Time) error + Retry func(ctx context.Context, lease LeaseRef, retry RetryRequest) error + type IErrorClassifier interface + Classify func(err error) DeliveryOutcome + func DefaultErrorClassifier() IErrorClassifier + type IHealthChecker interface + Health func(ctx context.Context) Health + type IMaintenanceStore interface + Cancel func(ctx context.Context, id ID, reason string) error + Purge func(ctx context.Context, request PurgeRequest) (int, error) + Requeue func(ctx context.Context, id ID, options RequeueOptions) error + Reschedule func(ctx context.Context, id ID, availableAt time.Time) error + type IPostgresDB interface + Begin func(ctx context.Context) (pgx.Tx, error) + Exec func(ctx context.Context, sql string, arguments ...any) (pgconn.CommandTag, error) + Query func(ctx context.Context, sql string, args ...any) (pgx.Rows, error) + QueryRow func(ctx context.Context, sql string, args ...any) pgx.Row + type IPostgresRowScanner interface + Scan func(dest ...any) error + type IReader interface + Find func(ctx context.Context, query Query) (Page, error) + Get func(ctx context.Context, id ID) (Record, error) + type IRetryPolicy interface + Next func(attempt int, failure error) (delay time.Duration, retry bool) + func DefaultRetryPolicy(maxAttempts int) IRetryPolicy + type ISink interface + Deliver func(ctx context.Context, record Record) error + type IStore interface + type LeaseRef struct + ID ID + Owner string + Token string + Version uint64 + type Limits struct + MaxAggregateIDBytes int + MaxAggregateTypeBytes int + MaxAttempts int + MaxBatchSize int + MaxClaimBatchSize int + MaxDestinationBytes int + MaxDestinationWorkers int + MaxErrorCodeBytes int + MaxErrorMessageBytes int + MaxHeaderBytes int + MaxHeaderKeyBytes int + MaxHeaderValueBytes int + MaxHeaders int + MaxIDBytes int + MaxIdempotencyKeyBytes int + MaxLeaseOwnerBytes int + MaxLeaseTokenBytes int + MaxMessageTypeBytes int + MaxOrderingKeyBytes int + MaxPageSize int + MaxPayloadBytes int + MaxPurgeSize int + MaxQueryIDs int + MaxQueryValues int + MaxWorkerCount int + func DefaultLimits() Limits + func (l Limits) Normalized() Limits + type NewRecord struct + AggregateID string + AggregateType string + AvailableAt time.Time + Destination string + Headers map[string]string + ID ID + IdempotencyKey string + MaxAttempts int + MessageType string + OrderingKey string + Payload []byte + func (r NewRecord) Clone() NewRecord + type OperationError struct + Err error + Operation string + func (e *OperationError) Error() string + func (e *OperationError) Unwrap() error + type Page struct + NextCursor string + Records []Record + type PostgresConfig struct + DefaultMaxAttempts int + DuplicateMode DuplicateMode + Limits Limits + MaxLeaseDuration time.Duration + Namespace string + Schema string + Table string + func DefaultPostgresConfig() PostgresConfig + type PostgresMigration struct + Name string + SQL string + Version int + type PostgresStore struct + func NewPostgresStore(db IPostgresDB, config PostgresConfig) (*PostgresStore, error) + func NewPostgresStoreFromDataSource(dataSource *pgxext.DataSource, config PostgresConfig) (*PostgresStore, error) + func (store *PostgresStore) Acknowledge(ctx context.Context, lease LeaseRef, _ DeliveryResult) error + func (store *PostgresStore) Append(ctx context.Context, records ...NewRecord) ([]Record, error) + func (store *PostgresStore) AppendBatch(ctx context.Context, request AppendRequest) ([]Record, error) + func (store *PostgresStore) Backlog(ctx context.Context) (Backlog, error) + func (store *PostgresStore) Bind(tx pgx.Tx) *PostgresTxAppender + func (store *PostgresStore) Cancel(ctx context.Context, id ID, reason string) error + func (store *PostgresStore) Claim(ctx context.Context, request ClaimRequest) ([]Record, error) + func (store *PostgresStore) Close() error + func (store *PostgresStore) DeadLetter(ctx context.Context, lease LeaseRef, failure Failure) error + func (store *PostgresStore) Find(ctx context.Context, query Query) (Page, error) + func (store *PostgresStore) Get(ctx context.Context, id ID) (Record, error) + func (store *PostgresStore) Health(ctx context.Context) Health + func (store *PostgresStore) Migrate(ctx context.Context) error + func (store *PostgresStore) Migrations() []PostgresMigration + func (store *PostgresStore) Purge(ctx context.Context, request PurgeRequest) (int, error) + func (store *PostgresStore) Release(ctx context.Context, lease LeaseRef, availableAt time.Time) error + func (store *PostgresStore) Renew(ctx context.Context, lease LeaseRef, until time.Time) error + func (store *PostgresStore) Requeue(ctx context.Context, id ID, options RequeueOptions) error + func (store *PostgresStore) Reschedule(ctx context.Context, id ID, availableAt time.Time) error + func (store *PostgresStore) Retry(ctx context.Context, lease LeaseRef, retry RetryRequest) error + type PostgresStoreService struct + Logger react.ILogger + func NewPostgresStoreService(injections gioc.Injections) (*PostgresStoreService, error) + func (service *PostgresStoreService) Close() error + type PostgresTxAppender struct + func (appender *PostgresTxAppender) Append(ctx context.Context, records ...NewRecord) ([]Record, error) + func (appender *PostgresTxAppender) AppendBatch(ctx context.Context, request AppendRequest) ([]Record, error) + type PurgeRequest struct + Before time.Time + Limit int + States []State + func NormalizePurgeRequest(request PurgeRequest, limits Limits) (PurgeRequest, error) + type Query struct + AggregateID string + AggregateType string + AvailableAt TimeRange + CreatedAt TimeRange + Cursor string + Destinations []string + Direction SortDirection + IDs []ID + IdempotencyKey string + Limit int + MessageTypes []string + OrderingKey string + Sort SortField + States []State + type Record struct + AggregateID string + AggregateType string + Attempts int + AvailableAt time.Time + CancelledAt *time.Time + ContentDigest string + CreatedAt time.Time + DeadAt *time.Time + DeliveredAt *time.Time + Destination string + Headers map[string]string + ID ID + IdempotencyKey string + LastErrorCode string + LastErrorMessage string + LeaseOwner string + LeaseToken string + LeaseUntil *time.Time + MaxAttempts int + MessageType string + OrderingKey string + Payload []byte + State State + UpdatedAt time.Time + Version uint64 + func PrepareRecord(input NewRecord, now time.Time, defaultMaxAttempts int, limits Limits) (Record, error) + func (r Record) Clone() Record + func (r Record) LeaseRef() LeaseRef + type RedisCompositionRequest struct + Append AppendRequest + ApplyLua string + DomainArguments []any + DomainKeys []string + ValidateLua string + type RedisConfig struct + AllowUnsafeEviction bool + DefaultMaxAttempts int + DuplicateMode DuplicateMode + DurabilityMode RedisDurabilityMode + Limits Limits + MaxAppendEncodedBytes int + MaxClaimResponseBytes int + MaxLeaseDuration time.Duration + Namespace string + RequireNoEviction bool + func DefaultRedisConfig() RedisConfig + type RedisDurabilityMode string + const RedisDurabilityRequireAOF + const RedisDurabilityUnchecked + const RedisDurabilityWarn + type RedisDurabilityReport struct + AOFEnabled bool + AOFFsync string + AOFLastWriteOK bool + Checked bool + EvictionPolicy string + Role string + Warnings []string + type RedisKeys struct + func NewRedisKeys(namespace string) (RedisKeys, error) + func (keys RedisKeys) Cancelled() string + func (keys RedisKeys) Dead() string + func (keys RedisKeys) Delivered() string + func (keys RedisKeys) Idempotency() string + func (keys RedisKeys) Leased() string + func (keys RedisKeys) Namespace() string + func (keys RedisKeys) Pending() string + func (keys RedisKeys) PendingDestinations() string + func (keys RedisKeys) QueryAll() string + func (keys RedisKeys) QueryCancelled() string + func (keys RedisKeys) QueryDead() string + func (keys RedisKeys) QueryDelivered() string + func (keys RedisKeys) QueryDestinations() string + func (keys RedisKeys) QueryLeased() string + func (keys RedisKeys) QueryPending() string + func (keys RedisKeys) QueryState(state State) (string, error) + func (keys RedisKeys) RecordKey(id ID) string + func (keys RedisKeys) Records() string + func (keys RedisKeys) ScriptKeys() []string + func (keys RedisKeys) State(state State) (string, error) + type RedisStore struct + func NewRedisStore(ctx context.Context, client goredis.UniversalClient, config RedisConfig) (*RedisStore, error) + func NewRedisStoreFromService(ctx context.Context, service *reactredis.Service, config RedisConfig) (*RedisStore, error) + func (store *RedisStore) Acknowledge(ctx context.Context, lease LeaseRef, _ DeliveryResult) error + func (store *RedisStore) Append(ctx context.Context, records ...NewRecord) ([]Record, error) + func (store *RedisStore) AppendBatch(ctx context.Context, request AppendRequest) ([]Record, error) + func (store *RedisStore) Backlog(ctx context.Context) (Backlog, error) + func (store *RedisStore) Cancel(ctx context.Context, id ID, reason string) error + func (store *RedisStore) CheckDurability(ctx context.Context) (report RedisDurabilityReport, resultErr error) + func (store *RedisStore) Claim(ctx context.Context, request ClaimRequest) ([]Record, error) + func (store *RedisStore) Close() error + func (store *RedisStore) Compose(ctx context.Context, request RedisCompositionRequest) ([]Record, error) + func (store *RedisStore) DeadLetter(ctx context.Context, lease LeaseRef, failure Failure) error + func (store *RedisStore) Find(ctx context.Context, query Query) (Page, error) + func (store *RedisStore) Get(ctx context.Context, id ID) (Record, error) + func (store *RedisStore) Health(ctx context.Context) Health + func (store *RedisStore) Keys() RedisKeys + func (store *RedisStore) LastDurabilityReport() RedisDurabilityReport + func (store *RedisStore) Purge(ctx context.Context, request PurgeRequest) (int, error) + func (store *RedisStore) Release(ctx context.Context, lease LeaseRef, availableAt time.Time) error + func (store *RedisStore) Renew(ctx context.Context, lease LeaseRef, until time.Time) error + func (store *RedisStore) Requeue(ctx context.Context, id ID, options RequeueOptions) error + func (store *RedisStore) Reschedule(ctx context.Context, id ID, availableAt time.Time) error + func (store *RedisStore) Retry(ctx context.Context, lease LeaseRef, retry RetryRequest) error + type RedisStoreService struct + Logger react.ILogger + func NewRedisStoreService(injections gioc.Injections) (*RedisStoreService, error) + func (service *RedisStoreService) Close() error + type RequeueOptions struct + AvailableAt time.Time + MaxAttempts int + ResetAttempts bool + type RetryRequest struct + AvailableAt time.Time + Failure Failure + type Service struct + ApplicationService *react.ApplicationService + Logger react.ILogger + func NewService(injections gioc.Injections) (*Service, error) + func (service *Service) Append(ctx context.Context, records ...NewRecord) ([]Record, error) + func (service *Service) AppendBatch(ctx context.Context, request AppendRequest) ([]Record, error) + func (service *Service) Destinations() []string + func (service *Service) Done() <-chan struct{} + func (service *Service) Register(sink ISink, config DestinationsConfig) error + func (service *Service) Shutdown(ctx context.Context) error + func (service *Service) String() string + type SinkFunc func(ctx context.Context, record Record) error + func (sink SinkFunc) Deliver(ctx context.Context, record Record) error + type SortDirection string + const SortAscending + const SortDescending + type SortField string + const SortAvailableAt + const SortCreatedAt + type State string + const StateCancelled + const StateDead + const StateDelivered + const StateLeased + const StatePending + func (s State) Terminal() bool + func (s State) Valid() bool + type StoreFeature struct + func DefineStoreFeature(name string, provider StoreProviderFactory) StoreFeature + func (feature StoreFeature) Name() string + type StoreProviderFactory func(token gioc.Token) gioc.IProvider + type TerminalError struct + Err error + func (e *TerminalError) Error() string + func (e *TerminalError) Unwrap() error + type TimeRange struct + From *time.Time + To *time.Time + type Tokens struct + Appender gioc.Token + DeliveryStore gioc.Token + MaintenanceStore gioc.Token + Reader gioc.Token + Service gioc.Token + Store gioc.Token + func ModuleTokens() Tokens + func NewTokens(name string) (Tokens, error)