undercurrent.router¶
Advanced API
Most users want ProbedModel, which
drives a router for you. Use the router directly when you embed
Undercurrent in your own serving stack.
The router dispatches activation records to per-request probe instances,
runs inline probes on the generation thread and async probes on bounded
worker pools, and records metrics. Router, RequestHandle, RouterError,
OverflowPolicy and the metrics types are also exported from undercurrent.
Router¶
Router(probe_registry=None, *, worker_pool_size=None, default_queue_depth=DEFAULT_QUEUE_DEPTH, default_overflow_policy=OverflowPolicy.DROP_OLDEST, drain_timeout=DEFAULT_DRAIN_TIMEOUT, metrics_sink=None, default_intervention_policy=None, circuit_breaker_threshold=DEFAULT_CIRCUIT_BREAKER_THRESHOLD)
¶
Dispatches each ActivationRecord to the probe instances of its request.
Inline extraction points run on the caller's thread and can return an
ABORT signal; async ones run on a shared, bounded worker pool and
only observe. Every request gets its own freshly spawned probe instances,
and nothing is shared across requests.
with Router() as router: # probes come from @register_probe
with router.request(extraction_points=spec) as req:
for record in activation_stream:
signal = req.route(record) # non-None only for inline points
if signal is not None and signal.action is ProbeAction.ABORT:
break
results = req.results # {extraction_point_name: ProbeResult}
router.request(...) wraps the lower-level
register_request /
route /
end_request calls. Call
shutdown (or use
with Router(...)) when done.
Each async extraction point of a live request holds one worker thread
until its request ends, which keeps its activations in order. If more
than worker_pool_size async bindings are alive at once, the extra
ones wait for a free worker (their bounded queues apply their overflow
policy meanwhile).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
probe_registry
|
Mapping[str, ProbeFactory | type[Probe]] | ProbeRegistry | None
|
where |
None
|
worker_pool_size
|
int | None
|
number of threads backing all async bindings,
shared across every request. Defaults to
|
None
|
default_queue_depth
|
int
|
queue depth for an async extraction point whose
|
DEFAULT_QUEUE_DEPTH
|
default_overflow_policy
|
OverflowPolicy
|
what a full async queue does; see
|
DROP_OLDEST
|
drain_timeout
|
float
|
seconds |
DEFAULT_DRAIN_TIMEOUT
|
metrics_sink
|
MetricsSink | None
|
an optional external
|
None
|
default_intervention_policy
|
InterventionPolicy | None
|
policy for an extraction point whose
|
None
|
circuit_breaker_threshold
|
int
|
number of consecutive timeouts or
exceptions (across requests) after which a |
DEFAULT_CIRCUIT_BREAKER_THRESHOLD
|
Raises:
| Type | Description |
|---|---|
ProbingValueError
|
|
attach_log_sink(sink)
¶
Forward async extraction points' signals and results to sink.
sink is a LogSink (anything with
write_signal and write_result). Every async binding calls
sink.write_signal(...) for each non-None signal its probe returns,
and sink.write_result(...) once end_request finalizes it. Inline
extraction points don't forward: their signals already go to
route()'s caller.
Signals are forwarded from the binding's own worker thread, never from
route()'s caller, so a slow or raising sink can't add latency to
dispatch or affect other bindings. Results are forwarded inside
end_request, after the binding's queue has drained.
Takes effect immediately for every binding, present and future, so it
may be called before or after register_request. Pass None to
detach.
attach_metrics_sink(sink)
¶
Forward every async-binding metrics event to an external MetricsSink.
Use it for a Prometheus or StatsD backend. Events (queue depth, drops,
activation latency, probe errors) still go to the router's own in-memory
registry that get_metrics reads; this adds a second destination.
Like attach_log_sink: events are forwarded from each binding's worker
thread, never route()'s caller; a raising sink is caught, logged and
ignored; it takes effect immediately for every binding. Pass None to
detach.
get_metrics(request_id, extraction_point_name)
¶
Return the current MetricsSnapshot for one async binding.
Only valid while request_id is registered: end_request tears the
metrics down with everything else for that request.
Raises:
| Type | Description |
|---|---|
RouterError
|
|
register_request(request_id, extraction_points, request_ctx)
¶
Spawn a fresh probe instance for each extraction point and start it.
All-or-nothing: every extraction point is validated before any probe is
spawned, so an invalid one never leaves the request partially
registered. Prefer request, which
also guarantees end_request runs.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
request_id
|
str
|
unique among in-flight requests. |
required |
extraction_points
|
ProbeSpec | Iterable[ExtractionPoint]
|
a |
required |
request_ctx
|
RequestContext
|
passed to every probe's |
required |
Raises:
| Type | Description |
|---|---|
RouterError
|
the router is shut down, |
route(record)
¶
Dispatch one activation to the probe registered for its (request_id, extraction_point_name).
- Inline, default policy (
mode=reject): calls the probe'son_activationon this thread and returns itsProbeSignal(or None). An engine adapter uses this to decide whether to abort generation. - Inline,
block_until_signal: waits up totimeout_msfor the same call and returns theon_timeoutfallback signal if it doesn't finish in time or raises. Once the extraction point's circuit breaker has tripped (seecircuit_breaker_threshold), it is dispatched as plainrejectinstead, permanently. - Async: pushes the record onto the binding's queue (subject to its overflow policy) and returns None immediately.
Raises:
| Type | Description |
|---|---|
RouterError
|
the request isn't registered, the extraction point is unknown, or the record doesn't match its extraction point. |
end_request(request_id)
¶
Finalize every probe registered for request_id and tear down its state.
Each async binding's queue is closed and drained (backlogged activations
are processed, up to drain_timeout) before on_end is called, so a
trajectory probe's verdict reflects everything it saw, including for a
request that was aborted mid-stream. Then every
on_request_end listener is
called with (request_id, results).
Returns:
| Type | Description |
|---|---|
dict[str, ProbeResult]
|
|
Raises:
| Type | Description |
|---|---|
RouterError
|
|
request(extraction_points, *, request_id=None, prompt_metadata=None)
¶
Context manager for one request: registers it on entry and always calls end_request on exit.
with router.request(extraction_points=spec, prompt_metadata={"prompt": p}) as req:
for record in stream:
sig = req.route(record)
req.results # {extraction_point_name: ProbeResult}
end_request runs even when the body raises; see
RequestHandle for how errors are
reported.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
extraction_points
|
ProbeSpec | Iterable[ExtractionPoint]
|
a |
required |
request_id
|
str | None
|
defaults to a fresh |
None
|
prompt_metadata
|
Mapping[str, Any] | None
|
becomes |
None
|
on_request_end(listener)
¶
Call listener(request_id, results) at the end of every end_request.
This is how a caller that doesn't call end_request itself (an engine
adapter does it internally) still gets each request's results.
Listeners run on end_request's calling thread, outside the router's
internal lock (so a listener may call back into the router), in
registration order. Each gets its own shallow copy of results. A
raising listener is logged and never affects end_request or other
listeners.
Returns:
| Type | Description |
|---|---|
Callable[[], None]
|
A callable that removes this listener; calling it more than once is a no-op. |
get_probe(request_id, extraction_point_name)
¶
Look up the live probe instance for one (request, extraction point).
For introspection and testing; ordinary usage doesn't need it.
Raises:
| Type | Description |
|---|---|
RouterError
|
the request or extraction point isn't registered. |
shutdown(wait=True, timeout=None)
¶
Stop every still-registered request's async bindings and shut down the worker pool.
Call it once, when the host process is stopping; use end_request for
per-request cleanup. Requests still registered are dropped without
calling their probes' on_end. shutdown is terminal: afterwards
register_request, route and end_request raise
RouterError.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
wait
|
bool
|
|
True
|
timeout
|
float | None
|
with |
None
|
RequestHandle(router, request_id, extraction_points, prompt_metadata=None)
¶
One request's lifecycle on a Router. Create it via Router.request.
__enter__ registers the request; __exit__ calls end_request,
whether or not the body raised. If the body raised and end_request
then also raises, the end_request error is logged and the body's
exception propagates unchanged -- cleanup never masks the real error.
Single use: a handle can be entered once.
ended
property
¶
True once end_request has been attempted for this request.
results
property
¶
{extraction_point_name: ProbeResult}, available after the with block.
Raises RouterError while the request is still active, or if
end_request failed and so produced no results.
route(record)
¶
Same as Router.route, but first checks that record belongs to this request.
RequestEndListener = Callable[[str, 'dict[str, ProbeResult]'], None]
module-attribute
¶
listener(request_id, results), registered with Router.on_request_end.
RouterError
¶
Bases: ProbingError
Raised for invalid Router usage.
For example: an unknown probe_type or a duplicate request_id at
registration, routing a record for a request that isn't registered (or
has ended), or a record that doesn't match the extraction point it names.
Async execution¶
OverflowPolicy
¶
Bases: str, Enum
What an async extraction point's bounded queue does when it is full.
DROP_OLDEST = 'drop_oldest'
class-attribute
instance-attribute
¶
Evict the oldest queued activation to make room (the default).
DROP_NEWEST = 'drop_newest'
class-attribute
instance-attribute
¶
Drop the incoming activation.
BLOCK = 'block'
class-attribute
instance-attribute
¶
Block route() until there is room. Applies backpressure to generation.
default_worker_pool_size()
¶
The default worker_pool_size: min(32, os.cpu_count() * 4).
Pool threads mostly wait (on their binding's queue, or inside a probe
blocked on I/O or a GPU call), so the pool is oversubscribed relative to
the core count. Each async binding holds one thread for its request's
whole lifetime, so this number is how many async bindings are serviced
at once. If you expect more concurrently live async bindings, pass
Router(worker_pool_size=...) explicitly.
DEFAULT_QUEUE_DEPTH = 32
module-attribute
¶
Default Router(default_queue_depth=...): queue depth for async points that don't set queue_depth.
DEFAULT_DRAIN_TIMEOUT = 30.0
module-attribute
¶
Default Router(drain_timeout=...), in seconds.
DEFAULT_CIRCUIT_BREAKER_THRESHOLD = 5
module-attribute
¶
Default Router(circuit_breaker_threshold=...).
Metrics¶
MetricsSink
¶
Bases: ABC
Pluggable destination for async-binding metrics events.
Subclass it to export metrics (Prometheus, StatsD, ...) and attach it with
Router(metrics_sink=...) or Router.attach_metrics_sink. Every
method is called from the async binding's own worker thread, never from
Router.route()'s caller, so a slow or raising sink only delays that
binding's processing, never dispatch.
record_queue_depth(request_id, extraction_point_name, depth)
abstractmethod
¶
Current number of items sitting in the binding's queue, sampled
immediately after it changes (a route() enqueue or a worker
dequeue).
record_drop(request_id, extraction_point_name)
abstractmethod
¶
Called once per item the queue's overflow policy discarded --
either an eviction (drop_oldest) or a rejection (drop_newest).
Never called for a put() rejected merely because the queue was
already closed; that's shutdown behavior, not overflow.
record_activation(request_id, extraction_point_name, latency_seconds)
abstractmethod
¶
Called once per on_activation call the worker made (whether it
returned normally or raised), with the wall-clock time it took.
record_probe_error(request_id, extraction_point_name)
abstractmethod
¶
Called once per on_activation call that raised.
InMemoryMetricsRegistry()
¶
Bases: MetricsSink
Thread-safe in-memory counters and gauges, keyed by (request_id, extraction_point_name).
Every Router keeps one; Router.get_metrics reads it. Updates to
different bindings never contend with each other.
register(request_id, extraction_point_name)
¶
Seed a zeroed entry, so a query before the first activation returns zeros.
forget(request_id, extraction_point_name)
¶
Drop the entry for one binding (end_request does this, so metrics don't outlive the request).
snapshot(request_id, extraction_point_name)
¶
The current metrics for one binding, or None if it has none.
MetricsSnapshot(queue_depth, drop_count, activation_count, error_count, avg_activation_latency_seconds)
dataclass
¶
Point-in-time read of one async binding's metrics.
Attributes:
| Name | Type | Description |
|---|---|---|
queue_depth |
int
|
activations waiting in the binding's queue. |
drop_count |
int
|
activations the overflow policy discarded. |
activation_count |
int
|
|
error_count |
int
|
|
avg_activation_latency_seconds |
float | None
|
mean |