spool

package
v0.7.0 Latest Latest
Warning

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

Go to latest
Published: Oct 1, 2026 License: Apache-2.0 Imports: 22 Imported by: 0

Documentation

Overview

Package spool is the writer's durable record store, the node agent queue, and directory locks.

Index

Constants

View Source
const (
	DefaultCapacityBytes int64   = 10 << 30
	DefaultWindowBytes   int64   = 8 << 20
	DefaultSegmentBytes  int64   = 64 << 20
	DefaultCoalesceAt    float64 = 0.90
)

Defaults applied to zero Options fields (PRD 8.2).

Variables

View Source
var (
	ErrClosed        = errors.New("spool: closed")
	ErrHalted        = errors.New("spool: writer halted")
	ErrNoEpoch       = errors.New("spool: no open epoch")
	ErrEpochSealed   = errors.New("spool: epoch sealed by discard; open a new epoch")
	ErrEpochMismatch = errors.New("spool: not the current epoch")
	ErrRecordsRemain = errors.New("spool: records of the current epoch remain spooled")
	ErrDivergence    = errors.New("spool: divergence")
	ErrNotSpooled    = errors.New("spool: record not spooled")
	ErrTxDone        = errors.New("spool: transaction finished")
)
View Source
var ErrCorrupt = errors.New("spool: corrupt data")

ErrCorrupt reports damaged data, or metadata that points at a missing or damaged body.

View Source
var ErrFull = errors.New("spool: queue full")

ErrFull reports that an append would exceed the queue capacity.

View Source
var ErrLocked = errors.New("spool: directory is locked")

ErrLocked reports that the directory lock is held, by another process or another open in this one.

View Source
var ErrRebaselineRequired = errors.New("spool: capacity exceeded after full coalescing; rebaseline required")

ErrRebaselineRequired: the fully coalesced never-transmitted suffix exceeds capacity (SPEC 8.7).

Functions

func LockDir

func LockDir(dir string) (unlock func() error, err error)

LockDir takes an exclusive non-blocking flock on dir/LOCK, creating dir if needed.

Types

type AppendOption

type AppendOption func(*appendOpts)

AppendOption modifies an append.

func WithCursor

func WithCursor(key string, value uint64) AppendOption

WithCursor persists an idempotency cursor atomically with the record.

type ClientStore

type ClientStore struct{ S *Spool }

ClientStore adapts a Spool to the session client's Store interface.

func (ClientStore) Commit

func (a ClientStore) Commit(epoch protocol.EpochID, seq uint64, chainHash protocol.Hash) error

func (ClientStore) DiscardAbove

func (a ClientStore) DiscardAbove(seq uint64) error

func (ClientStore) Do

func (a ClientStore) Do(fn func(tx client.Tx) error) error

func (ClientStore) Entries

func (a ClientStore) Entries(fromSeq uint64) []*client.Entry

func (ClientStore) Epoch

func (a ClientStore) Epoch() (client.EpochState, bool)

func (ClientStore) Halted

func (a ClientStore) Halted() (string, bool)

func (ClientStore) Identity

func (a ClientStore) Identity() client.Identity

func (ClientStore) Incarnation

func (a ClientStore) Incarnation() uint64

func (ClientStore) LastCommitted

func (a ClientStore) LastCommitted() (protocol.ChainPoint, bool)

func (ClientStore) MarkRegistered

func (a ClientStore) MarkRegistered() error

func (ClientStore) MarkTransmitted

func (a ClientStore) MarkTransmitted(seqs ...uint64) error

func (ClientStore) Notify

func (a ClientStore) Notify() <-chan struct{}

func (ClientStore) OpenEpoch

func (a ClientStore) OpenEpoch(reason string, prev *protocol.EpochID, prevHead *uint64) (protocol.EpochID, error)

func (ClientStore) SetHalted

func (a ClientStore) SetHalted(code string) error

func (ClientStore) SetIdentity

func (a ClientStore) SetIdentity(id client.Identity) error

func (ClientStore) WriterID

func (a ClientStore) WriterID() protocol.WriterID

type ClientTx

type ClientTx struct{ T *Tx }

ClientTx adapts a spool transaction to client.Tx and exposes cursor-bearing appends.

func (ClientTx) Append

func (a ClientTx) Append(t protocol.RecordType, build func(env protocol.Envelope) (*protocol.Record, error)) (*client.Entry, error)

func (ClientTx) AppendWithCursor

func (a ClientTx) AppendWithCursor(t protocol.RecordType, build func(env protocol.Envelope) (*protocol.Record, error), key string, value uint64) (*client.Entry, error)

AppendWithCursor appends and persists an idempotency cursor atomically with the record.

type Entry

type Entry struct {
	Seq       uint64
	Type      protocol.RecordType
	State     RecordState
	Bytes     []byte // exact record bytes
	Hash      protocol.Hash
	ChainHash protocol.Hash
}

Entry is one spooled record.

type EpochState

type EpochState struct {
	ID         protocol.EpochID  `json:"id"`
	TargetID   string            `json:"target_id"`
	OpenReason string            `json:"open_reason"`
	PrevEpoch  *protocol.EpochID `json:"prev_epoch,omitempty"`
	PrevHead   *uint64           `json:"prev_head,omitempty"`
	Registered bool              `json:"registered"`
	OpenedAt   time.Time         `json:"opened_at"`
	Sealed     bool              `json:"sealed"` // records above SealedAt were discarded; no further appends
	SealedAt   uint64            `json:"sealed_at"`
	Chain      protocol.Chain    `json:"chain"` // Head is the highest sequence ever assigned
}

EpochState is the current epoch and its chain position.

type Halt

type Halt struct {
	Code    string    `json:"code"`
	Message string    `json:"message"`
	At      time.Time `json:"at"`
}

Halt records a rejection that forbids writing until an operator clears it.

type Identity

type Identity struct {
	TargetID     string `json:"target_id"`
	TargetType   string `json:"target_type"`
	Credential   string `json:"credential"`
	CredentialID string `json:"credential_id"`
	MachineID    string `json:"machine_id"`
}

Identity is the enrolled identity stored with the spool.

type Options

type Options struct {
	Dir           string  // coordinator PVC mount or /var/lib/exitmesh/spool
	CapacityBytes int64   // hard budget for record bodies
	WindowBytes   int64   // transmitted-unconfirmed window
	SegmentBytes  int64   // segment file size before rollover
	CoalesceAt    float64 // fraction of capacity that triggers relief
	Clock         func() time.Time
	// CommitFault injects faults for tests: a non-nil error fails a transaction's metadata commit after its bodies were written.
	CommitFault func() error
}

Options configures a Spool.

type Queue

type Queue struct {
	// contains filtered or unexported fields
}

Queue is the node agent's durable, segmented, append-only queue.

func OpenQueue

func OpenQueue(dir string, capacityBytes int64) (*Queue, error)

OpenQueue locks dir and recovers the queue, dropping a torn tail frame left by a crash.

func (*Queue) Ack

func (q *Queue) Ack(seq uint64) error

Ack durably drops every item with sequence at most seq.

func (*Queue) Append

func (q *Queue) Append(b []byte) (uint64, error)

Append writes b as the next item and fsyncs it before returning its sequence.

func (*Queue) Close

func (q *Queue) Close() error

Close releases the queue and its lock.

func (*Queue) Err

func (q *Queue) Err() error

Err returns the first persistent queue fault, if any.

func (*Queue) ID

func (q *Queue) ID() string

ID is the random identity created with the queue; a queue recreated from nothing gets a new one.

func (*Queue) Peek

func (q *Queue) Peek(fromSeq uint64, maxBytes int) []QueueItem

Peek returns unacknowledged items from fromSeq on, up to maxBytes of data but at least one.

func (*Queue) Usage

func (q *Queue) Usage() QueueUsage

Usage reports queue occupancy.

type QueueItem

type QueueItem struct {
	Seq  uint64
	Data []byte
}

QueueItem is one queued payload and its sequence.

type QueueUsage

type QueueUsage struct {
	Bytes    int64 // segment file bytes, counted against Capacity
	Capacity int64
	Items    uint64 // appended and not acknowledged
	Acked    uint64
	Next     uint64 // sequence the next append receives
}

QueueUsage reports queue occupancy.

type RecordState

type RecordState uint8

RecordState is the delivery state of a spooled record (SPEC 8.1). Committed records are deleted.

const (
	NeverTransmitted       RecordState = 1
	TransmittedUnconfirmed RecordState = 2
)

func (RecordState) String

func (s RecordState) String() string

type Relief

type Relief struct {
	BytesBefore       int64
	BytesAfter        int64
	EvidenceEvicted   int             // records whose findings kept only one evidence sample
	SamplesCompacted  int             // records whose finding samples were compacted to counts
	RecordsCoalesced  int             // records folded into range records
	Unavailable       []protocol.Span // interiors of the range records written, ascending
	SegmentsCompacted int
	Exhausted         bool // every step ran and usage is still above the relief target
}

Relief reports one pressure relief pass (PRD S13).

type Spool

type Spool struct {
	// contains filtered or unexported fields
}

Spool is the writer's durable record store; only Tx methods may be called from inside Do.

func Open

func Open(opts Options) (*Spool, error)

Open locks Dir, creates the writer ID once, durably increments the incarnation, and recovers.

func (*Spool) ClearHalt

func (s *Spool) ClearHalt() error

ClearHalt is the operator path that allows writing again.

func (*Spool) Close

func (s *Spool) Close() error

Close releases the spool and its lock. It must not be called from inside Do.

func (*Spool) Commit

func (s *Spool) Commit(epoch protocol.EpochID, seq uint64, chainHash protocol.Hash) error

Commit deletes every record at or below a committed head, or returns ErrDivergence (SPEC 8.4 step 2).

func (*Spool) Cursor

func (s *Spool) Cursor(key string) uint64

Cursor returns an idempotency cursor persisted with WithCursor, or 0.

func (*Spool) Cursors

func (s *Spool) Cursors(prefix string) map[string]uint64

Cursors returns the persisted idempotency cursors whose keys start with prefix.

func (*Spool) DiscardAbove

func (s *Spool) DiscardAbove(seq uint64) error

DiscardAbove deletes records above seq and seals the epoch, whose sequences are never reissued (SPEC 8.6).

func (*Spool) Do

func (s *Spool) Do(fn func(tx *Tx) error) error

Do runs fn under the sequence lock, commits its appends durably, then relieves pressure if needed.

func (*Spool) Entries

func (s *Spool) Entries(fromSeq uint64) []*Entry

Entries pages spooled records from fromSeq in chain order (at most four windows, at least one) and pins them against rewrites.

func (*Spool) Epoch

func (s *Spool) Epoch() (EpochState, bool)

Epoch returns the current epoch, if one is open.

func (*Spool) Err

func (s *Spool) Err() error

Err returns the first storage fault or read-side corruption seen since Open, if any.

func (*Spool) Halted

func (s *Spool) Halted() (Halt, bool)

func (*Spool) Identity

func (s *Spool) Identity() Identity

func (*Spool) Incarnation

func (s *Spool) Incarnation() uint64

func (*Spool) KV

func (s *Spool) KV(bucket string) kv.Store

KV returns a durable store over its own bucket in the spool metadata database.

func (*Spool) LastCommitted

func (s *Spool) LastCommitted() (protocol.ChainPoint, bool)

LastCommitted returns the committed head last confirmed by the control plane for the current epoch.

func (*Spool) LoadRecoverySnapshot

func (s *Spool) LoadRecoverySnapshot() ([]byte, bool, error)

LoadRecoverySnapshot returns the local recovery snapshot, verified against its checksum.

func (*Spool) MarkRegistered

func (s *Spool) MarkRegistered() error

MarkRegistered records that the control plane registered the current epoch.

func (*Spool) MarkTransmitted

func (s *Spool) MarkTransmitted(seqs ...uint64) error

MarkTransmitted durably moves records to transmitted-unconfirmed; call it before sending them.

func (*Spool) Notify

func (s *Spool) Notify() <-chan struct{}

Notify is signaled, without blocking, after appends, commits, and discards.

func (*Spool) OpenEpoch

func (s *Spool) OpenEpoch(reason string, prev *protocol.EpochID, prevHead *uint64) (protocol.EpochID, error)

OpenEpoch starts a fresh chain; remaining records must be at or below prevHead of prev, and are dropped.

func (*Spool) Relieve

func (s *Spool) Relieve() (Relief, error)

Relieve evicts evidence, compacts samples, then coalesces, stopping below the relief target.

func (*Spool) SaveRecoverySnapshot

func (s *Spool) SaveRecoverySnapshot(b []byte) error

SaveRecoverySnapshot replaces the single local recovery snapshot in place.

func (*Spool) SetHalted

func (s *Spool) SetHalted(code, message string) error

SetHalted persists a rejection so the writer never resumes writing after restart until ClearHalt.

func (*Spool) SetIdentity

func (s *Spool) SetIdentity(id Identity) error

func (*Spool) Usage

func (s *Spool) Usage() Usage

Usage reports occupancy, append rate, and the projected window (free capacity / rate).

func (*Spool) WriterID

func (s *Spool) WriterID() protocol.WriterID

type Tx

type Tx struct {
	// contains filtered or unexported fields
}

Tx appends under the sequence lock; its records and cursors persist only if Do returns nil.

func (*Tx) Append

func (tx *Tx) Append(t protocol.RecordType, build func(env protocol.Envelope) (*protocol.Record, error), opts ...AppendOption) (*Entry, error)

Append builds, encodes if needed, and chains the next record; range records come only from Relieve.

func (*Tx) DeleteCursor

func (tx *Tx) DeleteCursor(key string)

DeleteCursor removes an idempotency cursor atomically with the transaction.

func (*Tx) OnAbort

func (tx *Tx) OnAbort(f func())

OnAbort runs f, in reverse registration order and still under the sequence lock, when fn fails or the commit fails.

func (*Tx) OnCommit

func (tx *Tx) OnCommit(f func())

OnCommit runs f after the transaction committed durably, still under the sequence lock; f never runs if Do fails.

type Usage

type Usage struct {
	Bytes              int64 // record body bytes, counted against Capacity
	DiskBytes          int64 // segment file bytes including superseded frames
	Capacity           int64
	Records            int
	InFlightBytes      int64 // transmitted-unconfirmed record bytes
	WindowBytes        int64
	Oldest             time.Time // emit time of the oldest spooled record, zero when empty
	AppendRate         float64   // bytes per second, exponentially weighted
	ProjectedWindow    time.Duration
	RebaselineRequired bool
	ReliefError        error // last automatic relief failure, if any
}

Usage reports spool occupancy and the projected outage window (PRD S13).

Jump to

Keyboard shortcuts

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