Skip to content

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 probe_type names are looked up. None (the default) uses the global registry that @register_probe fills (plus undercurrent.probes entry-point plugins), consulted at register_request time, so probes registered after the router is built still resolve. A ProbeRegistry is used as-is. A mapping of probe_type -> ProbeFactory or Probe subclass is copied and is authoritative: no fallback to the global registry or to plugins.

None
worker_pool_size int | None

number of threads backing all async bindings, shared across every request. Defaults to default_worker_pool_size().

None
default_queue_depth int

queue depth for an async extraction point whose queue_depth is None.

DEFAULT_QUEUE_DEPTH
default_overflow_policy OverflowPolicy

what a full async queue does; see OverflowPolicy.

DROP_OLDEST
drain_timeout float

seconds end_request waits for an async binding's queue to drain before finalizing anyway (logged, not raised).

DEFAULT_DRAIN_TIMEOUT
metrics_sink MetricsSink | None

an optional external MetricsSink, the same as calling attach_metrics_sink right after construction. get_metrics works with or without one.

None
default_intervention_policy InterventionPolicy | None

policy for an extraction point whose intervention is None. Defaults to InterventionPolicy() (mode=reject, no waiting).

None
circuit_breaker_threshold int

number of consecutive timeouts or exceptions (across requests) after which a block_until_signal extraction point is permanently downgraded to mode=reject by this router.

DEFAULT_CIRCUIT_BREAKER_THRESHOLD

Raises:

Type Description
ProbingValueError

worker_pool_size < 1.

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

request_id isn't registered (or has ended), extraction_point_name is unknown, or the extraction point is inline (nothing is queued or measured for those).

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 ProbeSpec or any iterable of ExtractionPoint.

required
request_ctx RequestContext

passed to every probe's on_start and on_end.

required

Raises:

Type Description
RouterError

the router is shut down, request_id is already registered, a probe_type can't be resolved, or an extraction point is invalid for the router.

route(record)

Dispatch one activation to the probe registered for its (request_id, extraction_point_name).

  • Inline, default policy (mode=reject): calls the probe's on_activation on this thread and returns its ProbeSignal (or None). An engine adapter uses this to decide whether to abort generation.
  • Inline, block_until_signal: waits up to timeout_ms for the same call and returns the on_timeout fallback signal if it doesn't finish in time or raises. Once the extraction point's circuit breaker has tripped (see circuit_breaker_threshold), it is dispatched as plain reject instead, 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]

{extraction_point_name: ProbeResult}.

Raises:

Type Description
RouterError

request_id isn't registered or has already ended.

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 ProbeSpec or any iterable of ExtractionPoint.

required
request_id str | None

defaults to a fresh uuid4().hex.

None
prompt_metadata Mapping[str, Any] | None

becomes RequestContext.prompt_metadata (an empty dict when omitted).

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 (default): close every live async binding's queue and drain it (like end_request does), and block until every worker has exited. False: discard whatever is queued and return promptly; an activation already in progress still runs to completion (a Python thread can't be preempted), and its worker exits after it.

True
timeout float | None

with wait=True, seconds to wait for each binding to drain (each binding gets its own budget). Defaults to the router's drain_timeout.

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

on_activation calls made (including ones that raised).

error_count int

on_activation calls that raised.

avg_activation_latency_seconds float | None

mean on_activation wall-clock time, or None before the first call.