Skip to content

Sinks and embedders

All sinks are imported from fastmcp_feedback.instrumentation. DatabaseSink needs the instrumentation extra; EmbeddingSink and OpenAIEmbedder need instrumentation,embeddings (install).

A sink is anything with async write(records). async aclose() is optional and is called by mw.aclose().

from collections.abc import Sequence
from fastmcp import FastMCP
from fastmcp_feedback.instrumentation import EventRecord, ToolCallRecord, instrument
class SlowCallsSink:
"""Print calls slower than a threshold."""
def __init__(self, threshold_ms: float):
self.threshold_ms = threshold_ms
async def write(self, records: Sequence[ToolCallRecord | EventRecord]) -> None:
for r in records:
if r.record_type == "tool_call" and r.duration_ms > self.threshold_ms:
print(f"slow: {r.tool} took {r.duration_ms:.0f} ms")
instrument(FastMCP("My Server"), [SlowCallsSink(threshold_ms=500)])

Batches mix calls and events. Each write is limited to the middleware’s sink_timeout; an exception or timeout is logged and counted in mw.dispatcher.sink_errors, and the next batch is delivered as usual.

JsonLinesSink(stream=None)

One JSON object per record, with a "record" field of "tool_call" or "event". stream defaults to sys.stderr, looked up on every write. The default sink when instrument gets no sinks.

CallbackSink(callback)

Passes each batch to callback(records), sync or async. A sync callback runs on the event loop, so keep it quick.

MemorySink()

Keeps records in .records (arrival order); .calls and .events filter by kind. Meant for tests.

DatabaseSink(engine, *, prefix="ffb_", create_tables=False, retention=None,
prune_interval=timedelta(hours=1), prune_batch=5000, embedding_dim=1024)
ParameterDescription
engineAn AsyncEngine, or an async URL such as sqlite+aiosqlite:///calls.db or postgresql+asyncpg://.... An engine built from a URL is disposed on aclose(); a passed-in engine is left alone.
prefixTable name prefix.
create_tablesCreate missing tables on first use. Never alters existing tables. ffb_embeddings is created only when an EmbeddingSink first writes.
retentionA positive timedelta. When set, rows older than this are pruned automatically. None keeps everything.
prune_intervalMinimum time between automatic prunes.
prune_batchRows deleted per transaction. At least 1.
embedding_dimVector size of ffb_embeddings; must equal the embedder’s dim.

Methods (async): write(records), prune(older_than=None) -> int, ensure_tables(), events_for(key, *, kinds=None, limit=1000), recent_calls(...), replace_links(...), linked_calls(feedback_ref), aclose(). Attributes: engine, metadata, table (calls), events, links, embeddings.

write stores calls and events in separate transactions, calls first. With retention set, it prunes after the insert at most once per prune_interval. prune deletes calls by started_at (except calls linked to feedback), events by occurred_at, and embeddings by created_at (except feedback embeddings), and returns the total deleted.

EmbeddingSink(embedder, database_sink, *, sources=("feedback", "call_error", "event", "llm"),
event_text_keys=("error_excerpt", "error", "message", "reason"),
max_chars=4000, max_queue=1000, batch_size=32, timeout=30.0,
scan_limit=5000, redactor=None)
ParameterDescription
embedderAnything with model, dim and async embed(texts).
database_sinkThe DatabaseSink whose engine and tables to use. Its embedding_dim must equal embedder.dim, or ValueError.
sourcesWhich source types to embed. Unknown names raise ValueError.
event_text_keysEvent attributes whose values make up an event text.
max_charsLongest text embedded and stored.
max_queueTexts waiting before new ones are dropped and counted.
batch_sizeTexts per embed call.
timeoutSeconds one embed call may take, and how long aclose() waits for the queue to drain.
scan_limitWithout pgvector, how many of the newest rows similar compares against.
redactorDefaults to the middleware’s when passed to instrument, else Redactor().

Counters: embedded, skipped (same text already stored), dropped, errors. Methods (async): write, flush(), aclose(), similar(text, ...), nearest(vector, ..., exclude=None), vector_for(source_type, source_id). EmbeddingSink.similar raises on failure; the middleware’s similar wraps it and never raises.

OpenAIEmbedder(base_url, api_key=None, model="mxbai-embed-large", dim=1024,
batch_size=32, timeout=30.0, transport=None)

Posts {"model": model, "input": [...]} to {base_url}/embeddings. base_url goes up to and including the version, such as https://api.example.com/v1. api_key is sent as Authorization: Bearer ... when given. A vector of any length other than dim, or a response with the wrong number of vectors, raises ValueError. Uses httpx (FastMCP 2 and 3) or httpx2 (FastMCP 4), whichever is installed; no extra dependency.

from fastmcp_feedback.instrumentation import Embedder
class HashEmbedder:
"""A deterministic stand-in, for tests: not a meaningful embedding."""
model = "hash-8"
dim = 8
async def embed(self, texts: list[str]) -> list[list[float]]:
return [[(hash(t) >> i) % 7 + 1.0 for i in range(self.dim)] for t in texts]
assert isinstance(HashEmbedder(), Embedder)

embed returns one vector of length dim per text, in order. Only vectors made by the same model are compared in a search.