Bounded thread-pool primitives — factory + safe submit/await.
Extends hive-weave with a pool abstraction so downstream code does not reach into java.util.concurrent directly (DIP).
Responsibilities:
ThreadPoolExecutor with CallerRunsPolicy
backpressure and a named thread factory (for JVM diagnostics).submit! returning an opaque Future-like handle.await! — submit + block up to a timeout, returning a
fallback on timeout/error. Never hangs.pool-stats and shutdown! for lifecycle.Callers keep pool instances in their own registry (e.g. named
io/compute/event/memory pools) and hand them to await! when they
need bounded, isolated execution for a piece of work.
Quick start: (require '[hive-weave.pool :as wp])
(def db-pool (wp/make-pool {:name "db" :size 8}))
(wp/await! db-pool (fn [] (query-database ...)) {:timeout-ms 5000 :fallback ::db-timeout}) ;; => result or ::db-timeout
Bounded thread-pool primitives — factory + safe submit/await.
Extends hive-weave with a pool abstraction so downstream code does
not reach into java.util.concurrent directly (DIP).
Responsibilities:
- Construct a bounded `ThreadPoolExecutor` with CallerRunsPolicy
backpressure and a named thread factory (for JVM diagnostics).
- Expose `submit!` returning an opaque Future-like handle.
- Expose `await!` — submit + block up to a timeout, returning a
fallback on timeout/error. Never hangs.
- Re-export `pool-stats` and `shutdown!` for lifecycle.
Callers keep pool *instances* in their own registry (e.g. named
io/compute/event/memory pools) and hand them to `await!` when they
need bounded, isolated execution for a piece of work.
Quick start:
(require '[hive-weave.pool :as wp])
(def db-pool (wp/make-pool {:name "db" :size 8}))
(wp/await! db-pool
(fn [] (query-database ...))
{:timeout-ms 5000 :fallback ::db-timeout})
;; => result or ::db-timeout(await! pool f {:keys [timeout-ms fallback name] :or {name "pool-task"}})Submit f to pool and block on its result up to :timeout-ms.
On timeout, cancels the task (with interrupt) and returns :fallback.
On exception during execution, logs and returns :fallback.
A saturated :abort pool also yields :fallback, logged as a rejection.
Never hangs indefinitely.
This BLOCKS the calling thread. Inside an async server (an aleph handler, a netty event-loop thread) that is the thing you are trying to avoid: there, submit! and compose on the Future instead.
Options: :timeout-ms — max wait in ms (required) :fallback — value returned on timeout, exception or rejection (default nil) :name — diagnostic label used in logs (default "pool-task")
Submit `f` to `pool` and block on its result up to `:timeout-ms`. On timeout, cancels the task (with interrupt) and returns `:fallback`. On exception during execution, logs and returns `:fallback`. A saturated `:abort` pool also yields `:fallback`, logged as a rejection. Never hangs indefinitely. This BLOCKS the calling thread. Inside an async server (an aleph handler, a netty event-loop thread) that is the thing you are trying to avoid: there, submit! and compose on the Future instead. Options: :timeout-ms — max wait in ms (required) :fallback — value returned on timeout, exception or rejection (default nil) :name — diagnostic label used in logs (default "pool-task")
(bound-future & body)Drop-in replacement for clojure.core/future that conveys the caller's
dynamic-var bindings through the active IBindingConveyor.
Drop-in replacement for `clojure.core/future` that conveys the caller's dynamic-var bindings through the active IBindingConveyor.
(capture-frame)Snapshot the current thread's binding frame. Pair with FixedFrameConveyor to inject this frame into work-fns running on other threads.
Snapshot the current thread's binding frame. Pair with FixedFrameConveyor to inject this frame into work-fns running on other threads.
(convey-fn f)Wrap thunk f via the active conveyor — the canonical entry point for any async submission that must preserve dynamic-var bindings.
Wrap thunk f via the active conveyor — the canonical entry point for any async submission that must preserve dynamic-var bindings.
(get-conveyor)Return active conveyor. Falls back to BoundFnConveyor if atom is nil.
Return active conveyor. Falls back to BoundFnConveyor if atom is nil.
(convey this f)Wrap thunk f so dynamic-var bindings transfer to the executing thread.
Wrap thunk f so dynamic-var bindings transfer to the executing thread.
(io-executor
{:keys [name fallback-size virtual?] :or {fallback-size 64} :as opts})The executor to run BLOCKING IO on: virtual threads where the JVM has them,
otherwise a bounded pool of :fallback-size platform threads (default 64).
Use it when the work waits on something and the thread itself is the only
thing being rationed. Keep make-pool for CPU-bound work, where a pool
sized near the core count is the point, and put a gate or budget in front of
whatever the work contends for.
Options: :name, :fallback-size, :virtual? (defaults to what this JVM can
do), plus anything make-pool takes. :virtual? false forces the pool,
which is how the fallback stays testable on a JVM that has virtual threads.
The executor to run BLOCKING IO on: virtual threads where the JVM has them, otherwise a bounded pool of `:fallback-size` platform threads (default 64). Use it when the work waits on something and the thread itself is the only thing being rationed. Keep `make-pool` for CPU-bound work, where a pool sized near the core count is the point, and put a gate or budget in front of whatever the work contends for. Options: :name, :fallback-size, :virtual? (defaults to what this JVM can do), plus anything `make-pool` takes. `:virtual? false` forces the pool, which is how the fallback stays testable on a JVM that has virtual threads.
(make-pool {:keys [name size queue-capacity keep-alive-s stack-bytes rejection]
:or {queue-capacity default-queue-capacity
keep-alive-s 60
rejection :caller-runs}})Create a bounded fixed-size ThreadPoolExecutor.
Options: :name thread-name prefix and diagnostic label (required) :size fixed pool size (required) :queue-capacity bounded LinkedBlockingQueue capacity (default 256) :keep-alive-s idle keep-alive in seconds (default 60) :stack-bytes explicit worker stack size (default: the JVM's) :rejection what happens once workers AND queue are full: :caller-runs (default) — the submitting thread runs it :abort — submit! throws RejectedExecutionException
:caller-runs never drops work but blocks the submitter, which is wrong when the submitter is a request thread that should shed load instead. Choose :abort there and answer the rejection.
Pass :stack-bytes when the tasks call into native code that recurses. A native stack overflow is a SIGSEGV, not an exception, so a pool sized for Clojure work will take the whole process down with it (hive-weave.stack). Note that CallerRunsPolicy runs a rejected task on the CALLER's stack, which this option cannot size: keep the queue big enough that native work is not pushed back onto the submitter.
Create a bounded fixed-size ThreadPoolExecutor.
Options:
:name thread-name prefix and diagnostic label (required)
:size fixed pool size (required)
:queue-capacity bounded LinkedBlockingQueue capacity (default 256)
:keep-alive-s idle keep-alive in seconds (default 60)
:stack-bytes explicit worker stack size (default: the JVM's)
:rejection what happens once workers AND queue are full:
:caller-runs (default) — the submitting thread runs it
:abort — submit! throws RejectedExecutionException
:caller-runs never drops work but blocks the submitter, which is wrong when
the submitter is a request thread that should shed load instead. Choose
:abort there and answer the rejection.
Pass :stack-bytes when the tasks call into native code that recurses. A
native stack overflow is a SIGSEGV, not an exception, so a pool sized for
Clojure work will take the whole process down with it (hive-weave.stack).
Note that CallerRunsPolicy runs a rejected task on the CALLER's stack, which
this option cannot size: keep the queue big enough that native work is not
pushed back onto the submitter.(pool-stats pool)Snapshot of an executor's runtime counters.
Every ExecutorService answers: :executor (class name), :shutdown?, :terminated?. A ThreadPoolExecutor additionally answers its counters and the :rejection policy it was built with.
Written against the INTERFACE because the executors weave has to live
beside are not all ThreadPoolExecutors: manifold's and aleph's are
dirigiste Executors, which extend AbstractExecutorService, and a
ThreadPoolExecutor-shaped read of one throws ClassCastException.
Snapshot of an executor's runtime counters. Every ExecutorService answers: :executor (class name), :shutdown?, :terminated?. A ThreadPoolExecutor additionally answers its counters and the :rejection policy it was built with. Written against the INTERFACE because the executors weave has to live beside are not all ThreadPoolExecutors: manifold's and aleph's are dirigiste `Executor`s, which extend AbstractExecutorService, and a ThreadPoolExecutor-shaped read of one throws ClassCastException.
(rejected? e)Whether e is a java.util.concurrent.RejectedExecutionException.
Tested by class name rather than caught by class because babashka does not expose that class: naming it in an :import or a catch clause makes the namespace unloadable there, and with it every bb tool that reaches hive-weave.parallel. Catch RuntimeException and ask this instead.
Whether `e` is a java.util.concurrent.RejectedExecutionException. Tested by class name rather than caught by class because babashka does not expose that class: naming it in an :import or a catch clause makes the namespace unloadable there, and with it every bb tool that reaches hive-weave.parallel. Catch RuntimeException and ask this instead.
(set-conveyor! c)Install conveyor c as the active binding conveyor. Returns prior conveyor.
Install conveyor c as the active binding conveyor. Returns prior conveyor.
(shutdown! pool & [{:keys [await-ms] :or {await-ms 5000}}])Orderly shutdown: stop accepting new tasks, wait up to
:await-ms for in-flight tasks, then force-shutdown.
Default :await-ms is 5000.
Orderly shutdown: stop accepting new tasks, wait up to `:await-ms` for in-flight tasks, then force-shutdown. Default `:await-ms` is 5000.
(submit! pool f)Submit f to pool, returning a java.util.concurrent.Future.
f is wrapped via the active IBindingConveyor (default
BoundFnConveyor, equivalent to clojure.core/bound-fn*) so the
caller's dynamic var frame is conveyed to the pool thread. This
matches the behaviour of clojure.core/future and avoids a silent
trap where code relying on binding loses its frame at the pool
boundary.
A rejection from a SHUT DOWN pool runs f on the caller thread and returns
an already-completed Future: the work was accepted before the shutdown race
and still has to happen. A rejection from a SATURATED :abort pool is
rethrown — that one is load shedding, and swallowing it would turn the
policy the caller asked for back into :caller-runs.
Submit `f` to `pool`, returning a java.util.concurrent.Future. `f` is wrapped via the active `IBindingConveyor` (default `BoundFnConveyor`, equivalent to `clojure.core/bound-fn*`) so the caller's dynamic var frame is conveyed to the pool thread. This matches the behaviour of `clojure.core/future` and avoids a silent trap where code relying on `binding` loses its frame at the pool boundary. A rejection from a SHUT DOWN pool runs `f` on the caller thread and returns an already-completed Future: the work was accepted before the shutdown race and still has to happen. A rejection from a SATURATED `:abort` pool is rethrown — that one is load shedding, and swallowing it would turn the policy the caller asked for back into :caller-runs.
(virtual-executor)An ExecutorService that starts one virtual thread per task.
A virtual thread costs a few hundred bytes and parks instead of holding a carrier while it blocks, so the thread stops being the scarce thing. What that removes is the THREAD ceiling; it removes no other ceiling, which is why this is normally paired with a gate or a budget over whatever is actually finite (connections, GPU memory, a provider's rate limit).
Throws on a JVM older than 21 rather than silently degrading, because a caller that asked for unbounded concurrency and quietly got 32 threads would be the worst of both.
Called through java.lang.reflect rather than clojure.lang.Reflector, which babashka does not expose.
An ExecutorService that starts one virtual thread per task. A virtual thread costs a few hundred bytes and parks instead of holding a carrier while it blocks, so the thread stops being the scarce thing. What that removes is the THREAD ceiling; it removes no other ceiling, which is why this is normally paired with a gate or a budget over whatever is actually finite (connections, GPU memory, a provider's rate limit). Throws on a JVM older than 21 rather than silently degrading, because a caller that asked for unbounded concurrency and quietly got 32 threads would be the worst of both. Called through java.lang.reflect rather than clojure.lang.Reflector, which babashka does not expose.
(virtual-threads?)Whether this JVM can start a virtual thread per task (Java 21+).
Asked reflectively so that loading this namespace on an older JVM costs
nothing: a direct call to Executors/newVirtualThreadPerTaskExecutor links
at class load and fails there rather than here, where the caller can choose.
Whether this JVM can start a virtual thread per task (Java 21+). Asked reflectively so that loading this namespace on an older JVM costs nothing: a direct call to `Executors/newVirtualThreadPerTaskExecutor` links at class load and fails there rather than here, where the caller can choose.
(with-pool-await pool opts & body)Submit body to pool, block up to (:timeout-ms opts), return
(:fallback opts) on timeout/exception.
(with-pool-await memory-pool {:timeout-ms 30000 :fallback ::failed} (chroma/add-entry! ...))
Submit body to `pool`, block up to (:timeout-ms opts), return
(:fallback opts) on timeout/exception.
(with-pool-await memory-pool {:timeout-ms 30000 :fallback ::failed}
(chroma/add-entry! ...))cljdoc builds & hosts documentation for Clojure/Script libraries
| Ctrl+k | Jump to recent docs |
| ← | Move to previous article |
| → | Move to next article |
| Ctrl+/ | Jump to the search field |