spool

package
v0.2.4 Latest Latest
Warning

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

Go to latest
Published: Sep 28, 2026 License: Apache-2.0 Imports: 26 Imported by: 0

Documentation

Index

Constants

View Source
const (

	// FrameOverheadBudget bounds the header plus encoded record metadata. It is
	// exported so composition roots that shrink a spool can also shrink the
	// payload limit without creating a record that can never fit.
	FrameOverheadBudget = int64(512)
)

Variables

View Source
var (
	// ErrSaturated means the quota could not be satisfied without deleting a
	// non-sheddable P0..P2 record. Producers should apply backpressure.
	ErrSaturated = ports.ErrTelemetrySpoolSaturated
	// ErrClosed means the spool no longer accepts or serves operations.
	ErrClosed = errors.New("telemetry spool closed")
	// ErrFailed means a durability operation had an ambiguous outcome. The
	// spool fails stop until it is closed and recovered from disk.
	ErrFailed = errors.New("telemetry spool durability failure")
	// ErrGapJournalFull means loss evidence cannot be retained within its
	// reserved share of the spool quota. The spool fails stop rather than erase
	// evidence or grow beyond its configured disk bound.
	ErrGapJournalFull = errors.New("telemetry spool gap journal full")
	// ErrLocked means another process already owns the spool directory.
	ErrLocked = errors.New("telemetry spool already open")
	// ErrACKAhead means an ACK claims a sequence which this spool has never assigned.
	ErrACKAhead = errors.New("telemetry spool ACK is ahead of assigned sequence")
	// ErrStaleACK means an ACK addresses an incarnation which this spool cannot own.
	ErrStaleACK = errors.New("telemetry spool ACK addresses a stale incarnation")
)

Functions

func RecommendedGapBytes

func RecommendedGapBytes(maxBytes int64) int64

RecommendedGapBytes returns the bounded loss-evidence reserve used when MaxGapBytes is zero. The reserve includes room for an atomic rewrite copy.

Types

type Collector

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

Collector exports spool health without registering globals. The agent composition root may register it on a private registry/listener.

func NewCollector

func NewCollector(source ports.TelemetrySpool) *Collector

NewCollector constructs an unregistered collector. A nil spool is rejected at collection time as an invalid metric rather than panicking a process.

func (*Collector) Collect

func (c *Collector) Collect(ch chan<- prometheus.Metric)

func (*Collector) Describe

func (c *Collector) Describe(ch chan<- *prometheus.Desc)

type Config

type Config struct {
	// Dir must be service-owned and non-world-writable. Open tightens it to
	// 0700 and refuses symlink paths before reading spool contents.
	Dir      string
	Session  fleetagent.SessionID
	Boot     fleetagent.BootID
	MaxBytes int64
	// MaxGapBytes is reserved inside MaxBytes for durable loss evidence and its
	// atomic compaction scratch file. Zero selects a size-derived default; WAL
	// records cannot consume this reserve.
	MaxGapBytes    int64
	SegmentBytes   int64
	MaxRecordBytes int64
	PeekRecords    int
	PeekBytes      int64
	BatchInterval  time.Duration
	BatchBytes     int64
	Sync           map[fleetagent.DeliveryPriority]SyncPolicy
	Now            func() time.Time
}

Config is the bounded resource and durability policy for one spool.

func DefaultConfig

func DefaultConfig() Config

DefaultConfig returns conservative production defaults. It deliberately leaves identity and directory unset so a caller cannot accidentally open a spool whose incarnation cannot be attributed.

type RetryDecision

type RetryDecision struct {
	Retry  bool
	Delay  time.Duration
	Reason RetryReason
}

RetryDecision tells a future transport whether and when to retry a batch.

type RetryPolicy

type RetryPolicy struct {
	Base   time.Duration
	Max    time.Duration
	Random func() float64
}

RetryPolicy implements capped exponential backoff with full jitter. Random is injectable for deterministic tests; values outside [0,1] are clamped.

func DefaultRetryPolicy

func DefaultRetryPolicy() RetryPolicy

DefaultRetryPolicy is intentionally moderate: retries begin quickly, while the cap prevents a long outage from creating synchronized request storms.

func (RetryPolicy) ClassifyHTTP

func (p RetryPolicy) ClassifyHTTP(status int, retryAfter string, now time.Time, attempt uint) (RetryDecision, error)

ClassifyHTTP applies the A2 retry contract. 429 and Retry-After take precedence; 408 and 5xx retry with jitter; other 4xx are permanent. A valid Retry-After delta/date is capped to Max to keep configuration authoritative.

func (RetryPolicy) Delay

func (p RetryPolicy) Delay(attempt uint) (time.Duration, error)

Delay returns full-jitter exponential backoff for a zero-based attempt.

func (RetryPolicy) NetworkFailure

func (p RetryPolicy) NetworkFailure(attempt uint) (RetryDecision, error)

NetworkFailure applies backoff to a transport error which has no HTTP response.

type RetryReason

type RetryReason string

RetryReason is a stable label for transport metrics and logs. The A3 wire client consumes this policy; A2 owns it because retry timing determines how quickly the durable spool grows while a control plane is unavailable.

const (
	RetryNone          RetryReason = "none"
	RetryRateLimited   RetryReason = "rate_limited"
	RetryServerFailure RetryReason = "server_failure"
	RetryTimeout       RetryReason = "request_timeout"
	RetryNetwork       RetryReason = "network_failure"
	RetryPermanent     RetryReason = "permanent_failure"
)

type SaturatedError

type SaturatedError struct {
	UsedBytes     int64
	MaxBytes      int64
	RequiredBytes int64
}

SaturatedError includes safe capacity information without leaking host paths.

func (*SaturatedError) Error

func (e *SaturatedError) Error() string

func (*SaturatedError) Unwrap

func (e *SaturatedError) Unwrap() error

type Spool

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

Spool is a priority-aware, crash-recoverable implementation of ports.TelemetrySpool. All mutable state is guarded by mu, including file offsets: callers may safely enqueue, peek, ACK, and scrape metrics concurrently.

func Open

func Open(input Config) (*Spool, error)

Open exclusively opens or creates a spool and repairs recoverable torn or corrupt WAL frames before returning. The lock is held until Close.

func (*Spool) Ack

Ack commits a highest-contiguous ACK and reclaims fully acknowledged segments.

func (*Spool) AckGap

func (s *Spool) AckGap(ctx context.Context, reported ports.SpoolGap) (bool, error)

AckGap removes one durable gap only when the local object still equals the exact snapshot the server acknowledged. This closes the coalescing race: if another loss extended Count/range/time while the older report was in flight, the newer evidence remains on disk and is shipped again instead of being deleted by a stale ACK.

func (*Spool) Close

func (s *Spool) Close() error

Close flushes data, closes descriptors, and releases directory ownership.

func (*Spool) Enqueue

Enqueue durably assigns and appends an item. Returning nil guarantees that a later restart can recover the record or an explicit sequence gap.

func (*Spool) Flush

func (s *Spool) Flush(ctx context.Context) error

Flush synchronizes all active segments. It is safe to call repeatedly.

func (*Spool) Gaps

func (s *Spool) Gaps(ctx context.Context) ([]ports.SpoolGap, error)

func (*Spool) Peek

Peek returns records in strict priority order and sequence order within a lane. It does not mutate delivery state; records remain until ACKed.

func (*Spool) PeekPriority

func (s *Spool) PeekPriority(ctx context.Context, priority fleetagent.DeliveryPriority, req ports.PeekSpoolRequest) ([]ports.SpoolRecord, error)

PeekPriority returns ordered live records from exactly one A2 lane without consuming them. A3 uses this instead of the global priority-ordered Peek so a future A4 P1 backlog cannot starve raw telemetry in P2/P3, while ACK ownership remains lane-specific.

func (*Spool) Stats

func (s *Spool) Stats(ctx context.Context) (ports.SpoolStats, error)

type SyncPolicy

type SyncPolicy string

SyncPolicy controls when an accepted record is forced to stable storage. Always is the default for the non-sheddable P0..P2 lanes. Batch is intended only for P3: sequence reservations make a crash-visible gap explicit even when the operating system loses the last accepted, not-yet-synced frames.

const (
	SyncAlways SyncPolicy = "always"
	SyncBatch  SyncPolicy = "batch"
)

Jump to

Keyboard shortcuts

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