Documentation
¶
Overview ¶
Package timebox implements an event sourcing toolkit with pluggable persistence backends such as memory, Redis/Valkey, PostgreSQL, and Raft. It couples an append-only event log, snapshots, optimistic concurrency, and append-time indexing into a library that can be embedded into services
Typical usage looks like:
- Open a Backend and create a Store from it
- Define Appliers that fold events into your aggregate state
- Optionally define an Indexer to project current status or tag indexes
- Use an Executor to run Commands that raise events on an Aggregator
- Save snapshots explicitly or let the Executor refresh them while loading
- Query the Store directly for events and aggregate state
The examples/ directory contains a runnable order workflow that exercises the API in a small domain
Index ¶
- Constants
- Variables
- func Configure[T With[T]](defaults T, others ...T) T
- func MakeDecodeAll[Value, Data any](codec Codec[Value, Data]) func([]Data) ([]Value, error)
- func MakeEncodeAll[Value, Data any](codec Codec[Value, Data]) func([]Value) ([]Data, error)
- type AggregateID
- type Aggregator
- type AlwaysReady
- type AppendRequest
- type Applier
- type Appliers
- type ArchiveHandler
- type ArchiveRecord
- type Archiver
- type Backend
- type Codec
- type Command
- type Config
- type Constructor
- type Empty
- type Event
- type EventCodec
- type EventType
- type EventsResult
- type Executor
- type Handler
- type ID
- type Index
- type Indexer
- type LoadEventsRequest
- type LoadSnapshotRequest
- type Message
- type Publisher
- type Queries
- type Schedule
- type ScheduleKey
- type ScheduleVersion
- type ScheduleVersionConflictError
- type SnapshotRecord
- type SnapshotRequest
- type SnapshotResult
- type StatusEntry
- type StatusQuery
- type Store
- func (s *Store) AppendEvents(id AggregateID, atSeq int64, evs []*Event) error
- func (s *Store) Archive(id AggregateID) error
- func (s *Store) Config() Config
- func (s *Store) ConsumeArchive(ctx context.Context, handler ArchiveHandler) error
- func (s *Store) Executor[T any](cons Constructor[T], apps Appliers[T], onSuccess ...SuccessAction[T]) *Executor[T]
- func (s *Store) GetEvents(id AggregateID, fromSeq int64) ([]*Event, error)
- func (s *Store) GetSnapshot(id AggregateID, target any) (*SnapshotResult, error)
- func (s *Store) ListSchedules(through time.Time) ([]*Schedule, error)
- func (s *Store) LoadSchedule(key ScheduleKey) (*Schedule, error)
- func (s *Store) PutSnapshot(id AggregateID, value any, sequence int64) error
- func (s *Store) Ready() <-chan struct{}
- func (s *Store) ScheduleChanges() <-chan struct{}
- func (s *Store) Transact(fn func(*Transaction) error) error
- func (s *Store) WaitReady(ctx context.Context) error
- type SuccessAction
- type Transaction
- func (t *Transaction) CancelSchedule(key ScheduleKey) error
- func (t *Transaction) CancelSchedulePrefix(prefix ScheduleKey) error
- func (t *Transaction) ConsumeSchedule(key ScheduleKey, version ScheduleVersion) error
- func (t *Transaction) Exec[T any](e *Executor[T], id AggregateID, cmd Command[T]) (T, error)
- func (t *Transaction) Schedule(key ScheduleKey, at time.Time, message *Message) error
- type VersionConflictError
- type With
Constants ¶
const ( // DefaultTrimEvents determines whether snapshots trim stored events DefaultTrimEvents = false // DefaultSnapshotRatio is the default EventsSize/SnapshotSize threshold DefaultSnapshotRatio = 1.0 // DefaultMaxRetries is the default optimistic concurrency retry count DefaultMaxRetries = 16 // DefaultCacheSize controls the projection LRU size DefaultCacheSize = 128 )
Variables ¶
var ( // ErrUnexpectedResult indicates data returned in an unexpected shape ErrUnexpectedResult = errors.New("unexpected result") // ErrArchivingDisabled indicates archiving is not enabled ErrArchivingDisabled = errors.New("archiving not enabled for this store") // ErrArchiveRecordMalformed indicates an archive record was malformed ErrArchiveRecordMalformed = errors.New("archive record malformed") // ErrArchiveHandlerMissing indicates a consume call is missing a handler ErrArchiveHandlerMissing = errors.New("archive handler is required") // ErrDuplicateAggregate indicates one append names an aggregate twice ErrDuplicateAggregate = errors.New("aggregate appended twice") )
var ( // JSONEvent encodes complete events as JSON JSONEvent EventCodec[[]byte] = jsonEventCodec{} // BinEvent encodes complete events using the internal binary format BinEvent interface { EventCodec[[]byte] AppendAll(buf []byte, evs []*Event) ([]byte, error) ReadAll(data []byte) ([]*Event, []byte, error) } = binEventCodec{} // EncodeJSONEvents encodes complete events as JSON bytes EncodeJSONEvents = MakeEncodeAll(JSONEvent) // DecodeJSONEvents decodes complete events from JSON bytes DecodeJSONEvents = MakeDecodeAll(JSONEvent) // EncodeBinEvents encodes complete events to the internal binary format EncodeBinEvents = MakeEncodeAll(BinEvent) // DecodeBinEvents decodes complete events from the internal binary format DecodeBinEvents = MakeDecodeAll(BinEvent) )
var ( // ErrInvalidMaxRetries indicates MaxRetries is below the allowed range ErrInvalidMaxRetries = errors.New("max retries must be > 0") // ErrInvalidCacheSize indicates CacheSize is below the allowed range ErrInvalidCacheSize = errors.New("cache size must be > 0") // ErrInvalidSnapshotRatio indicates SnapshotRatio is below allowed range ErrInvalidSnapshotRatio = errors.New("snapshot ratio must be >= 0") )
var ( // ErrScheduleKeyRequired indicates a schedule key was empty ErrScheduleKeyRequired = errors.New("schedule key is required") // ErrScheduleMessageRequired indicates a schedule has no deferred message ErrScheduleMessageRequired = errors.New("schedule message is required") // ErrScheduleMessageTypeRequired indicates an empty deferred message type ErrScheduleMessageTypeRequired = errors.New( "schedule message type is required", ) )
var ( // ErrStoreMismatch indicates a Transaction was given an Executor bound to // a different Store ErrStoreMismatch = errors.New("executor belongs to a different store") // ErrAggregateTypeConflict indicates one aggregate was joined twice under // conflicting state types ErrAggregateTypeConflict = errors.New( "aggregate joined with conflicting state types", ) )
var ( // ErrInvalidAggregateID indicates an encoded AggregateID did not decode to // a type and a key ErrInvalidAggregateID = errors.New( "aggregate id must have a type and a key", ) )
var ( // ErrMaxRetriesExceeded indicates optimistic concurrency retries were // exhausted while attempting to persist events ErrMaxRetriesExceeded = errors.New("max retries exceeded") )
Functions ¶
func Configure ¶
func Configure[T With[T]](defaults T, others ...T) T
Configure overlays each supplied value on top of defaults in order
func MakeDecodeAll ¶
MakeDecodeAll returns a batch decoder for the provided codec
func MakeEncodeAll ¶
MakeEncodeAll returns a batch encoder for the provided codec
Types ¶
type AggregateID ¶
type AggregateID struct {
Type ID
Key ID
}
AggregateID identifies an aggregate by type and key ("order", "123")
func NewAggregateID ¶
func NewAggregateID(typ, key ID) AggregateID
NewAggregateID builds an AggregateID from its type and key
func NewAggregateType ¶
func NewAggregateType(typ ID) AggregateID
NewAggregateType builds the AggregateID of the only aggregate of a type
func (AggregateID) MarshalJSON ¶
MarshalJSON encodes the AggregateID as a type and key pair
func (AggregateID) String ¶
func (id AggregateID) String() string
String returns a human-readable AggregateID representation
func (*AggregateID) UnmarshalJSON ¶
UnmarshalJSON decodes the AggregateID from a type and key pair
type Aggregator ¶
type Aggregator[T any] struct { // contains filtered or unexported fields }
Aggregator maintains aggregate state for a command and tracks events raised through Raise. It is not safe for concurrent use
func (*Aggregator[_]) ID ¶
func (a *Aggregator[_]) ID() AggregateID
ID returns the aggregate's identifier
func (*Aggregator[_]) NextSequence ¶
func (a *Aggregator[_]) NextSequence() int64
NextSequence returns the next sequence number that will be assigned to a new event
func (*Aggregator[T]) OnSuccess ¶
func (a *Aggregator[T]) OnSuccess(fn SuccessAction[T])
OnSuccess registers an action to run after Executor.Exec persists the raised events successfully
func (*Aggregator[T]) Raise ¶
func (a *Aggregator[T]) Raise[V any](typ EventType, value V) error
Raise marshals the value and enqueues a new event on the Aggregator
func (*Aggregator[_]) Transaction ¶
func (a *Aggregator[_]) Transaction() *Transaction
Transaction returns the Transaction this Aggregator's events commit in, so code holding only an Aggregator can enlist further aggregates
func (*Aggregator[T]) Value ¶
func (a *Aggregator[T]) Value() T
Value returns the aggregate's current state
type AlwaysReady ¶
type AlwaysReady struct{}
AlwaysReady can be embedded in a Backend to provide a Ready() implementation for backends that are immediately available
func (AlwaysReady) Ready ¶
func (AlwaysReady) Ready() <-chan struct{}
Ready returns a pre-closed channel, indicating immediate readiness
type AppendRequest ¶
type AppendRequest struct {
StatusAt time.Time
Status *string
Tags map[string]bool
ID AggregateID
Events []*Event
ExpectedSequence int64
TrimEvents bool
}
AppendRequest contains primitive inputs required for an atomic append
type Applier ¶
Applier applies an event to an aggregate state, returning the new state
func MakeApplier ¶
MakeApplier wraps a strongly typed applier that receives the event payload value and returns an Applier that works with Event
type ArchiveHandler ¶
type ArchiveHandler func(context.Context, *ArchiveRecord) error
ArchiveHandler handles a single archive record
type ArchiveRecord ¶
type ArchiveRecord struct {
StreamID string
AggregateID AggregateID
SnapshotData json.RawMessage
Events []*Event
SnapshotSequence int64
}
ArchiveRecord stores stream metadata and aggregate artifacts
type Archiver ¶
type Archiver interface {
// Archive moves an aggregate's persisted artifacts into archive storage
Archive(id AggregateID) error
// ConsumeArchive blocks until one archive record is available or ctx
// is done
ConsumeArchive(ctx context.Context, handler ArchiveHandler) error
}
Archiver provides optional archive lifecycle support for Store
type Backend ¶
type Backend interface {
io.Closer
Queries
// Ready reports when the backend can serve requests
Ready() <-chan struct{}
// Append atomically appends every distinctly named request if each
// expected sequence still matches. It returns a VersionConflictError
// naming the first request whose sequence check fails
Append(...AppendRequest) error
// LoadEvents loads raw persisted events starting at fromSeq
LoadEvents(LoadEventsRequest) (*EventsResult, error)
// LoadSnapshot loads the raw snapshot and any trailing raw events
LoadSnapshot(LoadSnapshotRequest) (*SnapshotRecord, error)
// SaveSnapshot stores raw snapshot data using the supplied semantics
SaveSnapshot(SnapshotRequest) error
}
Backend provides the low-level primitives Store uses to implement Store semantics, along with its query operations
type Codec ¶
type Codec[Value, Data any] interface { // Encode encodes a value to the codec data format Encode(Value) (Data, error) // Decode decodes codec data to a value Decode(Data) (Value, error) // Append appends an encoded value to a buffer Append(Data, Value) (Data, error) // Read reads a value from a buffer and returns the remainder Read(Data) (Value, Data, error) }
Codec encodes and decodes values
type Command ¶
type Command[T any] func(T, *Aggregator[T]) error
Command is user code that inspects state and raises events on an Aggregator. Returning an error aborts the operation
type Config ¶
type Config struct {
Indexer Indexer
SnapshotRatio float64
MaxRetries int
CacheSize int
TrimEvents bool
}
Config configures Store behavior
func DefaultConfig ¶
func DefaultConfig() Config
DefaultConfig returns a Config populated with sensible defaults
type Constructor ¶ added in v0.2.0
type Constructor[T any] func() T
Constructors instantiate initial aggregate state
type Event ¶
type Event struct {
Timestamp time.Time `json:"timestamp"`
Message
Sequence int64 `json:"sequence"`
Raised bool `json:"-"`
}
Event is a recorded Message with its sequence and timestamp
type EventCodec ¶
EventCodec encodes and decodes complete events
type EventsResult ¶
EventsResult contains raw persisted events and the sequence to assign to the first event in the slice
type Executor ¶
type Executor[T any] struct { // contains filtered or unexported fields }
Executor orchestrates loading aggregate state, executing commands, and persisting resulting events with optimistic retries
func (*Executor[T]) Exec ¶
Exec loads the aggregate state, executes the command, and persists raised events. It retries on version conflicts up to MaxRetries
func (*Executor[T]) SaveSnapshot ¶
SaveSnapshot forces an immediate snapshot save for the given Aggregate
type Handler ¶
Handler processes a single Message
func MakeDispatcher ¶
MakeDispatcher routes messages to handlers by type, ignoring unmatched types
type ID ¶
type ID string
ID is a single component of an AggregateID
const SingletonKey ID = "_"
SingletonKey is the Key of an aggregate that is the only one of its type
type Index ¶
type Index struct {
// Status represents the resultant aggregate status. nil means no
// status change, and "" clears any prior status
Status *string `json:"status,omitempty"`
// Tags updates aggregate tag membership. true adds and false removes
Tags map[string]bool `json:"tags,omitempty"`
}
Index stores optional projection metadata derived from an event
type LoadEventsRequest ¶
LoadEventsRequest contains primitive inputs required for an event load
type LoadSnapshotRequest ¶
type LoadSnapshotRequest struct {
ID AggregateID
TrimEvents bool
}
LoadSnapshotRequest contains primitive inputs required for a snapshot load
type Message ¶ added in v0.3.0
type Message struct {
Type EventType `json:"type"`
AggregateID AggregateID `json:"aggregate_id"`
Data json.RawMessage `json:"data,omitempty"`
// contains filtered or unexported fields
}
Message is a typed payload associated with an aggregate
type Publisher ¶ added in v0.2.0
type Publisher func(...*Event)
Publisher reports committed events
func (Publisher) PublishAppends ¶ added in v0.2.0
func (p Publisher) PublishAppends(reqs ...AppendRequest)
PublishAppends publishes the requests' events as one batch
type Queries ¶
type Queries interface {
// ListAggregates lists aggregate IDs of the provided type, or of
// every type when it is empty
ListAggregates(typ ID) ([]AggregateID, error)
// GetAggregateStatus loads the current indexed status for an aggregate
GetAggregateStatus(id AggregateID) (string, error)
// ListAggregatesByStatus lists aggregates matching the query, ordered
// by status time
ListAggregatesByStatus(StatusQuery) ([]StatusEntry, error)
// ListAggregatesByTag lists aggregates currently indexed by tag
ListAggregatesByTag(tag string) ([]AggregateID, error)
}
Queries provides aggregate and index query operations
type Schedule ¶ added in v0.3.0
type Schedule struct {
Message *Message
At time.Time
Key ScheduleKey
Version ScheduleVersion
}
Schedule is one durable deferred message and its delivery metadata
type ScheduleKey ¶ added in v0.3.0
type ScheduleKey string
ScheduleKey identifies one replaceable deferred message
type ScheduleVersion ¶ added in v0.3.0
type ScheduleVersion int64
ScheduleVersion identifies one incarnation within a schedule aggregate
type ScheduleVersionConflictError ¶ added in v0.3.0
type ScheduleVersionConflictError struct {
Key ScheduleKey
ExpectedVersion ScheduleVersion
}
ScheduleVersionConflictError indicates a stale schedule incarnation
func (*ScheduleVersionConflictError) Error ¶ added in v0.3.0
func (s *ScheduleVersionConflictError) Error() string
Error describes a stale schedule incarnation
type SnapshotRecord ¶
type SnapshotRecord struct {
Data json.RawMessage
Events []*Event
Sequence int64
}
SnapshotRecord contains raw snapshot data and any raw trailing events
type SnapshotRequest ¶
SnapshotRequest contains primitive inputs required for a snapshot save
type SnapshotResult ¶
type SnapshotResult struct {
AdditionalEvents []*Event
NextSequence int64
SnapshotSize int
EventsSize int
}
SnapshotResult holds the loaded snapshot, the sequence at which it was taken, and any events that need to be applied after it
type StatusEntry ¶
StatusEntry holds an aggregate ID and the time it entered a status
type StatusQuery ¶ added in v0.3.0
StatusQuery selects aggregates by status, optionally narrowed by aggregate type, key prefix, and a latest status time. Through is inclusive; a zero Through, Type, or KeyPrefix imposes no bound
type Store ¶
type Store struct {
Queries
// contains filtered or unexported fields
}
Store persists, queries, and snapshots aggregate events
func (*Store) AppendEvents ¶
AppendEvents atomically appends events for an aggregate if the expected sequence matches the current log sequence
func (*Store) ConsumeArchive ¶
func (s *Store) ConsumeArchive( ctx context.Context, handler ArchiveHandler, ) error
ConsumeArchive reads one archive record and invokes handler
func (*Store) Executor ¶
func (s *Store) Executor[T any]( cons Constructor[T], apps Appliers[T], onSuccess ...SuccessAction[T], ) *Executor[T]
Executor constructs an Executor bound to a Store with the given appliers and state constructor
func (*Store) GetSnapshot ¶
func (s *Store) GetSnapshot( id AggregateID, target any, ) (*SnapshotResult, error)
GetSnapshot loads the latest snapshot into target and returns any events stored after the snapshot sequence
func (*Store) ListSchedules ¶ added in v0.3.0
ListSchedules lists active schedules through the provided time, or all schedules when through is zero
func (*Store) LoadSchedule ¶ added in v0.3.0
func (s *Store) LoadSchedule(key ScheduleKey) (*Schedule, error)
LoadSchedule loads the active schedule for a key
func (*Store) PutSnapshot ¶
PutSnapshot saves a snapshot value and sequence if the provided sequence is newer than any stored snapshot
func (*Store) Ready ¶
func (s *Store) Ready() <-chan struct{}
Ready reports when the underlying Backend can serve requests
func (*Store) ScheduleChanges ¶ added in v0.3.0
func (s *Store) ScheduleChanges() <-chan struct{}
ScheduleChanges reports coalesced local schedule commit notifications
type SuccessAction ¶
SuccessAction receives the Aggregator's final value after Executor.Exec succeeds, as well as the Events persisted by that execution
type Transaction ¶
type Transaction struct {
// contains filtered or unexported fields
}
Transaction collects the append intent of several Aggregators and commits it as a single atomic append. It is not safe for concurrent use
func (*Transaction) CancelSchedule ¶ added in v0.3.0
func (t *Transaction) CancelSchedule(key ScheduleKey) error
CancelSchedule removes the current schedule for a key
func (*Transaction) CancelSchedulePrefix ¶ added in v0.3.0
func (t *Transaction) CancelSchedulePrefix(prefix ScheduleKey) error
CancelSchedulePrefix removes active schedules whose keys share prefix
func (*Transaction) ConsumeSchedule ¶ added in v0.3.0
func (t *Transaction) ConsumeSchedule( key ScheduleKey, version ScheduleVersion, ) error
ConsumeSchedule conditionally removes one observed schedule incarnation
func (*Transaction) Exec ¶
func (t *Transaction) Exec[T any]( e *Executor[T], id AggregateID, cmd Command[T], ) (T, error)
Exec runs cmd against the aggregate and enlists the resulting events in the Transaction, continuing any Aggregator already joined for it. The returned value holds only if the Transaction commits
func (*Transaction) Schedule ¶ added in v0.3.0
func (t *Transaction) Schedule( key ScheduleKey, at time.Time, message *Message, ) error
Schedule creates or replaces a durable deferred message
type VersionConflictError ¶
type VersionConflictError struct {
NewEvents []*Event
ID AggregateID
ExpectedSequence int64
ActualSequence int64
}
VersionConflictError is returned when AppendEvents encounters a sequence mismatch. NewEvents contains the conflicting events
func (*VersionConflictError) Error ¶
func (e *VersionConflictError) Error() string
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
internal
|
|
|
check
Package check enforces append contract rules every backend shares
|
Package check enforces append contract rules every backend shares |
|
Package raft implements Timebox persistence semantics using Raft for the replicated commit path, a durable write-ahead log, and a local materialized read store
|
Package raft implements Timebox persistence semantics using Raft for the replicated commit path, a durable write-ahead log, and a local materialized read store |
|
Package scheduler delivers durable Timebox messages when they become due
|
Package scheduler delivers durable Timebox messages when they become due |
