Documentation
¶
Overview ¶
Package outbox provides transport-neutral durable event publication contracts and an explicit at-least-once dispatcher.
Index ¶
- Variables
- type ClaimRequest
- type Completion
- type Delivery
- type Dispatcher
- type FailureDelay
- type Message
- type MessageSpec
- type Observation
- type Observer
- type Options
- type Publisher
- type Release
- type Result
- type SQLStatements
- type SQLStore
- func (store *SQLStore) Claim(ctx context.Context, request ClaimRequest) (deliveries []Delivery, resultErr error)
- func (store *SQLStore) Complete(ctx context.Context, completion Completion) error
- func (store *SQLStore) Enqueue(ctx context.Context, executor data.Executor, message Message) error
- func (store *SQLStore) Release(ctx context.Context, release Release) error
- type Store
Constants ¶
This section is empty.
Variables ¶
var ErrPublisherPanicked = errors.New("outbox publisher panicked")
ErrPublisherPanicked identifies an observed publisher panic. RunOnce reports it and then re-panics with the original value.
Functions ¶
This section is empty.
Types ¶
type ClaimRequest ¶
ClaimRequest asks a store to atomically lease the oldest available messages.
type Completion ¶
Completion identifies one lease that was published successfully.
type Delivery ¶
type Delivery struct {
// contains filtered or unexported fields
}
Delivery is an immutable leased message returned by a Store.
func NewDelivery ¶
NewDelivery validates and freezes one store lease.
type Dispatcher ¶
type Dispatcher struct {
// contains filtered or unexported fields
}
Dispatcher claims, publishes, and completes durable messages.
func NewDispatcher ¶
func NewDispatcher( store Store, publisher Publisher, options Options, observers ...Observer, ) (*Dispatcher, error)
NewDispatcher validates and freezes one dispatch worker.
type FailureDelay ¶
FailureDelay computes when a failed lease becomes available again.
type Message ¶
type Message struct {
// contains filtered or unexported fields
}
Message is an immutable serialized event prepared for durable storage.
func NewMessage ¶
func NewMessage(spec MessageSpec) (Message, error)
NewMessage validates and freezes one serialized event. IDs are caller-owned idempotency keys and must be unique within the application outbox.
func (Message) ContentType ¶
ContentType returns the normalized payload media type.
func (Message) OccurredAt ¶
OccurredAt returns the UTC event occurrence time.
type MessageSpec ¶
type MessageSpec struct {
ID string
Topic string
Module string
ContentType string
Payload []byte
OccurredAt time.Time
}
MessageSpec is the inspectable input to NewMessage.
type Observation ¶
type Observation struct {
Topic string
Module string
Attempt int
Duration time.Duration
Published bool
Completed bool
Released bool
Err error
Panicked bool
}
Observation contains bounded metadata and no payload or lease receipt.
type Observer ¶
type Observer func(context.Context, Observation)
Observer receives completed delivery attempts synchronously.
type Options ¶
type Options struct {
Owner string
BatchSize int
Lease time.Duration
Clock func() time.Time
FailureDelay FailureDelay
}
Options configures one instance-owned dispatcher.
type Publisher ¶
Publisher sends one message to an external transport. Implementations must use Message.ID as the downstream idempotency key.
type SQLStatements ¶
SQLStatements supplies dialect-owned, fixed SQL for the outbox protocol. Statement text is trusted startup configuration, never request input.
type SQLStore ¶
type SQLStore struct {
// contains filtered or unexported fields
}
SQLStore implements Store using standard database/sql contracts.
func NewSQLStore ¶
func NewSQLStore( executor data.Executor, statements SQLStatements, ) (*SQLStore, error)
NewSQLStore validates and freezes one driver-neutral SQL store. Construction performs no database operation.
func (*SQLStore) Claim ¶
func (store *SQLStore) Claim( ctx context.Context, request ClaimRequest, ) (deliveries []Delivery, resultErr error)
Claim atomically leases messages through the configured statement. Arguments are owner, current time, lease expiry, and limit. Rows must return ID, topic, module, content type, payload, occurrence time, receipt, and attempt.
func (*SQLStore) Complete ¶
func (store *SQLStore) Complete(ctx context.Context, completion Completion) error
Complete removes or marks one published lease using owner and receipt.
type Store ¶
type Store interface {
Enqueue(context.Context, data.Executor, Message) error
Claim(context.Context, ClaimRequest) ([]Delivery, error)
Complete(context.Context, Completion) error
Release(context.Context, Release) error
}
Store owns durable persistence and lease transitions. Enqueue must use the supplied executor so application state and its event can commit atomically.