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.
go get go.xchunk.org/anvil| Name | Package | Purpose |
|---|---|---|
Cache[K, V] | anvil | In-memory TTL cache, safe for concurrent use. |
WorkerPool[V] | anvil | Fixed workers fed from a bounded queue. |
Pipe, TransformPipe | anvil/pipes | Connect channels through a processing step. |
Result[T] | anvil | A value or an error. Experimental. |
FSM, FSMComparable | anvil | Deprecated. |
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.
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
Getreportsfalseif the key is missing or the item has expired. Concurrent Gets don't block each other.Setstores a value and restarts its ttl.SetWithTTLis likeSetbut the item lives for the giventtl. A non-positive value makes the item expire immediately.Invalidateremoves the item, if any.
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) // falseGetOrSet
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
loadis returned to all waiting callers and is not cached. - If
loadpanics, the panic propagates in the calling goroutine and waiting callers get an error.
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 7Cleanup, RunCleanup
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.
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.
NewWorkerPoolcreatessizeworkers with room forqueueSizepending tasks. It panics ifsizeis not positive. Workers don't run untilStartis called.Startlaunches the workers. Only the first call has an effect.Task.Execreceives the context passed toStart.Task.Resultreceives exactly oneResponsewhen the task completes. It may benilif the outcome isn't needed.- A panic in
Execand a cancelled pool context are also reported throughResponse.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
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.
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
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.
Without a middleware a Pipe forwards values unchanged. In async modes a middleware may be called from several goroutines at once.
Start
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 reachoutin completion order, not input order. - Returns early with the pool's error if a pool set with
WithWorkerPoolhas been shut down. - Returns
ErrMiddlewareRequiredif the pipe has no middleware andIndiffers fromOut.
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.
| Option | Effect |
|---|---|
WithMiddleware | Sets the function applied to every value. Alternative to WithMiddlewareErr; the last one wins. |
WithMiddlewareErr | Middleware that gets the context from Start and may fail. A failed value is dropped and its error goes to the error handler. |
WithErrorHandler | Called 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. |
WithAsync | Runs 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
Writesends to the input,Readreceives from the output (okisfalseonceoutis closed). Both block while the pipe is paused.PushandPulldo the same but ignorePause.Pause/Resumeonly affectReadandWrite; processing done byStartis 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.
Ofwraps the usual(value, error)pair, soOf(strconv.Atoi(s))works directly. Iferris non-nil the value is discarded.Mustreturns the value or panics with the error — meant for initialization that can't reasonably fail, likeMust(regexp.Compile(p)).
Methods
| Method | Description |
|---|---|
Unwrap() T | The 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) T | The value, or the fallback on error. |
IsOk() bool | Whether the result is successful. |
Error() error | The underlying error, or nil. |
Chaining
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.
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)) // -1FSM 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.
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.