timebox

package module
v0.3.0 Latest Latest
Warning

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

Go to latest
Published: Oct 7, 2026 License: MIT Imports: 14 Imported by: 2

README

Timebox

Build Status Code Coverage Maintainability GitHub

Timebox is a small, opinionated event sourcing library for Go with append-only event storage, optimistic concurrency, snapshotting, and durable message scheduling. It supports Redis, PostgreSQL, or Raft as its persistence backend.

Documentation

  • Getting started explains how to install Timebox and execute a command.
  • Order tutorial builds an application using aggregates and transactions.
  • Aggregates and events covers identities, payloads, and appliers.
  • Transactions shows how to commit changes across aggregates.
  • Scheduler explains how to deliver messages at a durable deadline.
  • Snapshots explains how to load state without replaying its full history.
  • Indexing covers queries by aggregate status and tags.
  • Archiving shows how to remove finished aggregates from live storage and export their records.
  • Backends helps you choose and configure persistence.
  • Production patterns covers the structure of a long-running service.

The order example shows an aggregate lifecycle using Redis.

Status

Work in progress. Not ready for production use.

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

View Source
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

View Source
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")
)
View Source
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)
)
View Source
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")
)
View Source
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",
	)
)
View Source
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",
	)
)
View Source
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",
	)
)
View Source
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

func MakeDecodeAll[Value, Data any](
	codec Codec[Value, Data],
) func([]Data) ([]Value, error)

MakeDecodeAll returns a batch decoder for the provided codec

func MakeEncodeAll

func MakeEncodeAll[Value, Data any](
	codec Codec[Value, Data],
) func([]Value) ([]Data, error)

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

func (id AggregateID) MarshalJSON() ([]byte, error)

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

func (id *AggregateID) UnmarshalJSON(data []byte) error

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

type Applier[T any] func(T, *Event) T

Applier applies an event to an aggregate state, returning the new state

func MakeApplier

func MakeApplier[T, Data any](fn func(T, *Event, Data) T) Applier[T]

MakeApplier wraps a strongly typed applier that receives the event payload value and returns an Applier that works with Event

type Appliers

type Appliers[T any] map[EventType]Applier[T]

Appliers is a map of EventType to Applier for a given aggregate

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

func (Config) Validate

func (cfg Config) Validate() error

Validate reports whether the configuration contains invalid values

func (Config) With

func (cfg Config) With(other Config) Config

With overlays the non-zero values from other onto cfg

type Constructor added in v0.2.0

type Constructor[T any] func() T

Constructors instantiate initial aggregate state

type Empty added in v0.3.0

type Empty struct{}

Empty is a payload with no fields

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

type EventCodec[Data any] interface {
	Codec[*Event, Data]
}

EventCodec encodes and decodes complete events

type EventType

type EventType string

EventType identifies the kind of an Event or Message

type EventsResult

type EventsResult struct {
	Events        []*Event
	StartSequence int64
}

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

func (e *Executor[T]) Exec(id AggregateID, cmd Command[T]) (T, error)

Exec loads the aggregate state, executes the command, and persists raised events. It retries on version conflicts up to MaxRetries

func (*Executor[T]) Get

func (e *Executor[T]) Get(id AggregateID) (T, error)

Get returns the current aggregate state

func (*Executor[T]) SaveSnapshot

func (e *Executor[T]) SaveSnapshot(id AggregateID) error

SaveSnapshot forces an immediate snapshot save for the given Aggregate

type Handler

type Handler func(*Message) error

Handler processes a single Message

func MakeDispatcher

func MakeDispatcher(handlers map[EventType]Handler) Handler

MakeDispatcher routes messages to handlers by type, ignoring unmatched types

func MakeHandler

func MakeHandler[T any](fn func(msg *Message, data T) error) Handler

MakeHandler decodes message data into the provided type before invoking fn

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 Indexer

type Indexer func([]*Event) []*Index

Indexer derives projection metadata for an event batch

type LoadEventsRequest

type LoadEventsRequest struct {
	ID         AggregateID
	FromSeq    int64
	TrimEvents bool
}

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

func (*Message) GetValue added in v0.3.0

func (m *Message) GetValue[T any]() (T, error)

GetValue unmarshals the message data into the requested type. It reuses a cached value when the requested type matches, and otherwise unmarshals without replacing the cached type. This is safe for concurrent access

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

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

type SnapshotRequest struct {
	ID         AggregateID
	Data       []byte
	Sequence   int64
	TrimEvents bool
}

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

type StatusEntry struct {
	Timestamp time.Time
	ID        AggregateID
}

StatusEntry holds an aggregate ID and the time it entered a status

type StatusQuery added in v0.3.0

type StatusQuery struct {
	Through   time.Time
	Status    string
	Type      ID
	KeyPrefix ID
}

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 NewStore

func NewStore(b Backend, cfgs ...Config) (*Store, error)

NewStore creates a Store backed by the supplied Backend

func (*Store) AppendEvents

func (s *Store) AppendEvents(id AggregateID, atSeq int64, evs []*Event) error

AppendEvents atomically appends events for an aggregate if the expected sequence matches the current log sequence

func (*Store) Archive

func (s *Store) Archive(id AggregateID) error

Archive moves aggregate artifacts to persistent archive storage

func (*Store) Config

func (s *Store) Config() Config

Config returns the Store configuration

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) GetEvents

func (s *Store) GetEvents(id AggregateID, fromSeq int64) ([]*Event, error)

GetEvents returns all events for an aggregate starting at fromSeq

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

func (s *Store) ListSchedules(through time.Time) ([]*Schedule, error)

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

func (s *Store) PutSnapshot(id AggregateID, value any, sequence int64) error

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

func (*Store) Transact

func (s *Store) Transact(fn func(*Transaction) error) error

Transact runs fn and commits every aggregate joined through Transaction.Exec as one atomic append. It retries fn on version conflict up to MaxRetries. An error returned from fn discards the transaction

func (*Store) WaitReady

func (s *Store) WaitReady(ctx context.Context) error

WaitReady blocks until the underlying Backend can serve requests

type SuccessAction

type SuccessAction[T any] func(T, []*Event)

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

type With

type With[T any] interface {
	With(T) T
}

With is a generic overlay contract used by Configure

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

Jump to

Keyboard shortcuts

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