Cache[K, V]
In-memory TTL cache, safe for concurrent use.
GetOrSetshares one load between concurrent callersSetWithTTLsets a per-item lifetime- Expired items are swept automatically
A TTL cache, a worker pool, channel pipes and a Result[T] type — concurrency-safe, context-aware and with no dependencies.
go get go.xchunk.org/anvilc := anvil.NewCache[string, User](time.Minute)
// The loader runs once, even if many
// goroutines ask for the same key.
u, err := c.GetOrSet("alice", func() (User, error) {
return loadUser("alice")
})
// Per-item lifetime.
c.SetWithTTL("session", s, 5*time.Second)wp := anvil.NewWorkerPool[int](4, 16) // 4 workers, queue of 16
wp.Start(ctx)
defer wp.Shutdown()
result := make(chan anvil.Response[int], 1)
wp.Submit(anvil.Task[int]{
Result: result,
Exec: compute,
})
res := <-result // res.Value, res.Errin, 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),
)
go p.Start(ctx)r := anvil.Of(strconv.Atoi("12"))
doubled := anvil.Map(r, func(n int) int {
return n * 2
})
v, err := doubled.Value() // 24, nil
anvil.Of(strconv.Atoi("x")).UnwrapOr(-1) // -1Each piece does one job, has a tiny API surface and documents exactly how it behaves under concurrency and cancellation.
Cache[K, V]In-memory TTL cache, safe for concurrent use.
GetOrSet shares one load between concurrent callersSetWithTTL sets a per-item lifetimeWorkerPool[V]A fixed number of workers fed from a bounded queue.
Submit / SubmitCtxShutdown drains the queuepipesConnect channels through a processing step — sequential, async or bounded concurrent, fallible or not.
Pipe and TransformPipe with type conversionPause / Resume flow controlResult[T]experimentalA value or an error, with chainable helpers.
Of, Map, AndThen, Must(value, error) pairA browser simulation of the behaviours that matter. Nothing here calls Go — it only mirrors the documented semantics.
Six goroutines ask for the same missing key. One becomes the leader and runs load; the rest wait for its result.
The cache has a default ttl; SetWithTTL overrides it per item. Expired items are never returned by Get.
The pool's queue is bounded: Submit blocks when it is full, and SubmitCtx gives up when your context does. A panicking task becomes an error, not a crashed process.
wp := anvil.NewWorkerPool[[]byte](8, 64)
wp.Start(ctx)
defer wp.Shutdown()
ctx, cancel := context.WithTimeout(ctx, time.Second)
defer cancel()
err := wp.SubmitCtx(ctx, anvil.Task[[]byte]{
Result: out, // buffered: the worker never blocks
Exec: func(ctx context.Context) ([]byte, error) {
return fetch(ctx, url)
},
})
if errors.Is(err, anvil.ErrPoolClosed) { /* shutting down */ }
Parse, enrich or fan out values between channels with one declarative setup. Run up to N steps at once, drop the ones that fail and observe why with an error handler.
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) // dropped value
}),
)
go p.Start(ctx) // closes out when in is closed
p.Pause() // Read / Write block...
p.Resume() // ...until resumed
Everything is documented with signatures, behaviour notes and runnable examples.
Open the documentation →