Metrics¶
The router measures every async binding: one probe instance for one
extraction point in one request. For each binding it records queue depth,
overflow drops, on_activation latency and probe errors. You can read the
numbers in two ways:
- Pull: call
router.get_metrics(request_id, extraction_point_name)for a point-in-timeMetricsSnapshot. This always works and needs no setup. - Push: attach your own
MetricsSinkwithrouter.attach_metrics_sink(sink)and forward each event to Prometheus, OpenTelemetry, StatsD or any other backend.
Inline extraction points aren't measured. Their probe runs synchronously in
route(), so you can time it yourself at the call site. All the names on
this page come from undercurrent.router.metrics and undercurrent.router.
The events: MetricsSink¶
MetricsSink is an abstract base class with four methods. Every call carries
the binding's request_id and extraction_point_name.
| Method | When it's called | What it means |
|---|---|---|
record_queue_depth(request_id, extraction_point_name, depth) |
Right after each change to the queue: after route() enqueues an activation, and after the worker takes one off. |
A gauge: how many activations are waiting for this probe. |
record_drop(request_id, extraction_point_name) |
Once for every activation that the queue's overflow policy discards. Under drop_oldest that is the evicted oldest item; under drop_newest it's the rejected new item. It's not called for a put that fails because the queue was already closed during shutdown. |
A counter: activations the probe never saw. |
record_activation(request_id, extraction_point_name, latency_seconds) |
Once for every on_activation call the worker made, whether it returned or raised. |
Wall-clock seconds spent inside on_activation. Sink forwarding and queue wait aren't included. |
record_probe_error(request_id, extraction_point_name) |
Once for every on_activation call that raised. |
A counter: probe failures. Each one also sends a synthetic error signal to the log sink. |
So the error rate of a binding is probe errors / activations, and the
drop rate is drops / (activations + drops + queue depth), roughly the
activations that were routed.
Which thread calls what¶
Your sink is called from two places:
record_activationandrecord_probe_error, plus the queue-depth reading after each dequeue, come from the binding's own worker thread.- The queue-depth reading after an enqueue, and
record_drop, come from the thread that calledroute(), which is the engine's generation path.record_dropis even called while the binding's queue lock is held.
Your MetricsSink must therefore be:
- Thread-safe. Many bindings call it concurrently.
- Fast and non-blocking. Increment a counter or set a gauge, and do no I/O. If you need to send data over the network, do it from a separate exporter thread, the way Prometheus' pull model and OpenTelemetry's periodic readers already do.
- Free of calls back into the router.
If your sink raises, the router catches the exception, logs it with
logger.exception, and carries on. Dispatch and probe processing are never
affected, but the event is lost for your sink. The built-in registry still
counts it.
The built-in registry and get_metrics¶
Every Router builds an InMemoryMetricsRegistry and always records into
it, whether or not you attach a sink. It can't be replaced; it exists so
that get_metrics works with no setup.
router.get_metrics(request_id, extraction_point_name) returns a frozen
MetricsSnapshot:
| Field | Type | Meaning |
|---|---|---|
queue_depth |
int |
The last value from record_queue_depth. |
drop_count |
int |
The number of record_drop calls. |
activation_count |
int |
The number of record_activation calls (completed on_activation calls). |
error_count |
int |
The number of record_probe_error calls. |
avg_activation_latency_seconds |
float or None |
The mean of the latency_seconds values, or None before the first activation. |
The registry only keeps state per request:
- Each binding's entry is created zeroed by
register_requestand deleted byend_request(and byshutdown). get_metricsraisesRouterErrorfor an unknown or already-endedrequest_id, an unknown extraction point name, or an inline extraction point.- Nothing is aggregated across requests.
If you need totals across requests or over time, which is what a dashboard
needs, attach a MetricsSink and aggregate there.
Example: watch a slow probe drop activations¶
The probe below blocks on its first activation until the example releases
it. While it's blocked, the example routes more activations than its queue
can hold, so the drop_oldest overflow policy has to evict some of them. A
tiny counting MetricsSink keeps totals that outlive the request.
import collections
import threading
from undercurrent.core import Probe, ProbeFactory, ProbeResult
from undercurrent.router import MetricsSink, Router
from undercurrent.spec import ActivationRecord, parse_yaml
SPEC = parse_yaml("""
version: "1"
extraction_points:
- name: slow_probe
layers: [4]
tensor_type: residual_stream
position: "generated[*]"
probe_type: slow
probe_kind: trajectory
execution_mode: async
queue_depth: 4
""")
started = threading.Event()
release = threading.Event()
class SlowProbe(Probe):
"""Blocks on its first activation until released; raises on negative inputs."""
probe_kind = "trajectory"
def on_start(self, request_ctx):
self.seen = 0
def on_activation(self, record):
started.set()
release.wait(timeout=10)
if record.tensor[0] < 0:
raise ValueError("negative activation")
self.seen += 1
return None
def on_end(self, request_ctx):
return ProbeResult(self.request_id, self.extraction_point_name, verdict={"seen": self.seen})
class CountingMetricsSink(MetricsSink):
"""Totals per extraction point, across all requests. Thread-safe and O(1) per event."""
def __init__(self):
self._lock = threading.Lock()
self.counters = collections.Counter()
self.max_queue_depth = collections.defaultdict(int)
def record_queue_depth(self, request_id, extraction_point_name, depth):
with self._lock:
self.max_queue_depth[extraction_point_name] = max(self.max_queue_depth[extraction_point_name], depth)
def record_drop(self, request_id, extraction_point_name):
with self._lock:
self.counters[(extraction_point_name, "drops")] += 1
def record_activation(self, request_id, extraction_point_name, latency_seconds):
with self._lock:
self.counters[(extraction_point_name, "activations")] += 1
def record_probe_error(self, request_id, extraction_point_name):
with self._lock:
self.counters[(extraction_point_name, "errors")] += 1
def activation(request_id, pos, value):
return ActivationRecord(
request_id=request_id,
extraction_point_name="slow_probe",
layer=4,
token_pos=pos,
tensor_type="residual_stream",
tensor=[value],
is_generated=True,
)
metrics = CountingMetricsSink()
with Router(probe_registry={"slow": ProbeFactory(SlowProbe)}) as router:
router.attach_metrics_sink(metrics) # or Router(..., metrics_sink=metrics)
with router.request(SPEC, request_id="req-1") as req:
req.route(activation("req-1", 0, 1.0))
started.wait(timeout=10) # the worker is now stuck inside on_activation
# 9 more activations into a queue of 4: the 5 oldest are evicted.
for pos in range(1, 10):
req.route(activation("req-1", pos, -1.0 if pos == 9 else 1.0))
snap = router.get_metrics("req-1", "slow_probe")
print(snap)
assert (snap.queue_depth, snap.drop_count, snap.activation_count) == (4, 5, 0)
assert snap.avg_activation_latency_seconds is None # nothing has finished yet
release.set() # let the probe catch up; leaving the block drains the queue
# The per-request entry is gone now; the external sink kept the totals.
print(dict(metrics.counters), dict(metrics.max_queue_depth))
assert metrics.counters[("slow_probe", "drops")] == 5
assert metrics.counters[("slow_probe", "activations")] == 5 # the first one plus the 4 that survived
assert metrics.counters[("slow_probe", "errors")] == 1 # the last activation was negative
assert metrics.max_queue_depth["slow_probe"] == 4
assert req.results["slow_probe"].verdict == {"seen": 4}
The code prints the following, plus a logged traceback for the deliberate
ValueError:
MetricsSnapshot(queue_depth=4, drop_count=5, activation_count=0, error_count=0, avg_activation_latency_seconds=None)
{('slow_probe', 'drops'): 5, ('slow_probe', 'activations'): 5, ('slow_probe', 'errors'): 1} {'slow_probe': 4}
get_metrics is useful for tests, debugging endpoints and adaptive logic
inside a request. For anything longer-lived, use a MetricsSink.
Sketch: export to Prometheus¶
Sketch, not shipped
Undercurrent doesn't include a Prometheus or OpenTelemetry exporter, and
doesn't depend on prometheus_client. The code below is an illustrative
starting point. Adapt and test it in your own stack.
Don't use request_id as a label. Every request creates new bindings,
and a per-request label set would grow without bound. Label by
extraction_point_name only, since it comes from your spec and has a fixed
set of values.
from prometheus_client import Counter, Gauge, Histogram, start_http_server
from undercurrent.router import MetricsSink
ACTIVATIONS = Counter("undercurrent_probe_activations_total", "on_activation calls completed", ["extraction_point"])
ERRORS = Counter("undercurrent_probe_errors_total", "on_activation calls that raised", ["extraction_point"])
DROPS = Counter("undercurrent_queue_drops_total", "Activations discarded by the overflow policy", ["extraction_point"])
LATENCY = Histogram(
"undercurrent_probe_activation_seconds",
"Wall-clock time inside on_activation",
["extraction_point"],
buckets=(0.0005, 0.001, 0.0025, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 1.0),
)
QUEUE_DEPTH = Gauge(
"undercurrent_queue_depth",
"Most recent queue depth reported by any binding of this extraction point",
["extraction_point"],
)
class PrometheusMetricsSink(MetricsSink):
# prometheus_client metrics are thread-safe and in-memory, so every method is cheap.
def record_queue_depth(self, request_id, extraction_point_name, depth):
QUEUE_DEPTH.labels(extraction_point_name).set(depth)
def record_drop(self, request_id, extraction_point_name):
DROPS.labels(extraction_point_name).inc()
def record_activation(self, request_id, extraction_point_name, latency_seconds):
ACTIVATIONS.labels(extraction_point_name).inc()
LATENCY.labels(extraction_point_name).observe(latency_seconds)
def record_probe_error(self, request_id, extraction_point_name):
ERRORS.labels(extraction_point_name).inc()
start_http_server(9400) # /metrics on its own thread
router.attach_metrics_sink(PrometheusMetricsSink())
QUEUE_DEPTH is "last write wins" across all concurrent requests for that
extraction point. It's a rough signal, good enough to spot a backlog. For
an exact per-binding maximum, keep a dict keyed by
(request_id, extraction_point_name) in the sink, remove keys you haven't
seen for a while, and export the maximum as the gauge.
With OpenTelemetry the shape is the same: create Counters and a
Histogram from a Meter, call .add(1, {"extraction_point": name}) and
.record(latency, {...}) in the four methods, and let the SDK's periodic
reader export them. For queue depth, use an UpDownCounter, or an
observable gauge that reads a dict your sink maintains.
Suggested alerts¶
Use these as starting points and tune the thresholds to your traffic. For how queue depth, overflow policies and the worker pool interact, see Async execution & backpressure.
| Alert | Example PromQL (with the sketch above) | What it usually means |
|---|---|---|
| Sustained drops | sum by (extraction_point) (rate(undercurrent_queue_drops_total[5m])) > 0 for 10m |
The probe can't keep up with the activation rate. Its results are computed from a subsample. Speed up the probe, add stride to the spec, raise queue_depth, or raise worker_pool_size if the bindings are waiting for a worker. |
| Drop ratio | rate(undercurrent_queue_drops_total[5m]) / (rate(undercurrent_queue_drops_total[5m]) + rate(undercurrent_probe_activations_total[5m])) > 0.01 |
Same as above, normalized for traffic. Use it to page if verdict quality depends on seeing every activation. |
| Queue depth near the limit | max by (extraction_point) (undercurrent_queue_depth) >= 0.8 * <queue_depth from your spec> for 5m |
Drops are imminent. Usually a slow probe or a slow custom log sink. Activation latency tells you which: if it's flat while depth climbs, look at the sink. |
| Probe error rate | rate(undercurrent_probe_errors_total[5m]) / rate(undercurrent_probe_activations_total[5m]) > 0.01 |
The probe is raising. The error signals in your observation sink carry error_type and the message. |
| Activation latency | histogram_quantile(0.99, sum by (le, extraction_point) (rate(undercurrent_probe_activation_seconds_bucket[5m]))) above your budget |
The probe is getting slower. Queue depth and drops follow if the latency doesn't fit between tokens. |
The router doesn't emit metrics for inline extraction points, intervention timeouts or circuit-breaker trips. Timeouts and trips are logged instead. See Intervention policies & timeouts.