Documentation

Anvil

Anvil is a small Go library of generic building blocks: a TTL cache, a worker pool, channel pipes and a Result[T] type. It has no dependencies and requires Go 1.24 or newer.

shell
go get go.xchunk.org/anvil
NamePackagePurpose
Cache[K, V]anvilIn-memory TTL cache, safe for concurrent use.
WorkerPool[V]anvilFixed workers fed from a bounded queue.
Pipe, TransformPipeanvil/pipesConnect channels through a processing step.
Result[T]anvilA value or an error. Experimental.
FSM, FSMComparableanvilDeprecated.
main.go
import (
    "go.xchunk.org/anvil"
    "go.xchunk.org/anvil/pipes"
)

Cache

An in-memory key-value store whose items expire ttl after they were set. It is safe for concurrent use.

type Cache[K comparable, V any] struct { /* ... */ } func NewCache[K comparable, V any](ttl time.Duration) *Cache[K, V]

NewCache returns a cache whose items live for ttl. A non-positive ttl makes every item expire immediately.

💡

Expiry model. Expiry is checked on Get. Memory of expired items is reclaimed by Set, which sweeps them at most once per max(ttl, one minute), and on demand by Cleanup or RunCleanup.

Get, Set, SetWithTTL, Invalidate

func (c *Cache[K, V]) Get(key K) (V, bool) func (c *Cache[K, V]) Set(key K, value V) func (c *Cache[K, V]) SetWithTTL(key K, value V, ttl time.Duration) func (c *Cache[K, V]) Invalidate(key K)
  • Get reports false if the key is missing or the item has expired. Concurrent Gets don't block each other.
  • Set stores a value and restarts its ttl.
  • SetWithTTL is like Set but the item lives for the given ttl. A non-positive value makes the item expire immediately.
  • Invalidate removes the item, if any.
example
c := anvil.NewCache[string, int](time.Minute)
c.Set("answer", 42)

v, ok := c.Get("answer")
fmt.Println(v, ok) // 42 true

c.Invalidate("answer")
_, ok = c.Get("answer")
fmt.Println(ok) // false

GetOrSet

func (c *Cache[K, V]) GetOrSet(key K, load func() (V, error)) (V, error)

Returns the item stored under key. If there is none, it calls load, stores the result with the cache's default ttl and returns it.

Concurrent calls for the same missing key share a single load: one caller runs it and the others wait for its result.

  • An error from load is returned to all waiting callers and is not cached.
  • If load panics, the panic propagates in the calling goroutine and waiting callers get an error.
example
c := anvil.NewCache[string, int](time.Minute)

load := func() (int, error) {
    fmt.Println("loading")
    return 7, nil
}
a, _ := c.GetOrSet("k", load) // loads
b, _ := c.GetOrSet("k", load) // served from the cache
fmt.Println(a, b)
// Output:
// loading
// 7 7

Cleanup, RunCleanup

func (c *Cache[K, V]) Cleanup() int func (c *Cache[K, V]) RunCleanup(ctx context.Context, interval time.Duration)

Cleanup removes all expired items and returns how many were removed. Since Set already sweeps periodically, you only need it to reclaim memory sooner, or when the cache is no longer written to.

RunCleanup calls Cleanup every interval until ctx is done. It blocks, so run it in its own goroutine. It panics if interval is not positive.

example
go c.RunCleanup(ctx, 30*time.Second)

WorkerPool

Runs tasks on a fixed number of goroutines fed from a bounded queue. Create it with NewWorkerPool, call Start once, submit tasks with Submit or SubmitCtx and finish with Shutdown.

func NewWorkerPool[V any](size, queueSize int) *WorkerPool[V] func (wp *WorkerPool[V]) Start(ctx context.Context) type Task[V any] struct { Exec func(ctx context.Context) (V, error) Result chan Response[V] } type Response[V any] struct { Value V Err error }
  • NewWorkerPool creates size workers with room for queueSize pending tasks. It panics if size is not positive. Workers don't run until Start is called.
  • Start launches the workers. Only the first call has an effect.
  • Task.Exec receives the context passed to Start.
  • Task.Result receives exactly one Response when the task completes. It may be nil if the outcome isn't needed.
  • A panic in Exec and a cancelled pool context are also reported through Response.Err.
⚠️

Use a buffered Result channel. The worker sends on it synchronously, so use a channel with capacity ≥ 1, or make sure a receiver is always waiting — otherwise the worker stays blocked.

Submit, SubmitCtx

var ErrPoolClosed = errors.New("anvil: worker pool is closed") func (wp *WorkerPool[V]) Submit(task Task[V]) error func (wp *WorkerPool[V]) SubmitCtx(ctx context.Context, task Task[V]) error

Submit enqueues a task, blocking while the queue is full. It returns ErrPoolClosed if the pool has been shut down. SubmitCtx is the same but gives up with ctx.Err() if ctx is done before the task could be enqueued.

example
wp := anvil.NewWorkerPool[int](2, 4)
wp.Start(context.Background())
defer wp.Shutdown()

result := make(chan anvil.Response[int], 1)
wp.Submit(anvil.Task[int]{
    Result: result,
    Exec:   func(ctx context.Context) (int, error) { return 6 * 7, nil },
})

res := <-result
fmt.Println(res.Value, res.Err) // 42 <nil>

Shutdown

func (wp *WorkerPool[V]) Shutdown()

Stops accepting tasks, lets queued tasks finish and waits for the workers to exit. It is safe to call more than once. Workers keep draining the queue after cancellation, so every submitted task still gets a response (with the context's error).

pipes

The pipes package connects channels through a processing step. A TransformPipe[In, Out] reads values from an input channel, passes them through a middleware and writes the results to an output channel. Pipe[T] is a TransformPipe[T, T] with a simpler set of options.

func NewPipe[T any](in chan T, out chan T, opts ...PipeOption[T]) *Pipe[T] func NewTransformPipe[In, Out any](in chan In, out chan Out, opts ...TransformPipeOption[In, Out]) *TransformPipe[In, Out] type Middleware[In, Out any] func(v In) Out type ErrMiddleware[In, Out any] func(ctx context.Context, v In) (Out, error)

Without a middleware a Pipe forwards values unchanged. In async modes a middleware may be called from several goroutines at once.

Start

func (p *TransformPipe[In, Out]) Start(ctx context.Context) error

Processes values until in is closed (returns nil) or ctx is done (returns ctx.Err()), then waits for work still in flight and closes out. It blocks, so run it in its own goroutine.

  • By default values are processed one at a time, in order.
  • In async modes (WithAsync, WithWorkerPool, WithConcurrency) calls overlap and results reach out in completion order, not input order.
  • Returns early with the pool's error if a pool set with WithWorkerPool has been shut down.
  • Returns ErrMiddlewareRequired if the pipe has no middleware and In differs from Out.
example
in, out := make(chan string), make(chan int)
p := pipes.NewTransformPipe(in, out,
    pipes.WithTransformMiddlewareErr(func(_ context.Context, s string) (int, error) {
        return strconv.Atoi(s)
    }),
    pipes.WithTransformConcurrency[string, int](4),
    pipes.WithTransformErrorHandler[string, int](func(err error) { log.Println(err) }),
)
go p.Start(ctx)

Options

Every option has a Pipe form and a WithTransform… form for TransformPipe.

OptionEffect
WithMiddlewareSets the function applied to every value. Alternative to WithMiddlewareErr; the last one wins.
WithMiddlewareErrMiddleware that gets the context from Start and may fail. A failed value is dropped and its error goes to the error handler.
WithErrorHandlerCalled with every error that makes the pipe drop a value, including worker-pool failures such as a panic. In async modes it may run concurrently, so it must be goroutine-safe.
WithAsyncRuns the middleware in a goroutine per value. Order is not kept. No effect without a middleware.
WithConcurrency(n)At most n middleware calls in flight. Enables async mode itself. Panics if n is not positive.
WithWorkerPool(wp)Runs the middleware of an async pipe on wp. Alternative to WithConcurrency; the last one wins.

Pause, Resume, Read, Write

func (p *TransformPipe[In, Out]) Write(v In) func (p *TransformPipe[In, Out]) Read() (Out, bool) func (p *TransformPipe[In, Out]) Push(v In) func (p *TransformPipe[In, Out]) Pull() (Out, bool) func (p *TransformPipe[In, Out]) Pause() func (p *TransformPipe[In, Out]) Resume()
  • Write sends to the input, Read receives from the output (ok is false once out is closed). Both block while the pipe is paused.
  • Push and Pull do the same but ignore Pause.
  • Pause / Resume only affect Read and Write; processing done by Start is not affected.
⚠️

Lock and Unlock are deprecated aliases of Pause and Resume. Despite the names they are not a mutual-exclusion lock.

Result experimental

A Result[T] holds either a successful value or an error.

🧪

Experimental. Result is a new feature. Its API may change or be removed in any release without notice.

func Ok[T any](value T) Result[T] func Err[T any](err error) Result[T] func Of[T any](value T, err error) Result[T] func Must[T any](value T, err error) T
  • Of wraps the usual (value, error) pair, so Of(strconv.Atoi(s)) works directly. If err is non-nil the value is discarded.
  • Must returns the value or panics with the error — meant for initialization that can't reasonably fail, like Must(regexp.Compile(p)).

Methods

MethodDescription
Unwrap() TThe value, or panics with an error wrapping the result's error (inspect with errors.Is / errors.As).
Value() (T, error)The usual Go pair. The value is the zero value on error.
UnwrapOr(fallback T) TThe value, or the fallback on error.
IsOk() boolWhether the result is successful.
Error() errorThe underlying error, or nil.

Chaining

func Map[T, U any](r Result[T], f func(T) U) Result[U] func AndThen[T, U any](r Result[T], f func(T) Result[U]) Result[U]

These are functions rather than methods because Go methods can't declare their own type parameters. If r holds an error, f is not called and the error is passed through. AndThen lets fallible steps be chained.

example
r := anvil.Of(strconv.Atoi("12"))
doubled := anvil.Map(r, func(n int) int { return n * 2 })

v, err := doubled.Value()
fmt.Println(v, err) // 24 <nil>

fmt.Println(anvil.Of(strconv.Atoi("x")).UnwrapOr(-1)) // -1

FSM deprecated

⚠️

FSM and FSMComparable are deprecated. They are a mutex-guarded map of per-key states with a default state — not a state machine (no transitions or validation) — and may be removed in a future release.

func NewFSM[K comparable, V any](initState V) *FSM[K, V] func NewFSMComparable[K comparable, V comparable](initState V) *FSMComparable[K, V] // both types: Get(id K) (V, bool) Set(id K, val V) Init(id K) // sets the initial state if there is none Reset(id K) // sets the initial state // FSMComparable only: IsInit(id K) bool

Versioning

The module follows semantic versioning. Deprecated and experimental APIs are marked as such in their doc comments. More runnable examples live in example_test.go and pipes/example_test.go; see also go doc and pkg.go.dev.