undercurrent.sinks¶
Advanced API
For a simple results hook, pass on_result= to
ProbedModel. Sinks are for shipping
observations out of process from a Router.
Observation sinks receive the signals and final results of async
(observe-mode) extraction points. The sinks and redaction helpers are also
exported from undercurrent.
Sinks¶
LogSink(*, redact=None)
¶
Bases: ABC
Out-of-band destination for a probe's intermediate signals and final verdict.
A router calls a sink for extraction points in observe mode
(execution_mode=async), whose output shouldn't block the response.
Attach one with
Router.attach_log_sink
(or ProbedModel(log_sink=...)). The router calls write_signal /
write_result from an async binding's own worker thread, so neither
should block for long: a fast local append is fine, but a network call
should hand off to a background thread, as
WebhookLogSink does.
To write a custom sink, subclass LogSink, implement the two abstract
methods, and convert payloads with
to_jsonable. To support a redact=
argument, call super().__init__(redact=redact) and pass every record
through self._apply_redaction(record) just before writing it (it
returns None when the record must be dropped).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
redact
|
RedactFn | None
|
an optional |
None
|
redaction_error_count
property
¶
Records dropped because the redact function raised (or returned a non-dict).
write_signal(request_id, extraction_point_name, signal)
abstractmethod
¶
Called whenever a probe emits an intermediate ProbeSignal,
e.g. a trajectory probe's running-score update from on_activation.
write_result(request_id, extraction_point_name, result)
abstractmethod
¶
Called once, after a probe's on_end finalizes its ProbeResult.
FileLogSink(path, *, redact=None)
¶
Bases: LogSink
Appends one NDJSON line per signal/result to path.
Each line is a JSON object with kind ("signal"/"result"),
request_id, extraction_point_name, timestamp and payload
(the signal or result, converted with
to_jsonable, so tensors are
summarized rather than dumped).
Writes are serialized, so lines from different async bindings never
interleave. Each write opens, appends to and closes the file, so there
is nothing to close; path's parent directory must exist.
Unlike WebhookLogSink, a
FileLogSink keeps data on the local machine and doesn't redact prompt
text by default.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str | Path
|
the NDJSON file to append to. |
required |
redact
|
RedactFn | None
|
an optional |
None
|
WebhookLogSink(url, *, max_retries=2, backoff_base=0.1, request_timeout=5.0, dead_letter_path=None, queue_maxsize=1000, post_fn=None, redact=None, include_prompt_text=False)
¶
Bases: LogSink
POSTs each signal/result as JSON to url from a background thread.
write_signal / write_result only enqueue and return, so a slow or
unreachable endpoint never adds latency to the caller.
A failed POST is retried max_retries times (3 attempts in total by
default) with exponential backoff (backoff_base * 2**attempt
seconds). A record that exhausts its retries is written to
dead_letter_path (NDJSON) if one was given, and otherwise logged with
the standard logging module. Errors are never raised to the caller.
If the bounded queue fills up (the endpoint is down and records are
waiting on retries), new records are dropped and counted in
dropped_count.
Because records leave the machine, every key in
DEFAULT_PROMPT_TEXT_KEYS
(prompt, text, generated_text, ...) is replaced with
"[REDACTED]" by default, wherever it appears. Redaction runs before
the first send attempt, so the dead-letter file and the local error log
only ever see the redacted record.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str
|
the endpoint to POST to. |
required |
max_retries
|
int
|
retries after the first failed attempt. |
2
|
backoff_base
|
float
|
base of the exponential backoff, in seconds. |
0.1
|
request_timeout
|
float
|
per-request timeout, in seconds. |
5.0
|
dead_letter_path
|
str | Path | None
|
NDJSON file for records that exhaust their retries. |
None
|
queue_maxsize
|
int
|
capacity of the background queue. |
1000
|
post_fn
|
PostFn | None
|
replaces the HTTP POST: called as |
None
|
redact
|
RedactFn | None
|
a custom |
None
|
include_prompt_text
|
bool
|
send the |
False
|
dropped_count
property
¶
Records discarded because the background queue was full. For tests/introspection.
close(timeout=None)
¶
Stop the background worker after everything already queued has been sent (or dead-lettered). Mainly for tests/clean shutdown -- the worker thread is a daemon, so it won't keep the process alive on its own.
With a timeout (seconds), returns within it even if records are
still being sent; without one, waits for the queue to drain.
wire_router(router, sink)
¶
Attach sink to router; the same as router.attach_log_sink(sink).
Every async ("observe mode") binding's signals and final results are then
forwarded to sink. See
Router.attach_log_sink.
Redaction¶
redact_keys(*keys, replacement=DEFAULT_REPLACEMENT)
¶
Replace the value of every matching key with replacement.
keys are plain key names (matched at any depth) or dotted paths
("payload.metadata.prompt", matched wherever that chain of keys
appears). The key itself stays in the record, so consumers can see that
something was removed.
drop_keys(*keys)
¶
Remove every matching key (and its value) from the record.
Key matching works exactly as in redact_keys.
chain(*fns)
¶
Compose redaction functions left to right. If any returns None, the record is dropped.
RedactFn = Callable[[dict[str, Any]], dict[str, Any] | None]
module-attribute
¶
A redaction hook: takes a JSON-able record, returns it (possibly changed) or None to drop it.
DEFAULT_PROMPT_TEXT_KEYS = ('prompt', 'prompts', 'prompt_text', 'text', 'input_text', 'output_text', 'generated_text', 'completion', 'response', 'messages', 'input_ids', 'output_ids', 'token_ids')
module-attribute
¶
Keys that can carry prompt or generated text, redacted by WebhookLogSink by default.
Writing a custom sink¶
to_jsonable(value)
¶
Recursively convert value into something json.dumps can handle.
- dataclass instances -> dict of their fields, recursively converted.
Enummembers ->.value.- dict -> dict with keys coerced to
str(JSON object keys must be strings) and values recursively converted. - list/tuple/set -> list, recursively converted.
- tensor-like objects (anything with both
.shapeand.dtype, e.g. atorch.Tensorornumpy.ndarray) -> a small summary dict, never the raw values. bytes/bytearray-> summarized the same way as a tensor would be, since dumping raw bytes into a log line is just as undesirable.- anything else JSON-native (str/int/float/bool/None) -> passed through.
- anything else ->
str(value), as a last-resort fallback rather than raising.