Documentation
¶
Overview ¶
Package collection provides concurrency-safe generic collections used by the resilience primitives: a rolling window for time-bucketed statistics and a safe map that avoids the Go map-delete memory-growth issue.
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ( // ErrClosed indicates the TimingWheel was already stopped. ErrClosed = errors.New("collection: timing wheel is closed") // ErrArgument indicates an invalid timer argument. ErrArgument = errors.New("collection: incorrect timer argument") )
Functions ¶
This section is empty.
Types ¶
type BucketInterface ¶
type BucketInterface[T Numerical] interface { Add(v T) Reset() }
BucketInterface is the interface that buckets must satisfy.
type Cache ¶
type Cache struct {
// contains filtered or unexported fields
}
Cache is an in-memory cache with timing-wheel expiry, optional LRU eviction and singleflight-protected read-through.
func NewCache ¶
func NewCache(expire time.Duration, opts ...CacheOption) (*Cache, error)
NewCache returns a Cache whose entries expire after expire. The expiry is jittered slightly to avoid thundering-herd expirations.
func (*Cache) Close ¶
func (c *Cache) Close()
Close stops the cache's timing wheel and releases its goroutine.
func (*Cache) SetWithExpire ¶
SetWithExpire stores value under key with the given expiry.
type CacheOption ¶
type CacheOption func(*Cache)
CacheOption customizes a Cache.
func WithLimit ¶
func WithLimit(limit int) CacheOption
WithLimit caps the number of cached items via LRU eviction.
type Numerical ¶
type Numerical interface {
~int | ~int8 | ~int16 | ~int32 | ~int64 |
~uint | ~uint8 | ~uint16 | ~uint32 | ~uint64 |
~float32 | ~float64
}
Numerical is a constraint that permits any numeric type.
type Queue ¶
type Queue struct {
// contains filtered or unexported fields
}
Queue is an auto-growing FIFO queue.
func NewQueue ¶
NewQueue returns a Queue with the given initial capacity. It panics on a non-positive size.
type Ring ¶
type Ring struct {
// contains filtered or unexported fields
}
Ring is a fixed-size ring buffer. It overwrites the oldest element once full.
type RollingWindow ¶
type RollingWindow[T Numerical, B BucketInterface[T]] struct { // contains filtered or unexported fields }
RollingWindow keeps a fixed number of time buckets and aggregates the values added to them, sliding the bucket boundary forward as time advances.
func NewRollingWindow ¶
func NewRollingWindow[T Numerical, B BucketInterface[T]](newBucket func() B, size int, interval time.Duration, opts ...RollingWindowOption[T, B]) *RollingWindow[T, B]
NewRollingWindow returns a RollingWindow with size buckets each covering interval. newBucket creates the concrete bucket instances.
func (*RollingWindow[T, B]) Add ¶
func (rw *RollingWindow[T, B]) Add(v T)
Add adds a value to the current bucket.
func (*RollingWindow[T, B]) Reduce ¶
func (rw *RollingWindow[T, B]) Reduce(fn func(b B))
Reduce runs fn on all buckets, ignoring the current bucket if IgnoreCurrentBucket was set.
type RollingWindowOption ¶
type RollingWindowOption[T Numerical, B BucketInterface[T]] func(*RollingWindow[T, B])
RollingWindowOption customizes a RollingWindow.
func IgnoreCurrentBucket ¶
func IgnoreCurrentBucket[T Numerical, B BucketInterface[T]]() RollingWindowOption[T, B]
IgnoreCurrentBucket makes Reduce ignore the current (partial) bucket.
type SafeMap ¶
type SafeMap[K comparable, V any] struct { // contains filtered or unexported fields }
SafeMap is a concurrency-safe generic map that avoids Go's map-delete memory retention (the backing array never shrinks on delete, see issue #20135). Once deletions outnumber live entries on a large map, the next Set rebuilds the backing map so the oversized allocation can be reclaimed.
func NewSafeMap ¶
func NewSafeMap[K comparable, V any]() *SafeMap[K, V]
NewSafeMap returns an empty SafeMap.
type Set ¶
type Set[T comparable] map[T]struct{}
Set is a concurrency-unsafe generic set. Guard it externally when shared across goroutines.
type TimingWheel ¶
type TimingWheel struct {
// contains filtered or unexported fields
}
TimingWheel is a timing wheel that schedules tasks by their delay. A single goroutine drives it: ticks advance the wheel, and a handful of channels serialize Set/Move/Remove/Drain/Stop operations.
func NewTimingWheel ¶
NewTimingWheel returns a TimingWheel driven by a real ticker.
func NewTimingWheelWithTicker ¶
func NewTimingWheelWithTicker(interval time.Duration, numSlots int, execute Execute, ticker timex.Ticker) (*TimingWheel, error)
NewTimingWheelWithTicker returns a TimingWheel driven by the given ticker.
func (*TimingWheel) Drain ¶
func (tw *TimingWheel) Drain(fn func(key, value any)) error
Drain fires all pending tasks through fn and returns once finished.
func (*TimingWheel) MoveTimer ¶
func (tw *TimingWheel) MoveTimer(key any, delay time.Duration) error
MoveTimer reschedules the task with the given key to a new delay.
func (*TimingWheel) RemoveTimer ¶
func (tw *TimingWheel) RemoveTimer(key any) error
RemoveTimer cancels the task with the given key.
func (*TimingWheel) SetTimer ¶
func (tw *TimingWheel) SetTimer(key, value any, delay time.Duration) error
SetTimer schedules value to fire after delay, keyed by key.
func (*TimingWheel) Stop ¶
func (tw *TimingWheel) Stop()
Stop stops the wheel and its ticker. It is idempotent: repeated calls are safe and only the first one closes the stop channel.