Documentation
¶
Overview ¶
Package raft implements Timebox persistence semantics using Raft for the replicated commit path, a durable write-ahead log, and a local materialized read store
The core correctness boundary remains the Timebox append contract: aggregate-local optimistic concurrency, atomic event-batch append, and aligned derived index updates. Timebox snapshots are stored through the replicated state machine as accelerative state, not as the authoritative source of truth
Archive lifecycle support is replicated through Raft and stored in the local materialized state
Index ¶
- Constants
- Variables
- func AggregateMetaKey(encodedID string) []byte
- func AggregateMetaPrefix() []byte
- type AggregateMeta
- type ApplyResult
- type ArchiveCommand
- type Backend
- func (b *Backend) Append(reqs ...timebox.AppendRequest) error
- func (b *Backend) Archive(id timebox.AggregateID) error
- func (b *Backend) Close() error
- func (b *Backend) ConsumeArchive(ctx context.Context, h timebox.ArchiveHandler) error
- func (b *Backend) GetAggregateStatus(id timebox.AggregateID) (string, error)
- func (b *Backend) LeaderWithID() (ServerAddress, ServerID)
- func (b *Backend) ListAggregates(typ timebox.ID) ([]timebox.AggregateID, error)
- func (b *Backend) ListAggregatesByStatus(q timebox.StatusQuery) ([]timebox.StatusEntry, error)
- func (b *Backend) ListAggregatesByTag(tag string) ([]timebox.AggregateID, error)
- func (b *Backend) LoadEvents(req timebox.LoadEventsRequest) (*timebox.EventsResult, error)
- func (b *Backend) LoadSnapshot(req timebox.LoadSnapshotRequest) (*timebox.SnapshotRecord, error)
- func (b *Backend) NewStore(cfgs ...timebox.Config) (*timebox.Store, error)
- func (b *Backend) Ready() <-chan struct{}
- func (b *Backend) SaveSnapshot(req timebox.SnapshotRequest) error
- func (b *Backend) State() State
- type Command
- func MakeAppendCommand(proposalID uint64, reqs []timebox.AppendRequest) (Command, error)
- func MakeArchiveCommand(proposalID uint64, ac *ArchiveCommand) Command
- func MakeConsumeArchiveCommand(proposalID uint64, ac *ConsumeArchiveCommand) Command
- func MakeSnapshotCommand(proposalID uint64, sc *SnapshotCommand) Command
- func (c Command) AppendRequests() ([]*timebox.AppendRequest, error)
- func (c Command) ArchiveRequest() (*ArchiveCommand, error)
- func (c Command) ConsumeArchiveRequest() (*ConsumeArchiveCommand, error)
- func (c Command) ProposalID() (uint64, error)
- func (c Command) SnapshotRequest() (*SnapshotCommand, error)
- func (c Command) Type() int
- type Config
- type ConsumeArchiveCommand
- type Server
- type ServerAddress
- type ServerID
- type SnapshotCommand
- type State
Constants ¶
const ( CmdTypeAppend = 0 CmdTypeSnapshot = 1 CmdTypeArchive = 2 CmdTypeConsumeArchive = 3 )
const ( // DefaultLogTailSize is the default hot retained WAL cache size DefaultLogTailSize = 20480 // MinLogTailSize is the smallest allowed hot retained WAL cache size MinLogTailSize = 2048 )
const DefaultApplyTimeout = 10 * time.Second
DefaultApplyTimeout bounds one local proposal round trip
Variables ¶
var ( // ErrUnexpectedApplyResult indicates the FSM returned an unexpected result ErrUnexpectedApplyResult = errors.New("unexpected raft apply result") // ErrCommandTypeUnknown indicates the FSM received an unknown command type ErrCommandTypeUnknown = errors.New("unknown command type") )
var ( // ErrLocalIDRequired indicates the local Raft server ID is required ErrLocalIDRequired = errors.New("raft local ID is required") // ErrDataDirRequired indicates durable local storage must be configured ErrDataDirRequired = errors.New("raft data directory is required") // ErrAddressRequired indicates the Raft TCP listener is required ErrAddressRequired = errors.New("raft address is required") // ErrInvalidAddress indicates a raft address must be a valid host:port ErrInvalidAddress = errors.New("raft address must be a valid host:port") // ErrBootstrapMissingLocalServer indicates the bootstrap voter set must // include the local node ErrBootstrapMissingLocalServer = errors.New( "bootstrap servers must include the local raft ID", ) )
var (
ErrTransportClosed = errors.New("transport closed")
)
Functions ¶
func AggregateMetaKey ¶
AggregateMetaKey returns the metadata key for one encoded aggregate ID
func AggregateMetaPrefix ¶
func AggregateMetaPrefix() []byte
AggregateMetaPrefix returns the key prefix for aggregate metadata
Types ¶
type AggregateMeta ¶
type AggregateMeta struct {
Tags map[string]bool
Status string
CurrentSequence int64
BaseSequence int64
SnapshotSequence int64
StatusAt int64
}
AggregateMeta stores the derived aggregate state needed for reads
type ApplyResult ¶
type ApplyResult struct {
Appends []*timebox.AppendRequest
Error error
}
ApplyResult reports the local outcome of one applied Raft command
type ArchiveCommand ¶
type ArchiveCommand struct {
ID timebox.AggregateID
}
ArchiveCommand carries one Timebox archive mutation through Raft
type Backend ¶
type Backend struct {
// contains filtered or unexported fields
}
Backend applies Timebox writes through Raft and serves reads from local materialized state
func (*Backend) Append ¶
func (b *Backend) Append(reqs ...timebox.AppendRequest) error
Append proposes every append mutation through the local Raft node
func (*Backend) Archive ¶
func (b *Backend) Archive(id timebox.AggregateID) error
Archive archives an aggregate and removes it from active storage
func (*Backend) ConsumeArchive ¶
ConsumeArchive blocks until one archive record is available or ctx is done
func (*Backend) GetAggregateStatus ¶
func (b *Backend) GetAggregateStatus( id timebox.AggregateID, ) (string, error)
GetAggregateStatus returns the current derived status for one aggregate
func (*Backend) LeaderWithID ¶
func (b *Backend) LeaderWithID() (ServerAddress, ServerID)
LeaderWithID returns the current leader address and server ID
func (*Backend) ListAggregates ¶
ListAggregates lists known aggregate IDs of the given type, or of every type when it is empty
func (*Backend) ListAggregatesByStatus ¶
func (b *Backend) ListAggregatesByStatus( q timebox.StatusQuery, ) ([]timebox.StatusEntry, error)
ListAggregatesByStatus lists aggregates matching the query, ordered by status time. Type narrows the status-index scan; KeyPrefix and Through filter each entry
func (*Backend) ListAggregatesByTag ¶
func (b *Backend) ListAggregatesByTag( tag string, ) ([]timebox.AggregateID, error)
ListAggregatesByTag lists aggregates currently indexed by tag
func (*Backend) LoadEvents ¶
func (b *Backend) LoadEvents( req timebox.LoadEventsRequest, ) (*timebox.EventsResult, error)
LoadEvents loads events for one aggregate starting at the requested sequence
func (*Backend) LoadSnapshot ¶
func (b *Backend) LoadSnapshot( req timebox.LoadSnapshotRequest, ) (*timebox.SnapshotRecord, error)
LoadSnapshot returns the latest snapshot and tail events for one aggregate
func (*Backend) Ready ¶
func (b *Backend) Ready() <-chan struct{}
Ready closes once the node is ready to serve leader-directed traffic
func (*Backend) SaveSnapshot ¶
func (b *Backend) SaveSnapshot(req timebox.SnapshotRequest) error
SaveSnapshot proposes one Timebox snapshot mutation through Raft
type Command ¶
type Command []byte
Command is the encoded form of one replicated Timebox mutation
func MakeAppendCommand ¶
func MakeAppendCommand( proposalID uint64, reqs []timebox.AppendRequest, ) (Command, error)
MakeAppendCommand encodes a set of append mutations into a Raft command
func MakeArchiveCommand ¶
func MakeArchiveCommand(proposalID uint64, ac *ArchiveCommand) Command
MakeArchiveCommand encodes one archive mutation into a Raft command
func MakeConsumeArchiveCommand ¶
func MakeConsumeArchiveCommand( proposalID uint64, ac *ConsumeArchiveCommand, ) Command
MakeConsumeArchiveCommand encodes one archive ack into a Raft command
func MakeSnapshotCommand ¶
func MakeSnapshotCommand(proposalID uint64, sc *SnapshotCommand) Command
MakeSnapshotCommand encodes one snapshot mutation into a Raft command
func (Command) AppendRequests ¶
func (c Command) AppendRequests() ([]*timebox.AppendRequest, error)
AppendRequests decodes the append requests from the command payload
func (Command) ArchiveRequest ¶
func (c Command) ArchiveRequest() (*ArchiveCommand, error)
ArchiveRequest decodes an archive request from the command payload
func (Command) ConsumeArchiveRequest ¶
func (c Command) ConsumeArchiveRequest() (*ConsumeArchiveCommand, error)
ConsumeArchiveRequest decodes an archive ack from the command payload
func (Command) ProposalID ¶
ProposalID returns the encoded proposal ID
func (Command) SnapshotRequest ¶
func (c Command) SnapshotRequest() (*SnapshotCommand, error)
SnapshotRequest decodes a snapshot request from the command payload
type Config ¶
type Config struct {
// Local state
LocalID string
DataDir string
// LogTailSize is the hot retained WAL suffix cache size in entries
LogTailSize int
// Cluster identity
Address string
Servers []Server
Publisher timebox.Publisher
}
Config defines one opinionated Raft Backend node
func DefaultConfig ¶
func DefaultConfig() Config
DefaultConfig returns the opinionated defaults for one Raft node
func (Config) LocalServer ¶
LocalServer returns the local server entry derived from this config
type ConsumeArchiveCommand ¶
type ConsumeArchiveCommand struct {
StreamID string
}
ConsumeArchiveCommand carries one archive acknowledgement through Raft
type ServerAddress ¶
type ServerAddress string
ServerAddress is the advertised Raft transport address for one node
type SnapshotCommand ¶
type SnapshotCommand struct {
ID timebox.AggregateID
Data []byte
Sequence int64
TrimEvents bool
}
SnapshotCommand carries one Timebox snapshot mutation through Raft