Set up embeddings and similarity search
“Has anyone reported this before?” is hard to answer with LIKE. An
EmbeddingSink turns the text in your records into vectors and stores them next
to the records, so you can search by meaning.
What gets embedded
Section titled “What gets embedded”source_type | Text | source_id |
|---|---|---|
feedback | title and description of a feedback.submitted event | the feedback id |
call_error | "{tool}: {error_type}: {error_message}" of calls with outcome error or soft_error | the call id |
event | the values of error_excerpt, error, message and reason in an event’s attrs | the event id |
llm | "prompt: ...\ncompletion: ..." of an llm.call event, when capture_llm_text=True | the event id |
Successful calls are not embedded. List values are embedded one item per line
and dict values as sorted key: value lines, rather than as Python reprs.
Install
Section titled “Install”uv add "fastmcp-feedback[instrumentation,embeddings]" asyncpgThe embeddings extra adds pgvector, needed only on PostgreSQL. The database
needs the vector extension; if your Postgres runs on Alpine, see
Add pgvector to an existing Alpine Postgres. SQLite
works too, with the search done in Python (see Search).
Wire it up
Section titled “Wire it up”OpenAIEmbedder talks to any OpenAI-compatible /v1/embeddings endpoint
(OpenAI, Ollama, LiteLLM, vLLM). The sink needs a DatabaseSink to store into,
and both must agree on the vector size:
import os
from fastmcp import FastMCPfrom fastmcp_feedback import add_feedback_toolsfrom fastmcp_feedback.instrumentation import ( DatabaseSink, EmbeddingSink, OpenAIEmbedder, instrument,)
app = FastMCP("Render Farm")
embedder = OpenAIEmbedder( os.environ["EMBED_BASE_URL"], # e.g. https://api.example.com/v1 api_key=os.environ.get("FFB_EMBED_API_KEY"), model="mxbai-embed-large", dim=1024,)db = DatabaseSink(os.environ["DATABASE_URL"], create_tables=True, embedding_dim=1024)mw = instrument(app, [db, EmbeddingSink(embedder, db)])add_feedback_tools(app, database_url="sqlite:///feedback.db", instrumentation=mw)With create_tables=True, ffb_embeddings is created when the embedding sink
first writes. On PostgreSQL that step first runs
CREATE EXTENSION IF NOT EXISTS vector, which needs a superuser or the database
owner on most setups. Without that privilege it logs a warning, keeps running,
counts the failed writes and retries at most once a minute; a DBA can create
the extension once and the next attempt succeeds.
If embedding_dim does not equal the embedder’s dim, EmbeddingSink raises
ValueError at construction.
Generate something to find
Section titled “Generate something to find”import asyncio
from fastmcp import Client
@app.tooldef render(scene: str) -> dict: if scene == "city": return {"status": "failed", "error": "GPU out of memory while loading textures"} return {"status": "done"}
async def generate(): async with Client(app) as client: await client.call_tool("render", {"scene": "city"}) await client.call_tool("submit_feedback", {"request": { "type": "bug", "title": "Renders of large scenes crash", "description": "The city scene stops with a memory error on the GPU.", "submitter": "docs-example", }}) await mw.flush() # records reach the sinks await mw.embedding_sink.flush() # queued texts are embedded and storedSearch
Section titled “Search”async def search(): for hit in await mw.similar("video card ran out of memory", k=5): first_line = hit["text"].splitlines()[0] print(f"{hit['score']:.2f}", hit["source_type"], hit["source_id"][:8], first_line)
# reports and errors like feedback 1, leaving 1 itself out related = await mw.similar_feedback("1", source_types=["call_error"]) print([hit["text"] for hit in related])
async def main(): await generate() await search() await mw.aclose()
asyncio.run(main())0.73 feedback 1 Renders of large scenes crash0.70 call_error 7ece68f2 render: SoftError: status=failed: GPU out of memory while loading textures['render: SoftError: status=failed: GPU out of memory while loading textures']Scores depend on the model; rank by them rather than relying on a fixed threshold.
Each hit has source_type, source_id, text (what was embedded), score
and created_at. score is 1 minus the cosine distance, clamped to [0, 1].
min_score= drops weaker hits and source_types= limits the search. Only
vectors from the embedder’s model are compared. similar_feedback reuses the
stored vector of that feedback, so it embeds nothing, and returns [] until the
feedback has been embedded. Neither raises; on any failure they log and return
[].
- PostgreSQL ranks in the database (
ORDER BY embedding <=> :q) on an HNSW index. pgvector appliessource_typesand the model filter after the index scan, so a narrow filter over a large table can return fewer thankhits;SET hnsw.iterative_scan = relaxed_order(pgvector 0.8) fixes that. - Other databases compare in Python against the newest
scan_limit(5000) rows, and log once when there are more.
Feedback from your own tools
Section titled “Feedback from your own tools”With instrumentation=mw, submit_feedback records a feedback.submitted
event keyed by the feedback id. If you keep feedback elsewhere, record the same
event with your own id and it is embedded the same way:
import secrets
@app.toolasync def report_problem(title: str, description: str) -> dict: ref = f"bug-{secrets.token_hex(3)}" # your own id and storage mw.record_event( "feedback.submitted", key=ref, attrs={"type": "bug", "title": title, "description": description}, ) return {"id": ref}It stays off the call path
Section titled “It stays off the call path”EmbeddingSink.write only picks out the text and puts it on the sink’s own
bounded queue (max_queue, 1000), so a slow embedding endpoint never holds up
the dispatcher or a tool call. A background task embeds up to batch_size (32)
texts per request, each request limited to timeout (30 seconds), and upserts
the rows. A full queue drops texts and counts them in sink.dropped; embedder
and database failures are logged and counted in sink.errors. A text whose row
already holds the same text (by text_hash) is not embedded again; a changed
text replaces its row. On shutdown, mw.aclose() lets the queue drain for up
to timeout.
Existing data
Section titled “Existing data”The sink embeds records as they arrive. To embed what was stored before it was running, see Backfill embeddings.