Architecture
The package has two halves that share a database at most: the feedback tools, which are ordinary FastMCP tools backed by synchronous SQLAlchemy, and the instrumentation, a FastMCP middleware that records every tool call on the server. This page is about the second half, since that is where the design decisions are.
flowchart LR
C[MCP client] --> M[middleware]
M --> T[your tool]
T --> M
M --> Q[(record queue)]
E[record_event] --> Q
Q --> D[dispatcher task]
D --> S1[JsonLinesSink]
D --> S2[DatabaseSink]
D --> S3[EmbeddingSink]
S3 --> Q2[(text queue)]
Q2 --> X[embedder]
S2 --> DB[(ffb_ tables)]
X --> DB
The middleware hands each record to a bounded queue and never waits on it.
Events from record_event, on any thread, join the same queue. A dispatcher
task drains it into the sinks, and the embedding sink keeps a second queue of
its own in front of the embedder.
The middleware
Section titled “The middleware”instrument(app) adds one InstrumentationMiddleware with
app.add_middleware. FastMCP calls its on_call_tool for every tool call, so
every tool on the server is covered, including tools registered later, mounted
servers’ tools, and the feedback tools themselves.
For each call it:
- picks the record’s id and opens a call scope (so events recorded inside the tool can carry that id),
- awaits the rest of the chain and the tool,
- in a
finallyblock, classifies the result, applies sampling, runs the identity and enrichment hooks, redacts, builds aToolCallRecord, and hands it to the dispatcher, - returns the tool’s result or re-raises its exception, unchanged.
Step 3 is wrapped so that any failure in it is logged and swallowed. The client receives exactly what the tool produced.
The dispatcher
Section titled “The dispatcher”RecordDispatcher owns an asyncio.Queue with max_queue slots (1000) and a
background task that drains it in batches of up to batch_size (100), handing
each batch to every sink in turn.
submit uses put_nowait. When the queue is full the record is dropped and
counted in dispatcher.dropped, with a warning on the first drop and every
thousandth. A tool call never waits for room.
Each sink write is bounded by sink_timeout (10 seconds). A sink that raises
or times out is logged and counted in dispatcher.sink_errors; the other sinks
still get the batch, and the next batch is delivered as usual.
Events take the same path. record_event is synchronous; from a thread other
than the server’s loop it hands the record over with call_soon_threadsafe.
The guarantees, and their cost
Section titled “The guarantees, and their cost”Never block. Recording costs the tool call the work of building a record:
hooks (each capped at hook_timeout, 0.25 seconds), redaction, and JSON-size
measurement. Delivery happens later, off the call path.
Never raise. Hook errors, classifier errors, sink errors, a full queue, a missing event loop: each is logged and counted, none reaches the client.
The cost of these guarantees is that records can be lost. A full queue drops
them, a crash loses whatever was still queued, and a sink that is down loses
the batches it failed to write. The package chooses a working server with gaps
in its records over a server that slows down or fails because its telemetry
did. Watch dispatcher.dropped and dispatcher.sink_errors if gaps matter to
you, and call await mw.aclose() on shutdown so the queue is flushed.
Embedding stays off the call path twice
Section titled “Embedding stays off the call path twice”An EmbeddingSink is a sink like any other, but embedding means a network
round trip that can take seconds. So its write only extracts text and puts it
on a second bounded queue of its own, then returns. Its own background task
calls the embedder and writes ffb_embeddings. A slow embedding endpoint
therefore fills the embedding queue and drops texts there, without delaying the
dispatcher or the other sinks.
Lazy imports
Section titled “Lazy imports”The base install has no async database stack. DatabaseSink, EmbeddingSink
and the backfill functions load on first access; the middleware recognizes them
among its sinks without importing their modules. That is why
instrument(app) with JsonLinesSink works on a base install, and why the
database helpers return empty results rather than failing when the extra is
missing.
Where things live
Section titled “Where things live”| Module | Contents |
|---|---|
instrumentation/middleware.py | InstrumentationMiddleware, instrument, session keys |
instrumentation/dispatch.py | RecordDispatcher |
instrumentation/record.py, events.py | ToolCallRecord, EventRecord |
instrumentation/sinks.py | JsonLinesSink, CallbackSink, MemorySink |
instrumentation/storage.py | DatabaseSink, build_metadata |
instrumentation/embeddings.py | EmbeddingSink, OpenAIEmbedder |
instrumentation/correlation.py | Feedback-to-call linking rules |
instrumentation/classify.py, redaction.py | Soft-error classifier, Redactor |
tools.py, mixins/ | The feedback tools |