Skip to content

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 RedactFn applied to each record before it is written.

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 RedactFn applied to each record just before it is written.

None

path property

The file records are appended to.

write_raw(record)

Append an already-built record (still passed through redact, if set).

WebhookLogSink uses this for its dead-letter file.

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 post_fn(url, record) and expected to raise on failure. Useful for tests or a custom HTTP client.

None
redact RedactFn | None

a custom RedactFn, run after the default prompt-text redaction.

None
include_prompt_text bool

send the DEFAULT_PROMPT_TEXT_KEYS as-is (a custom redact then runs alone).

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.
  • Enum members -> .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 .shape and .dtype, e.g. a torch.Tensor or numpy.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.