Skip to content

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.

source_typeTextsource_id
feedbacktitle and description of a feedback.submitted eventthe feedback id
call_error"{tool}: {error_type}: {error_message}" of calls with outcome error or soft_errorthe call id
eventthe values of error_excerpt, error, message and reason in an event’s attrsthe event id
llm"prompt: ...\ncompletion: ..." of an llm.call event, when capture_llm_text=Truethe 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.

Terminal window
uv add "fastmcp-feedback[instrumentation,embeddings]" asyncpg

The 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).

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 FastMCP
from fastmcp_feedback import add_feedback_tools
from 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.

import asyncio
from fastmcp import Client
@app.tool
def 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 stored
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 crash
0.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 applies source_types and the model filter after the index scan, so a narrow filter over a large table can return fewer than k hits; 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.

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.tool
async 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}

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.

The sink embeds records as they arrive. To embed what was stored before it was running, see Backfill embeddings.