collection

package
v0.0.2 Latest Latest
Warning

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

Go to latest
Published: Sep 21, 2026 License: MIT Imports: 9 Imported by: 0

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

View Source
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 Bucket

type Bucket[T Numerical] struct {
	Sum   T
	Count int64
}

Bucket holds the sum and count of the additions.

func (*Bucket[T]) Add

func (b *Bucket[T]) Add(v T)

Add adds v to the bucket sum.

func (*Bucket[T]) Reset

func (b *Bucket[T]) Reset()

Reset clears the bucket.

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) Del

func (c *Cache) Del(key string)

Del deletes the item with the given key.

func (*Cache) Get

func (c *Cache) Get(key string) (any, bool)

Get returns the item with the given key.

func (*Cache) Set

func (c *Cache) Set(key string, value any)

Set stores value under key with the cache's default expiry.

func (*Cache) SetWithExpire

func (c *Cache) SetWithExpire(key string, value any, expire time.Duration)

SetWithExpire stores value under key with the given expiry.

func (*Cache) Take

func (c *Cache) Take(key string, fetch func() (any, error)) (any, error)

Take returns the item under key, fetching and caching it when absent. The fetch is deduplicated across concurrent callers via singleflight.

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 Execute

type Execute func(key, value any)

Execute runs a fired timer task.

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

func NewQueue(size int) *Queue

NewQueue returns a Queue with the given initial capacity. It panics on a non-positive size.

func (*Queue) Empty

func (q *Queue) Empty() bool

Empty reports whether the queue has no elements.

func (*Queue) Put

func (q *Queue) Put(element any)

Put appends element to the queue, growing it as needed.

func (*Queue) Take

func (q *Queue) Take() (any, bool)

Take removes and returns the oldest element, or false if the queue is empty.

type Ring

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

Ring is a fixed-size ring buffer. It overwrites the oldest element once full.

func NewRing

func NewRing(n int) *Ring

NewRing returns a Ring of size n. It panics on a non-positive size.

func (*Ring) Add

func (r *Ring) Add(v any)

Add appends v to the ring, overwriting the oldest element when full.

func (*Ring) Take

func (r *Ring) Take() []any

Take returns the ring's elements in insertion order (oldest first).

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.

func (*SafeMap[K, V]) Del

func (sm *SafeMap[K, V]) Del(key K)

Del removes key from the map.

func (*SafeMap[K, V]) Get

func (sm *SafeMap[K, V]) Get(key K) (V, bool)

Get returns the value for key and whether it was present.

func (*SafeMap[K, V]) Len

func (sm *SafeMap[K, V]) Len() int

Len returns the number of live entries.

func (*SafeMap[K, V]) Range

func (sm *SafeMap[K, V]) Range(fn func(key K, value V) bool)

Range calls fn for each entry; iteration stops when fn returns false.

func (*SafeMap[K, V]) Set

func (sm *SafeMap[K, V]) Set(key K, value V)

Set stores key -> value.

type Set

type Set[T comparable] map[T]struct{}

Set is a concurrency-unsafe generic set. Guard it externally when shared across goroutines.

func NewSet

func NewSet[T comparable]() Set[T]

NewSet returns an empty Set.

func (Set[T]) Add

func (s Set[T]) Add(v T)

Add inserts v into the set.

func (Set[T]) Contains

func (s Set[T]) Contains(v T) bool

Contains reports whether v is in the set.

func (Set[T]) Len

func (s Set[T]) Len() int

Len returns the number of elements.

func (Set[T]) Remove

func (s Set[T]) Remove(v T)

Remove deletes v from the set.

func (Set[T]) Values

func (s Set[T]) Values() []T

Values returns the elements in an unspecified order.

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

func NewTimingWheel(interval time.Duration, numSlots int, execute Execute) (*TimingWheel, error)

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.

Jump to

Keyboard shortcuts

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