Skip to content

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.

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:

  1. picks the record’s id and opens a call scope (so events recorded inside the tool can carry that id),
  2. awaits the rest of the chain and the tool,
  3. in a finally block, classifies the result, applies sampling, runs the identity and enrichment hooks, redacts, builds a ToolCallRecord, and hands it to the dispatcher,
  4. 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.

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.

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.

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.

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.

ModuleContents
instrumentation/middleware.pyInstrumentationMiddleware, instrument, session keys
instrumentation/dispatch.pyRecordDispatcher
instrumentation/record.py, events.pyToolCallRecord, EventRecord
instrumentation/sinks.pyJsonLinesSink, CallbackSink, MemorySink
instrumentation/storage.pyDatabaseSink, build_metadata
instrumentation/embeddings.pyEmbeddingSink, OpenAIEmbedder
instrumentation/correlation.pyFeedback-to-call linking rules
instrumentation/classify.py, redaction.pySoft-error classifier, Redactor
tools.py, mixins/The feedback tools