Documentation
¶
Index ¶
- Constants
- Variables
- func RecommendedGapBytes(maxBytes int64) int64
- type Collector
- type Config
- type RetryDecision
- type RetryPolicy
- type RetryReason
- type SaturatedError
- type Spool
- func (s *Spool) Ack(ctx context.Context, ack ports.SpoolACK) (ports.SpoolACKResult, error)
- func (s *Spool) AckGap(ctx context.Context, reported ports.SpoolGap) (bool, error)
- func (s *Spool) Close() error
- func (s *Spool) Enqueue(ctx context.Context, item ports.SpoolItem) (fleetagent.StreamPosition, error)
- func (s *Spool) Flush(ctx context.Context) error
- func (s *Spool) Gaps(ctx context.Context) ([]ports.SpoolGap, error)
- func (s *Spool) Peek(ctx context.Context, req ports.PeekSpoolRequest) ([]ports.SpoolRecord, error)
- func (s *Spool) PeekPriority(ctx context.Context, priority fleetagent.DeliveryPriority, ...) ([]ports.SpoolRecord, error)
- func (s *Spool) Stats(ctx context.Context) (ports.SpoolStats, error)
- type SyncPolicy
Constants ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
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) AckGap ¶
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) Enqueue ¶
func (s *Spool) Enqueue(ctx context.Context, item ports.SpoolItem) (fleetagent.StreamPosition, error)
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) Peek ¶
func (s *Spool) Peek(ctx context.Context, req ports.PeekSpoolRequest) ([]ports.SpoolRecord, error)
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.
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" )