Sinks and embedders
All sinks are imported from fastmcp_feedback.instrumentation. DatabaseSink
needs the instrumentation extra; EmbeddingSink and OpenAIEmbedder need
instrumentation,embeddings (install).
The Sink protocol
Section titled “The Sink protocol”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 FastMCPfrom 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
Section titled “JsonLinesSink”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
Section titled “CallbackSink”CallbackSink(callback)Passes each batch to callback(records), sync or async. A sync callback runs on
the event loop, so keep it quick.
MemorySink
Section titled “MemorySink”MemorySink()Keeps records in .records (arrival order); .calls and .events filter by
kind. Meant for tests.
DatabaseSink
Section titled “DatabaseSink”DatabaseSink(engine, *, prefix="ffb_", create_tables=False, retention=None, prune_interval=timedelta(hours=1), prune_batch=5000, embedding_dim=1024)| Parameter | Description |
|---|---|
engine | An 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. |
prefix | Table name prefix. |
create_tables | Create missing tables on first use. Never alters existing tables. ffb_embeddings is created only when an EmbeddingSink first writes. |
retention | A positive timedelta. When set, rows older than this are pruned automatically. None keeps everything. |
prune_interval | Minimum time between automatic prunes. |
prune_batch | Rows deleted per transaction. At least 1. |
embedding_dim | Vector 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
Section titled “EmbeddingSink”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)| Parameter | Description |
|---|---|
embedder | Anything with model, dim and async embed(texts). |
database_sink | The DatabaseSink whose engine and tables to use. Its embedding_dim must equal embedder.dim, or ValueError. |
sources | Which source types to embed. Unknown names raise ValueError. |
event_text_keys | Event attributes whose values make up an event text. |
max_chars | Longest text embedded and stored. |
max_queue | Texts waiting before new ones are dropped and counted. |
batch_size | Texts per embed call. |
timeout | Seconds one embed call may take, and how long aclose() waits for the queue to drain. |
scan_limit | Without pgvector, how many of the newest rows similar compares against. |
redactor | Defaults 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
Section titled “OpenAIEmbedder”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.
Embedder protocol
Section titled “Embedder protocol”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.