October 1, 2026

RAG with Postgres and pgvector: a FastAPI tutorial with Claude

Build RAG with Postgres, pgvector and FastAPI: chunk, embed with Voyage, HNSW search, cited Claude answers, retrieval eval and cost. Runnable code.

Recruiting

Tech

Guide

Backend

This pgvector RAG tutorial builds a FastAPI service that chunks your documents, embeds them with Voyage, stores the vectors in Postgres, and answers questions with Claude, quoting the passages it used. You'll end up with a /ask endpoint on Postgres that answers with citations to your own docs.

Scoping a hire for this kind of work instead of building it yourself? Start with Hire developers by tech stack: rates, vetting and interview guides. Everything below targets current versions: pgvector 0.8.6, Postgres 18, Voyage 4 embeddings and claude-sonnet-5-5.

What you'll build

A docs Q&A service for an invented product, "Acme Billing". The corpus is six short markdown help pages that you'll create in step 1. Nothing in them is real, which is the point: you can check every answer against text you can read in a minute. Plan on about an hour and a half. The steps build every file, and the Full code section is a copy-paste backstop.

Here's the finished endpoint, called with curl. The response is an example of the shape, not captured output:

Bash
curl -X POST http://localhost:8000/ask \
  -H "content-type: application/json" \
  -H "x-api-key: $API_KEY" \
  -d '{"question": "How long do refunds take?"}'
JSON
{
  "answer": "Approved refunds reach your original payment method in 5 to 10 business days. Bank transfers can take up to 14 days.",
  "citations": [
    {
      "source": "docs/refunds.md",
      "title": "Refunds",
      "cited_text": "Approved refunds reach the original payment method in 5 to 10 business days. Bank transfers can take up to 14 days."
    }
  ]
}

And the project layout you'll end with:

text
rag-pgvector-fastapi/
  .env
  .env.example
  .gitignore
  requirements.txt
  app/
    config.py
    db.py
    chunking.py
    embed.py
    ingest.py
    retrieve.py
    answer.py
    main.py
  eval/
    questions.json
    run_eval.py
  docs/
    refunds.md
    invoices.md
    payment-methods.md
    subscriptions.md
    taxes.md
    usage-billing.md

What you need before you start

  • A macOS or Linux shell. On Windows, use WSL. Every command below assumes a Bash-style shell.
  • Python 3.12 or later. Python 3.10 reaches end of life on 2026-10-31 and 3.12 runs to 2028-10-31, per endoflife.date's Python table. The packages used here all require 3.10 or later, except voyageai, which supports >=3.9,<3.15.
  • Docker, running, to start Postgres 18 with pgvector. PostgreSQL 18 is supported until 2030-11-14, according to endoflife.date's PostgreSQL table.
  • An Anthropic API key. The click path for creating one is in step 1 of How to use the Claude API: build an AI feature step by step.
  • A Voyage AI API key. Voyage's pricing page says the first 200M tokens on an account are free, which is far more than this tutorial uses.
  • Comfort reading Python and SQL. Nothing else.

Quickstart

Already have both keys? This is the whole idea in one script: start Postgres with pgvector, store one embedded sentence, and find it with a question.

Bash
docker pull pgvector/pgvector:pg18-bookworm
docker run --name acme-pg -e POSTGRES_PASSWORD=choose-a-password -p 127.0.0.1:5432:5432 -d pgvector/pgvector:pg18-bookworm

python3.12 -m venv .venv && source .venv/bin/activate
pip install "psycopg[binary]" pgvector voyageai
export VOYAGE_API_KEY="your-voyage-key"
export DATABASE_URL="postgresql://postgres:choose-a-password@localhost:5432/postgres"
Python
import os

import psycopg
import voyageai
from pgvector import Vector
from pgvector.psycopg import register_vector

vo = voyageai.Client()  # reads VOYAGE_API_KEY from the environment

with psycopg.connect(os.environ["DATABASE_URL"], autocommit=True) as conn:
    conn.execute("CREATE EXTENSION IF NOT EXISTS vector")
    register_vector(conn)
    conn.execute("CREATE TABLE IF NOT EXISTS items (id bigserial PRIMARY KEY, body text, embedding vector(1024))")

    text = "Approved refunds reach the original payment method in 5 to 10 business days."
    doc = vo.embed([text], model="voyage-4", input_type="document").embeddings[0]
    conn.execute("INSERT INTO items (body, embedding) VALUES (%s, %s)", (text, Vector(doc)))

    query = vo.embed(["How long do refunds take?"], model="voyage-4", input_type="query").embeddings[0]
    row = conn.execute("SELECT body FROM items ORDER BY embedding <=> %s LIMIT 1", (Vector(query),)).fetchone()
    print(row[0])

Save it as quickstart.py and run python quickstart.py in the same terminal. You should see the refund sentence printed back.

Three things to get right before you build on it:

  • Keep both API keys on the server, in environment variables. Never put them in frontend code or source control.
  • Embed documents and queries with the same model, and make the vector(N) column match its output dimension. Voyage 4 returns 1024 dimensions by default.
  • Query with the operator that matches the index. A cosine index (vector_cosine_ops) pairs with the <=> operator.

If you ran the Quickstart, its acme-pg container is already running. Step 2 tells you what to do about that.

How to build RAG with pgvector and FastAPI

1. Create the project, the environment files and the sample docs

This step makes the folder, the virtual environment, every dependency, the .env files and the six help pages, so nothing later depends on a file you haven't made.

In a terminal, create the folder and the environment:

Bash
mkdir -p rag-pgvector-fastapi/app rag-pgvector-fastapi/eval rag-pgvector-fastapi/docs
cd rag-pgvector-fastapi
python3.12 -m venv .venv
source .venv/bin/activate
printf '.env\n.venv/\n__pycache__/\n' > .gitignore

Create requirements.txt in the project folder:

text
anthropic==1.11.0
voyageai==0.5.0
pgvector==0.5.0
psycopg[binary,pool]==3.3.6
psycopg-pool==3.3.3
fastapi==0.142.2
uvicorn==0.54.0
pydantic-settings==2.15.0

Create .env.example with empty values, then copy it to .env, which .gitignore keeps out of source control:

text
ANTHROPIC_API_KEY=
VOYAGE_API_KEY=
DATABASE_URL=
API_KEY=
Bash
cp .env.example .env

Open .env and fill in the four values. DATABASE_URL is postgresql://postgres:choose-a-password@localhost:5432/postgres, and the password must match the one you give Docker in step 2, so pick your own now. API_KEY is any long random string you make up (python3 -c "import secrets; print(secrets.token_urlsafe(32))" prints one). Keys stay in .env and are loaded by the config module in step 3, never hard-coded.

Install everything the article uses:

Bash
pip install -r requirements.txt

Now create the six help pages. Acme Billing is invented, and so is every number in them. Create docs/refunds.md:

markdown
# Refunds

## Refund eligibility

You can request a refund within 30 days of any charge. Annual plans are refundable within 30 days of the first payment or of a renewal.

Overage charges from usage-based billing can't be refunded once the invoice is paid.

## Refund timing

Approved refunds reach the original payment method in 5 to 10 business days. Bank transfers can take up to 14 days.

You get an email when a refund is approved.

## How to request a refund

Open Billing, then Refunds, pick the charge and select Request refund. Support reviews every request within 2 business days.

Create docs/invoices.md:

markdown
# Invoices

## Finding your invoices

Open Billing, then Invoices. Every invoice downloads as a PDF.

## Changing invoice details

Edit the company name and address under Billing, then Company details. Changes apply to future invoices only.

Issued invoices can't be edited. Ask support for a credit note and a reissued invoice.

## Invoice recipients

Add finance contacts under Billing, then Invoice recipients. Each contact receives a copy of every new invoice by email.

Create docs/payment-methods.md:

markdown
# Payment methods

## Accepted payment methods

We accept Visa, Mastercard and American Express cards. Annual plans can also be paid by bank transfer.

## Failed payments

We retry a failed card payment after 3 days and again after 7 days. If both retries fail, the account becomes read-only until a payment succeeds.

## Updating your card

Open Billing, then Payment methods, and select Replace card. The new card is used for the next charge.

Create docs/subscriptions.md:

markdown
# Subscriptions

## Plans

Acme Billing has three plans: Starter, Team and Business.

## Changing plans

Upgrades take effect immediately, and the price difference is prorated for the rest of the billing period. Downgrades take effect at the next renewal.

To upgrade, open Billing, then Subscription, and choose a new plan.

## Cancelling

Cancel at any time from Billing, then Subscription. Access continues until the end of the paid period, and no further charges are made.

Create docs/taxes.md:

markdown
# Taxes

## VAT and sales tax

Acme Billing adds VAT to invoices for customers in EU countries. VAT isn't added when a business customer has a valid VAT number on file.

## Adding a VAT number

Open Billing, then Company details, and enter your VAT number. Validation takes up to 1 business day, and the number applies to future invoices.

## Tax on refunds

Refunds include the tax charged on the original payment.

Create docs/usage-billing.md:

markdown
# Usage-based billing

## Overage charges

Usage beyond your plan's included volume is billed as overage once a month, in arrears, on the first day of the next month.

## Usage alerts

Set an alert under Billing, then Usage alerts. We email you when usage reaches 80% and 100% of the amount you choose.

Checkpoint. From the project folder, with the virtual environment active:

Bash
ls docs
python -c "import fastapi, voyageai, anthropic, psycopg, pgvector, pydantic_settings; print('imports ok')"

You should see something like:

text
invoices.md  payment-methods.md  refunds.md  subscriptions.md  taxes.md  usage-billing.md
imports ok

If you see ModuleNotFoundError, the virtual environment isn't active: run source .venv/bin/activate and the install again.

2. Start Postgres with pgvector in Docker

Run the pgvector image so there's a Postgres 18 database with the vector extension available. The image comes from the pgvector repository and its Docker Hub tags page. Tags follow pgXX and 0.8.6-pgXX-[distro], cover Postgres 13 to 18, and ship for amd64 and arm64.

In any terminal, use the same password you put in DATABASE_URL:

Bash
docker pull pgvector/pgvector:pg18-bookworm
docker run --name acme-pg -e POSTGRES_PASSWORD=choose-a-password -p 127.0.0.1:5432:5432 -d pgvector/pgvector:pg18-bookworm

The 127.0.0.1: prefix publishes the port only on your machine, so nothing else on your network can reach the database. Skip the docker run if the Quickstart's acme-pg container is still up and its password matches .env.

If you see port is already allocated or address already in use for 5432, a local Postgres owns the port. Remove the failed container with docker rm -f acme-pg, rerun the command with -p 127.0.0.1:5433:5432, and change 5432 to 5433 in DATABASE_URL.

Checkpoint. Ask the database whether pgvector is available:

Bash
docker ps
docker exec acme-pg psql -U postgres -c "SELECT name, default_version FROM pg_available_extensions WHERE name = 'vector';"

You should see something like this (example output, the version may differ):

text
  name  | default_version
--------+-----------------
 vector | 0.8.6
(1 row)

3. Create the config, the schema and the HNSW index

Two files. app/config.py loads your four settings from .env. app/db.py holds the schema (a documents table, a chunks table with a vector(1024) column, and the HNSW index), plus the connection pool.

Create app/config.py:

Python
from functools import lru_cache

from pydantic import Field
from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    model_config = SettingsConfigDict(env_file=".env", extra="ignore")

    anthropic_api_key: str = Field(min_length=1)
    voyage_api_key: str = Field(min_length=1)
    database_url: str = Field(min_length=1)
    api_key: str = Field(min_length=1)


@lru_cache
def get_settings() -> Settings:
    return Settings()

Create app/db.py:

Python
from pgvector.psycopg import register_vector_async
from psycopg import AsyncConnection
from psycopg_pool import AsyncConnectionPool

SCHEMA = """
CREATE EXTENSION IF NOT EXISTS vector;

CREATE TABLE IF NOT EXISTS documents (
    id bigserial PRIMARY KEY,
    source text NOT NULL UNIQUE,
    title text NOT NULL,
    content_hash text NOT NULL
);

CREATE TABLE IF NOT EXISTS chunks (
    id bigserial PRIMARY KEY,
    document_id bigint NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
    chunk_index int NOT NULL,
    content text NOT NULL,
    embedding vector(1024) NOT NULL
);

CREATE INDEX IF NOT EXISTS chunks_embedding_hnsw
    ON chunks USING hnsw (embedding vector_cosine_ops)
    WITH (m = 16, ef_construction = 64);
"""


async def init_schema(database_url: str) -> None:
    async with await AsyncConnection.connect(database_url, autocommit=True) as conn:
        await conn.execute(SCHEMA)


async def _configure(conn: AsyncConnection) -> None:
    await register_vector_async(conn)


def make_pool(database_url: str) -> AsyncConnectionPool:
    return AsyncConnectionPool(
        database_url, min_size=1, max_size=5, open=False, configure=_configure
    )

The content_hash column lets step 6 skip files that haven't changed. The index uses vector_cosine_ops, so searches must use the <=> cosine operator (step 7). HNSW trades a little recall, because it can miss a true nearest neighbour, for speed: without an index, pgvector compares the query against every row. The m and ef_construction values are the pgvector README's defaults, written out so you can see where to change them. Leave them alone until step 8 shows recall is your problem.

pgvector-python's psycopg 3 section gives the registration pattern: register_vector_async(conn), and for pools a configure function passed to the pool constructor. The psycopg pool documentation adds that async pools should be created with open=False and opened with await pool.open(), because opening in the constructor is deprecated. Registration looks up the vector type in the database, so the extension must exist before the pool opens. That's why init_schema uses its own one-off connection.

Checkpoint. Run the schema, then list the indexes on chunks:

Bash
python -c "import asyncio; from app.config import get_settings; from app.db import init_schema; asyncio.run(init_schema(get_settings().database_url))"
docker exec acme-pg psql -U postgres -c "SELECT indexname FROM pg_indexes WHERE tablename = 'chunks';"

You should see something like:

text
       indexname
-----------------------
 chunks_pkey
 chunks_embedding_hnsw
(2 rows)

If you see a ValidationError naming missing fields, .env has an empty value or you're not in the project folder.

4. Chunk the documents

Create app/chunking.py. It splits each markdown file by heading, then packs paragraphs into chunks up to a size limit. A chunk should be small enough to embed one idea and big enough to make sense on its own.

Python
import re
from dataclasses import dataclass


@dataclass
class Chunk:
    index: int
    text: str


def parse_title(markdown: str) -> str:
    match = re.search(r"(?m)^# (.+)$", markdown)
    if match is None:
        raise ValueError("document has no '# ' title line")
    return match.group(1).strip()


def chunk_markdown(markdown: str, max_chars: int = 800) -> list[Chunk]:
    body = re.sub(r"(?m)^# .*\n", "", markdown, count=1)
    chunks: list[Chunk] = []
    for section in re.split(r"(?m)^(?=## )", body):
        section = section.strip()
        if not section:
            continue
        if section.startswith("## "):
            heading_line, _, rest = section.partition("\n")
            heading = heading_line[3:].strip()
        else:
            heading, rest = "", section
        paragraphs = [p.strip() for p in re.split(r"\n\s*\n", rest) if p.strip()]
        pieces: list[str] = []
        current = ""
        for paragraph in paragraphs:
            if current and len(current) + len(paragraph) + 2 > max_chars:
                pieces.append(current)
                current = paragraph
            else:
                current = f"{current}\n\n{paragraph}" if current else paragraph
        if current:
            pieces.append(current)
        for piece in pieces:
            text = f"{heading}\n\n{piece}" if heading else piece
            chunks.append(Chunk(index=len(chunks), text=text))
    return chunks

This is tutorial code, not vendor-documented, and the 800-character limit is a starting point, not a recommendation. Nobody can tell you the best chunk size without your corpus and questions, which is why step 8 measures it. Each chunk keeps its section heading, so a paragraph about "5 to 10 business days" still carries the words "Refund timing" when it's embedded. A single paragraph longer than max_chars stays whole rather than being cut mid-sentence.

Checkpoint. Count the chunks in each doc:

Bash
python -c "from pathlib import Path; from app.chunking import chunk_markdown; print([len(chunk_markdown(p.read_text(encoding='utf-8'))) for p in sorted(Path('docs').glob('*.md'))])"

You should see this, in the order invoices, payment-methods, refunds, subscriptions, taxes, usage-billing:

text
[3, 3, 3, 3, 3, 2]

5. Embed the chunks with Voyage

Create app/embed.py. It sends chunk text to Voyage's voyage-4 model with input_type="document" for stored text and input_type="query" for questions. Anthropic's embeddings guide is direct about the split of work: "Anthropic does not offer its own embedding model," and it recommends Voyage AI.

Python
import asyncio
from functools import lru_cache

import voyageai

from app.config import get_settings

MODEL = "voyage-4"
DIMENSIONS = 1024
BATCH_SIZE = 128


@lru_cache
def get_client() -> voyageai.Client:
    return voyageai.Client(api_key=get_settings().voyage_api_key, max_retries=3)


def _embed_sync(texts: list[str], input_type: str) -> tuple[list[list[float]], int]:
    client = get_client()
    vectors: list[list[float]] = []
    tokens = 0
    for start in range(0, len(texts), BATCH_SIZE):
        result = client.embed(
            texts[start : start + BATCH_SIZE],
            model=MODEL,
            input_type=input_type,
            output_dimension=DIMENSIONS,
        )
        vectors.extend(result.embeddings)
        tokens += result.total_tokens
    if any(len(vector) != DIMENSIONS for vector in vectors):
        raise ValueError(f"expected {DIMENSIONS}-dimension vectors from {MODEL}")
    return vectors, tokens


async def embed_documents(texts: list[str]) -> tuple[list[list[float]], int]:
    return await asyncio.to_thread(_embed_sync, texts, "document")


async def embed_query(text: str) -> tuple[list[float], int]:
    vectors, tokens = await asyncio.to_thread(_embed_sync, [text], "query")
    return vectors[0], tokens

Two details from the same guide. Don't omit input_type: Voyage prepends a different instruction to each ("Represent the document for retrieval:" and "Represent the query for retrieving supporting documents:"), and skipping the parameter drops that. And Voyage embeddings are normalised to length 1, so cosine similarity equals dot product.

voyageai.Client() reads VOYAGE_API_KEY from the shell, but a .env file isn't the shell environment, so the code passes the key in explicitly. The batch size of 128 is a tutorial choice, and output_dimension=1024 keeps the column size and the model output from drifting apart. In the installed 0.5.0 source, max_retries retries rate-limit errors, service-unavailable errors and timeouts with exponential backoff, and nothing else. The client is synchronous, so the async app calls it through asyncio.to_thread. Both functions return the token count, which the cost log in step 10 uses.

Checkpoint. Embed one question:

Bash
python -c "import asyncio; from app.embed import embed_query; v, t = asyncio.run(embed_query('How long do refunds take?')); print(len(v), t)"

You should see something like this (example output, the token count varies):

text
1024 7

If the call fails with an authentication error, recheck VOYAGE_API_KEY in .env. The client doesn't retry a bad key.

6. Ingest the docs into Postgres

Create app/ingest.py, the command-line script that reads the markdown files, chunks and embeds them, and writes them to Postgres. It hashes each file, skips a file whose hash matches the stored content_hash, and otherwise replaces that document's chunks inside one transaction. Re-running costs nothing for unchanged files, and an interrupted run leaves the old chunks in place instead of half of the new ones.

Python
import asyncio
import hashlib
import sys
from pathlib import Path

from pgvector import Vector

from app.chunking import chunk_markdown, parse_title
from app.config import get_settings
from app.db import init_schema, make_pool
from app.embed import embed_documents


async def ingest(docs_dir: str) -> None:
    settings = get_settings()
    await init_schema(settings.database_url)
    pool = make_pool(settings.database_url)
    await pool.open(wait=True)
    try:
        for path in sorted(Path(docs_dir).glob("*.md")):
            markdown = path.read_text(encoding="utf-8")
            digest = hashlib.sha256(markdown.encode("utf-8")).hexdigest()
            source = f"{Path(docs_dir).name}/{path.name}"

            async with pool.connection() as conn:
                cursor = await conn.execute(
                    "SELECT content_hash FROM documents WHERE source = %s", (source,)
                )
                row = await cursor.fetchone()
            if row is not None and row[0] == digest:
                print(f"unchanged  {source}")
                continue

            title = parse_title(markdown)
            chunks = chunk_markdown(markdown)
            vectors, tokens = await embed_documents([chunk.text for chunk in chunks])

            async with pool.connection() as conn:
                cursor = await conn.execute(
                    "INSERT INTO documents (source, title, content_hash) VALUES (%s, %s, %s) "
                    "ON CONFLICT (source) DO UPDATE "
                    "SET title = EXCLUDED.title, content_hash = EXCLUDED.content_hash "
                    "RETURNING id",
                    (source, title, digest),
                )
                doc_id = (await cursor.fetchone())[0]
                await conn.execute("DELETE FROM chunks WHERE document_id = %s", (doc_id,))
                async with conn.cursor() as cur:
                    await cur.executemany(
                        "INSERT INTO chunks (document_id, chunk_index, content, embedding) "
                        "VALUES (%s, %s, %s, %s)",
                        [
                            (doc_id, chunk.index, chunk.text, Vector(vector))
                            for chunk, vector in zip(chunks, vectors)
                        ],
                    )
            print(f"ingested   {source} ({len(chunks)} chunks, {tokens} embedding tokens)")
    finally:
        await pool.close()


if __name__ == "__main__":
    asyncio.run(ingest(sys.argv[1] if len(sys.argv) > 1 else "docs"))

Passing a vector in is a one-liner: wrap the list of floats in pgvector's Vector, as the insert above does.

Checkpoint. Ingest, count the rows, then ingest again:

Bash
python -m app.ingest docs
docker exec acme-pg psql -U postgres -c "SELECT (SELECT count(*) FROM documents) AS documents, (SELECT count(*) FROM chunks) AS chunks;"
python -m app.ingest docs

You should see something like this (example output, the token counts vary):

text
ingested   docs/invoices.md (3 chunks, 120 embedding tokens)
ingested   docs/payment-methods.md (3 chunks, 110 embedding tokens)
ingested   docs/refunds.md (3 chunks, 130 embedding tokens)
ingested   docs/subscriptions.md (3 chunks, 115 embedding tokens)
ingested   docs/taxes.md (3 chunks, 105 embedding tokens)
ingested   docs/usage-billing.md (2 chunks, 70 embedding tokens)

 documents | chunks
-----------+--------
         6 |     17
(1 row)

unchanged  docs/invoices.md
unchanged  docs/payment-methods.md
unchanged  docs/refunds.md
unchanged  docs/subscriptions.md
unchanged  docs/taxes.md
unchanged  docs/usage-billing.md

If you see connection refused, the container isn't running: check docker ps and start it with docker start acme-pg.

Create app/retrieve.py. It orders chunks by cosine distance to the question's vector and takes the top k. Smaller distance means closer, and 0 means identical direction. The <=> operator matches the vector_cosine_ops index from step 3.

Python
from dataclasses import dataclass

from pgvector import Vector
from psycopg_pool import AsyncConnectionPool


@dataclass
class Hit:
    source: str
    title: str
    content: str
    distance: float


async def search(
    pool: AsyncConnectionPool,
    query_vector: list[float],
    k: int = 4,
    ef_search: int = 40,
    source: str | None = None,
) -> list[Hit]:
    vector = Vector(query_vector)
    sql = (
        "SELECT d.source, d.title, c.content, c.embedding <=> %s AS distance "
        "FROM chunks c JOIN documents d ON d.id = c.document_id "
    )
    params: list = [vector]
    if source is not None:
        sql += "WHERE d.source = %s "
        params.append(source)
    sql += "ORDER BY c.embedding <=> %s LIMIT %s"
    params += [vector, k]

    async with pool.connection() as conn:
        async with conn.transaction():
            await conn.execute("SELECT set_config('hnsw.ef_search', %s, true)", (str(ef_search),))
            if source is not None:
                await conn.execute("SELECT set_config('hnsw.iterative_scan', 'strict_order', true)")
            cursor = await conn.execute(sql, params)
            rows = await cursor.fetchall()
    return [Hit(source=r[0], title=r[1], content=r[2], distance=r[3]) for r in rows]

This function is tutorial glue. The ORDER BY ... LIMIT shape is what lets Postgres use the index. It sets hnsw.ef_search with set_config(..., true), which lasts only for the current transaction, so one request's tuning can't leak into the next request that reuses the connection.

The optional source filter exists for one reason. An approximate index finds its nearest candidates first and applies your WHERE clause afterwards, so a restrictive filter can leave you with fewer than k rows. The README's fix is SET hnsw.iterative_scan = strict_order;, which the function turns on whenever a filter is present. It needs pgvector 0.8.0 or later. On an older pgvector, a filtered /search fails with invalid configuration parameter name and returns a 503. On a corpus this small, Postgres may ignore the index and scan the table, which is fine. Run EXPLAIN on the query to see which it picked.

Checkpoint. In the project folder, search for the refund question (a heredoc, so no new file is created):

Bash
python - <<'EOF'
import asyncio

from app.config import get_settings
from app.db import make_pool
from app.embed import embed_query
from app.retrieve import search


async def main():
    pool = make_pool(get_settings().database_url)
    await pool.open(wait=True)
    vector, _ = await embed_query("How long do refunds take?")
    for hit in await search(pool, vector, k=3):
        print(f"{hit.distance:.3f}  {hit.source}  {hit.content.splitlines()[0]}")
    await pool.close()


asyncio.run(main())
EOF

You should see something like this (example output, your distances differ). The refund-timing chunk should rank first:

text
0.312  docs/refunds.md  Refund timing
0.455  docs/refunds.md  Refund eligibility
0.601  docs/refunds.md  How to request a refund

8. Evaluate retrieval

Measure retrieval before you trust it. You'll write 12 labelled questions, then compute hit@k (did the right document appear in the top k?) and MRR (mean reciprocal rank: 1 for first place, 0.5 for second, 0.33 for third, 0 for a miss, averaged). Both files are glue, written for this tutorial.

Create eval/questions.json:

JSON
[
  {"question": "How long do refunds take?", "source": "docs/refunds.md"},
  {"question": "Can I get my money back on an annual plan?", "source": "docs/refunds.md"},
  {"question": "Where do I download an invoice as a PDF?", "source": "docs/invoices.md"},
  {"question": "How do I change the company name on my invoices?", "source": "docs/invoices.md"},
  {"question": "Which cards and payment methods do you accept?", "source": "docs/payment-methods.md"},
  {"question": "What happens when my card payment fails?", "source": "docs/payment-methods.md"},
  {"question": "How do I upgrade from Starter to Team?", "source": "docs/subscriptions.md"},
  {"question": "Can I cancel in the middle of a billing period?", "source": "docs/subscriptions.md"},
  {"question": "Is VAT added to my invoice?", "source": "docs/taxes.md"},
  {"question": "How do I add my VAT number?", "source": "docs/taxes.md"},
  {"question": "When are overage charges billed?", "source": "docs/usage-billing.md"},
  {"question": "Can I get an alert before my usage gets expensive?", "source": "docs/usage-billing.md"}
]

Create eval/run_eval.py. It embeds each question the same way /ask will, calls search, and finds the rank of the first chunk from the labelled source:

Python
import argparse
import asyncio
import json
from pathlib import Path

from app.config import get_settings
from app.db import make_pool
from app.embed import embed_query
from app.retrieve import search


async def run(k: int, ef_search: int) -> None:
    questions = json.loads(Path("eval/questions.json").read_text(encoding="utf-8"))
    pool = make_pool(get_settings().database_url)
    await pool.open(wait=True)
    hits_at_k = 0
    reciprocal_ranks: list[float] = []
    missed: list[dict] = []
    try:
        for item in questions:
            vector, _ = await embed_query(item["question"])
            hits = await search(pool, vector, k=k, ef_search=ef_search)
            rank = next(
                (i for i, hit in enumerate(hits, start=1) if hit.source == item["source"]), None
            )
            if rank is None:
                missed.append(item)
                reciprocal_ranks.append(0.0)
            else:
                hits_at_k += 1
                reciprocal_ranks.append(1 / rank)
    finally:
        await pool.close()

    total = len(questions)
    print(f"questions={total} k={k} ef_search={ef_search}")
    print(f"hit@{k}={hits_at_k}/{total}")
    print(f"MRR={sum(reciprocal_ranks) / total:.3f}")
    for item in missed:
        print(f"missed: {item['question']} (expected {item['source']})")


if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--k", type=int, default=4)
    parser.add_argument("--ef-search", type=int, default=40)
    args = parser.parse_args()
    asyncio.run(run(args.k, args.ef_search))

Checkpoint. Run it, change one thing, run it again:

Bash
python -m eval.run_eval --k 4 --ef-search 40
python -m eval.run_eval --k 2 --ef-search 40

The output has this shape, with your own numbers where the placeholders are:

text
questions=12 k=4 ef_search=40
hit@4=<hits>/12
MRR=<number between 0 and 1>
missed: <question> (expected docs/<file>.md)

Two honest limits. The six-document corpus is tiny, so changing ef_search probably won't move the numbers at all. On a real corpus with thousands of chunks it will, and this is the script you'll use to find out. And no benchmark number belongs in this article: the figures that matter are the ones you measure on your own documents and questions.

The list of missed questions is the useful output. Read each miss and ask whether the chunk was cut in a bad place, the wording of the question was far from the document, or k was too small. That's also what a good retrieval review looks like in an interview: How to hire AI engineers lists reviewing a broken RAG pipeline as one of its exercises.

9. Answer with Claude and read the citations

Create app/answer.py. It turns each hit into a search_result block, sends the blocks plus the question to Claude, and reads the citations off the reply. This is how you get native citations instead of asking the model to write [1] markers and parsing them back out. It also holds the per-query cost estimator used in step 10.

Anthropic's search results documentation defines the block: type: "search_result", a source string, a title string, a content list of text blocks, and citations: {"enabled": true}. In the Python SDK it's the SearchResultBlockParam type.

Python
from dataclasses import dataclass

import anthropic
from anthropic.types import SearchResultBlockParam

from app.retrieve import Hit

MODEL = "claude-sonnet-5-5"
MAX_TOKENS = 2048
SYSTEM_PROMPT = (
    "You answer questions about Acme Billing using only the search results provided. "
    "Answer in one to three sentences. If the search results do not contain the answer, "
    "say that the docs don't cover it. Never use outside knowledge."
)

SONNET_INPUT_USD_PER_MTOK = 2.0
SONNET_OUTPUT_USD_PER_MTOK = 10.0
VOYAGE_4_USD_PER_MTOK = 0.06


class AnswerError(Exception):
    pass


@dataclass
class Answer:
    text: str
    citations: list[dict[str, str]]
    input_tokens: int
    output_tokens: int


def estimate_cost_usd(embed_tokens: int, input_tokens: int, output_tokens: int) -> float:
    return (
        embed_tokens * VOYAGE_4_USD_PER_MTOK
        + input_tokens * SONNET_INPUT_USD_PER_MTOK
        + output_tokens * SONNET_OUTPUT_USD_PER_MTOK
    ) / 1_000_000


def to_search_results(hits: list[Hit]) -> list[SearchResultBlockParam]:
    return [
        {
            "type": "search_result",
            "source": hit.source,
            "title": hit.title,
            "content": [{"type": "text", "text": paragraph} for paragraph in hit.content.split("\n\n")],
            "citations": {"enabled": True},
        }
        for hit in hits
    ]


async def answer_question(client: anthropic.AsyncAnthropic, question: str, hits: list[Hit]) -> Answer:
    message = await client.messages.create(
        model=MODEL,
        max_tokens=MAX_TOKENS,
        system=SYSTEM_PROMPT,
        messages=[
            {
                "role": "user",
                "content": [*to_search_results(hits), {"type": "text", "text": question}],
            }
        ],
    )
    if message.stop_reason != "end_turn":
        raise AnswerError(f"stop_reason={message.stop_reason}, request {message._request_id}")

    parts: list[str] = []
    citations: dict[tuple[str, str], dict[str, str]] = {}
    for block in message.content:
        if block.type != "text":
            continue
        parts.append(block.text)
        for citation in block.citations or []:
            if citation.type == "search_result_location":
                citations.setdefault(
                    (citation.source, citation.cited_text),
                    {"source": citation.source, "title": citation.title or "", "cited_text": citation.cited_text},
                )
    return Answer("".join(parts), list(citations.values()), message.usage.input_tokens, message.usage.output_tokens)

The block is the smallest thing Claude can cite, so a one-block chunk can only be cited whole. Splitting each chunk into paragraph blocks, as to_search_results does, gives quote-sized citations. Citations aren't a beta feature: no header is needed, and the same documentation says all active models support search results except Claude Haiku 3.

Each response citation of type search_result_location carries source, title, cited_text, search_result_index (0-based, in the order you sent the blocks), and start_block_index and end_block_index (end exclusive). The cited_text isn't counted toward your output tokens, per the same documentation.

The model is claude-sonnet-5-5. Anthropic's models overview tells you to start with Opus 5.5 for most workloads. Sonnet 5.5 costs half as much, so this tutorial uses it and you can test Opus on your own questions. The system prompt makes the service refuse politely: when the results don't contain the answer you get an answer with an empty citations list, which is a useful signal in itself.

Read the response defensively. Adaptive thinking is on by default on Sonnet 5.5, so the code loops over content blocks and checks type instead of reading content[0]. It checks stop_reason too: anything other than end_turn (a cut-off at max_tokens, a refusal) is an error here, not an answer. MAX_TOKENS = 2048 is deliberately generous, because thinking tokens count toward it. Don't combine this call with structured outputs: Anthropic's citations documentation lists the two features as incompatible.

Checkpoint. Retrieve and answer in one script (a heredoc again):

Bash
python - <<'EOF'
import asyncio

import anthropic

from app.answer import answer_question
from app.config import get_settings
from app.db import make_pool
from app.embed import embed_query
from app.retrieve import search


async def main():
    settings = get_settings()
    pool = make_pool(settings.database_url)
    await pool.open(wait=True)
    client = anthropic.AsyncAnthropic(api_key=settings.anthropic_api_key)
    question = "How long do refunds take?"
    vector, _ = await embed_query(question)
    answer = await answer_question(client, question, await search(pool, vector))
    print(answer.text)
    print(answer.citations)
    await client.close()
    await pool.close()


asyncio.run(main())
EOF

You should see something like this (example output, the wording varies):

text
Approved refunds reach your original payment method in 5 to 10 business days. Bank transfers can take up to 14 days.
[{'source': 'docs/refunds.md', 'title': 'Refunds', 'cited_text': 'Approved refunds reach the original payment method in 5 to 10 business days. Bank transfers can take up to 14 days.'}]

If the call fails with a 401 authentication error, recheck ANTHROPIC_API_KEY in .env.

10. Serve it with FastAPI and verify

Create app/main.py. It exposes /ask and a debug /search route, manages the pool and the Claude client in a lifespan handler, and puts the basics in front: an API-key check, input caps and typed error handling. FastAPI's lifespan events guide shows the pattern: an @asynccontextmanager function that sets things up, yields, and cleans up, passed as FastAPI(lifespan=lifespan). The older on_event handlers are deprecated, and you can't use both styles at once. The routes and handler are web-framework glue, written for this tutorial.

Python
import hmac
import logging
from contextlib import asynccontextmanager, contextmanager

import anthropic
import psycopg
from fastapi import Depends, FastAPI, Header, HTTPException, Request
from pydantic import BaseModel, Field
from voyageai.error import VoyageError

from app.answer import AnswerError, answer_question, estimate_cost_usd
from app.config import get_settings
from app.db import init_schema, make_pool
from app.embed import embed_query
from app.retrieve import search

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("app")


@asynccontextmanager
async def lifespan(app: FastAPI):
    settings = get_settings()
    await init_schema(settings.database_url)
    pool = make_pool(settings.database_url)
    await pool.open(wait=True)
    app.state.pool = pool
    app.state.claude = anthropic.AsyncAnthropic(api_key=settings.anthropic_api_key)
    try:
        yield
    finally:
        await app.state.claude.close()
        await pool.close()


app = FastAPI(lifespan=lifespan)


def require_api_key(x_api_key: str | None = Header(default=None)) -> None:
    expected = get_settings().api_key.encode("utf-8")
    supplied = (x_api_key or "").encode("utf-8")
    if not hmac.compare_digest(supplied, expected):
        raise HTTPException(status_code=401, detail="Invalid API key.")


class Question(BaseModel):
    question: str = Field(min_length=1, max_length=500)
    k: int = Field(default=4, ge=1, le=10)
    source: str | None = Field(default=None, max_length=200)


class Citation(BaseModel):
    source: str
    title: str
    cited_text: str


class AskResponse(BaseModel):
    answer: str
    citations: list[Citation]


class HitOut(BaseModel):
    source: str
    title: str
    content: str
    distance: float


@contextmanager
def translate_errors():
    try:
        yield
    except AnswerError as exc:
        logger.warning("unusable answer: %s", exc)
        raise HTTPException(status_code=502, detail="The model did not return a usable answer.") from exc
    except (anthropic.APIError, VoyageError) as exc:
        logger.exception("upstream failure")
        raise HTTPException(status_code=502, detail="An upstream service failed.") from exc
    except psycopg.Error as exc:
        logger.exception("database failure")
        raise HTTPException(status_code=503, detail="The database is unavailable.") from exc


@app.post("/search", response_model=list[HitOut], dependencies=[Depends(require_api_key)])
async def search_route(body: Question, request: Request):
    with translate_errors():
        query_vector, _ = await embed_query(body.question)
        hits = await search(request.app.state.pool, query_vector, k=body.k, source=body.source)
    return [HitOut(source=h.source, title=h.title, content=h.content, distance=h.distance) for h in hits]


@app.post("/ask", response_model=AskResponse, dependencies=[Depends(require_api_key)])
async def ask(body: Question, request: Request):
    with translate_errors():
        query_vector, embed_tokens = await embed_query(body.question)
        hits = await search(request.app.state.pool, query_vector, k=body.k, source=body.source)
        answer = await answer_question(request.app.state.claude, body.question, hits)
    logger.info(
        "ask: embed_tokens=%d input_tokens=%d output_tokens=%d est_usd=%.6f citations=%d",
        embed_tokens,
        answer.input_tokens,
        answer.output_tokens,
        estimate_cost_usd(embed_tokens, answer.input_tokens, answer.output_tokens),
        len(answer.citations),
    )
    return AskResponse(answer=answer.text, citations=[Citation(**c) for c in answer.citations])

The decisions inside it:

  • Auth: /ask and /search need an x-api-key header equal to API_KEY, compared with hmac.compare_digest. It's a shared secret for server-to-server callers, so don't ship it in a browser bundle. If a team shares keys like this one, AI in the engineering workflow: tools, permissions and policy covers who should hold them.
  • Caps: the question is limited to 500 characters and k to 10, so a leaked key can't push huge prompts through your Anthropic account.
  • Errors: Voyage and Anthropic failures become a 502, database failures a 503, an unusable model reply a 502 with a generic message. Details go to the log, never to the client, and keys are never logged.
  • CORS: absent on purpose, because the service is called server to server. If you later need a browser caller, list exact origins, never ["*"].
  • Retries: the Anthropic SDK retries transient failures by default, and the Voyage client was set up for it in step 5.

Checkpoint. In one terminal, from the project folder with the virtual environment active, start the server (if port 8000 is taken, add --port 8001 and use that port in the curl commands):

Bash
uvicorn app.main:app --reload

You should see something like:

text
INFO:     Uvicorn running on http://127.0.0.1:8000 (Press CTRL+C to quit)
INFO:     Application startup complete.

In a second terminal, from the same folder, export the same API_KEY value that's in .env and ask the question from What you'll build:

Bash
export API_KEY="the-value-from-your-.env"
curl -X POST http://localhost:8000/ask \
  -H "content-type: application/json" \
  -H "x-api-key: $API_KEY" \
  -d '{"question": "How long do refunds take?"}'

You should see something like this (example output, the wording varies):

JSON
{"answer":"Approved refunds reach your original payment method in 5 to 10 business days. Bank transfers can take up to 14 days.","citations":[{"source":"docs/refunds.md","title":"Refunds","cited_text":"Approved refunds reach the original payment method in 5 to 10 business days. Bank transfers can take up to 14 days."}]}

The server terminal also logs the cost of the request as an ask: embed_tokens=... est_usd=... line. Now the negative checks. A question the docs can't answer should come back with an empty citations list:

Bash
curl -X POST http://localhost:8000/ask \
  -H "content-type: application/json" \
  -H "x-api-key: $API_KEY" \
  -d '{"question": "Do you offer a student discount?"}'
JSON
{"answer":"The docs don't cover student discounts.","citations":[]}

And a call without the key should be rejected:

Bash
curl -i -X POST http://localhost:8000/ask \
  -H "content-type: application/json" \
  -d '{"question": "How long do refunds take?"}'
text
HTTP/1.1 401 Unauthorized
{"detail":"Invalid API key."}

If a call you expected to work returns 401, the API_KEY exported in this terminal doesn't match .env.

Architecture

text
client
  |  POST /ask  (x-api-key header, question)
  v
FastAPI app (app/main.py): auth, input caps, lifespan-managed pool
  |  1. embed the question (Voyage, input_type="query")
  v
Voyage embeddings API
  |  vector(1024)
  v
Postgres + pgvector: ORDER BY embedding <=> query LIMIT k   (HNSW index)
  |  top-k chunks
  v
FastAPI app: chunks become search_result blocks (citations enabled)
  |  Messages API call (claude-sonnet-5-5)
  v
Claude API
  |  text blocks carrying search_result_location citations
  v
client gets {answer, citations}

The ingest script (app/ingest.py) is a separate path into the same database: read markdown, chunk, embed with input_type="document", write. Both API keys live in .env on the server.

Clean up

Stop the server with Ctrl+C, then remove the database container (this also deletes the Quickstart's items table and every chunk):

Bash
docker rm -f acme-pg
deactivate

Delete the project folder if you're done with it, and revoke the Anthropic and Voyage keys if you created them only for this tutorial.

Full code

Everything below runs as pasted, and matches the files the steps built. The files are tutorial code except where the steps cite vendor documentation. Keys load from .env through pydantic-settings and are passed to each SDK explicitly.

The project layout:

text
rag-pgvector-fastapi/
  .env
  .env.example
  .gitignore
  requirements.txt
  app/
    config.py
    db.py
    chunking.py
    embed.py
    ingest.py
    retrieve.py
    answer.py
    main.py
  eval/
    questions.json
    run_eval.py
  docs/
    refunds.md
    invoices.md
    payment-methods.md
    subscriptions.md
    taxes.md
    usage-billing.md

Install and run, from the project folder. Start the database first (step 2), then:

Bash
python3.12 -m venv .venv && source .venv/bin/activate
pip install -r requirements.txt
cp .env.example .env   # then fill in the four values in .env
python -m app.ingest docs
uvicorn app.main:app --reload

In a second terminal, with API_KEY exported to the same value you put in .env:

Bash
curl -X POST http://localhost:8000/ask \
  -H "content-type: application/json" \
  -H "x-api-key: $API_KEY" \
  -d '{"question": "How long do refunds take?"}'

curl -X POST http://localhost:8000/ask \
  -H "content-type: application/json" \
  -H "x-api-key: $API_KEY" \
  -d '{"question": "Do you offer a student discount?"}'

The second question isn't in the docs, so the answer should say so and citations should come back empty. Then run python -m eval.run_eval --k 4.

.gitignore:

text
.env
.venv/
__pycache__/

.env.example:

text
ANTHROPIC_API_KEY=
VOYAGE_API_KEY=
DATABASE_URL=
API_KEY=

For the local database from step 2, DATABASE_URL is postgresql://postgres:choose-a-password@localhost:5432/postgres, with your own password. API_KEY is any long random string you make up. .env is a copy of .env.example with your four values filled in.

requirements.txt:

text
anthropic==1.11.0
voyageai==0.5.0
pgvector==0.5.0
psycopg[binary,pool]==3.3.6
psycopg-pool==3.3.3
fastapi==0.142.2
uvicorn==0.54.0
pydantic-settings==2.15.0

app/config.py:

Python
from functools import lru_cache

from pydantic import Field
from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    model_config = SettingsConfigDict(env_file=".env", extra="ignore")

    anthropic_api_key: str = Field(min_length=1)
    voyage_api_key: str = Field(min_length=1)
    database_url: str = Field(min_length=1)
    api_key: str = Field(min_length=1)


@lru_cache
def get_settings() -> Settings:
    return Settings()

app/db.py:

Python
from pgvector.psycopg import register_vector_async
from psycopg import AsyncConnection
from psycopg_pool import AsyncConnectionPool

SCHEMA = """
CREATE EXTENSION IF NOT EXISTS vector;

CREATE TABLE IF NOT EXISTS documents (
    id bigserial PRIMARY KEY,
    source text NOT NULL UNIQUE,
    title text NOT NULL,
    content_hash text NOT NULL
);

CREATE TABLE IF NOT EXISTS chunks (
    id bigserial PRIMARY KEY,
    document_id bigint NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
    chunk_index int NOT NULL,
    content text NOT NULL,
    embedding vector(1024) NOT NULL
);

CREATE INDEX IF NOT EXISTS chunks_embedding_hnsw
    ON chunks USING hnsw (embedding vector_cosine_ops)
    WITH (m = 16, ef_construction = 64);
"""


async def init_schema(database_url: str) -> None:
    async with await AsyncConnection.connect(database_url, autocommit=True) as conn:
        await conn.execute(SCHEMA)


async def _configure(conn: AsyncConnection) -> None:
    await register_vector_async(conn)


def make_pool(database_url: str) -> AsyncConnectionPool:
    return AsyncConnectionPool(
        database_url, min_size=1, max_size=5, open=False, configure=_configure
    )

app/chunking.py:

Python
import re
from dataclasses import dataclass


@dataclass
class Chunk:
    index: int
    text: str


def parse_title(markdown: str) -> str:
    match = re.search(r"(?m)^# (.+)$", markdown)
    if match is None:
        raise ValueError("document has no '# ' title line")
    return match.group(1).strip()


def chunk_markdown(markdown: str, max_chars: int = 800) -> list[Chunk]:
    body = re.sub(r"(?m)^# .*\n", "", markdown, count=1)
    chunks: list[Chunk] = []
    for section in re.split(r"(?m)^(?=## )", body):
        section = section.strip()
        if not section:
            continue
        if section.startswith("## "):
            heading_line, _, rest = section.partition("\n")
            heading = heading_line[3:].strip()
        else:
            heading, rest = "", section
        paragraphs = [p.strip() for p in re.split(r"\n\s*\n", rest) if p.strip()]
        pieces: list[str] = []
        current = ""
        for paragraph in paragraphs:
            if current and len(current) + len(paragraph) + 2 > max_chars:
                pieces.append(current)
                current = paragraph
            else:
                current = f"{current}\n\n{paragraph}" if current else paragraph
        if current:
            pieces.append(current)
        for piece in pieces:
            text = f"{heading}\n\n{piece}" if heading else piece
            chunks.append(Chunk(index=len(chunks), text=text))
    return chunks

app/embed.py:

Python
import asyncio
from functools import lru_cache

import voyageai

from app.config import get_settings

MODEL = "voyage-4"
DIMENSIONS = 1024
BATCH_SIZE = 128


@lru_cache
def get_client() -> voyageai.Client:
    return voyageai.Client(api_key=get_settings().voyage_api_key, max_retries=3)


def _embed_sync(texts: list[str], input_type: str) -> tuple[list[list[float]], int]:
    client = get_client()
    vectors: list[list[float]] = []
    tokens = 0
    for start in range(0, len(texts), BATCH_SIZE):
        result = client.embed(
            texts[start : start + BATCH_SIZE],
            model=MODEL,
            input_type=input_type,
            output_dimension=DIMENSIONS,
        )
        vectors.extend(result.embeddings)
        tokens += result.total_tokens
    if any(len(vector) != DIMENSIONS for vector in vectors):
        raise ValueError(f"expected {DIMENSIONS}-dimension vectors from {MODEL}")
    return vectors, tokens


async def embed_documents(texts: list[str]) -> tuple[list[list[float]], int]:
    return await asyncio.to_thread(_embed_sync, texts, "document")


async def embed_query(text: str) -> tuple[list[float], int]:
    vectors, tokens = await asyncio.to_thread(_embed_sync, [text], "query")
    return vectors[0], tokens

app/ingest.py:

Python
import asyncio
import hashlib
import sys
from pathlib import Path

from pgvector import Vector

from app.chunking import chunk_markdown, parse_title
from app.config import get_settings
from app.db import init_schema, make_pool
from app.embed import embed_documents


async def ingest(docs_dir: str) -> None:
    settings = get_settings()
    await init_schema(settings.database_url)
    pool = make_pool(settings.database_url)
    await pool.open(wait=True)
    try:
        for path in sorted(Path(docs_dir).glob("*.md")):
            markdown = path.read_text(encoding="utf-8")
            digest = hashlib.sha256(markdown.encode("utf-8")).hexdigest()
            source = f"{Path(docs_dir).name}/{path.name}"

            async with pool.connection() as conn:
                cursor = await conn.execute(
                    "SELECT content_hash FROM documents WHERE source = %s", (source,)
                )
                row = await cursor.fetchone()
            if row is not None and row[0] == digest:
                print(f"unchanged  {source}")
                continue

            title = parse_title(markdown)
            chunks = chunk_markdown(markdown)
            vectors, tokens = await embed_documents([chunk.text for chunk in chunks])

            async with pool.connection() as conn:
                cursor = await conn.execute(
                    "INSERT INTO documents (source, title, content_hash) VALUES (%s, %s, %s) "
                    "ON CONFLICT (source) DO UPDATE "
                    "SET title = EXCLUDED.title, content_hash = EXCLUDED.content_hash "
                    "RETURNING id",
                    (source, title, digest),
                )
                doc_id = (await cursor.fetchone())[0]
                await conn.execute("DELETE FROM chunks WHERE document_id = %s", (doc_id,))
                async with conn.cursor() as cur:
                    await cur.executemany(
                        "INSERT INTO chunks (document_id, chunk_index, content, embedding) "
                        "VALUES (%s, %s, %s, %s)",
                        [
                            (doc_id, chunk.index, chunk.text, Vector(vector))
                            for chunk, vector in zip(chunks, vectors)
                        ],
                    )
            print(f"ingested   {source} ({len(chunks)} chunks, {tokens} embedding tokens)")
    finally:
        await pool.close()


if __name__ == "__main__":
    asyncio.run(ingest(sys.argv[1] if len(sys.argv) > 1 else "docs"))

app/retrieve.py:

Python
from dataclasses import dataclass

from pgvector import Vector
from psycopg_pool import AsyncConnectionPool


@dataclass
class Hit:
    source: str
    title: str
    content: str
    distance: float


async def search(
    pool: AsyncConnectionPool,
    query_vector: list[float],
    k: int = 4,
    ef_search: int = 40,
    source: str | None = None,
) -> list[Hit]:
    vector = Vector(query_vector)
    sql = (
        "SELECT d.source, d.title, c.content, c.embedding <=> %s AS distance "
        "FROM chunks c JOIN documents d ON d.id = c.document_id "
    )
    params: list = [vector]
    if source is not None:
        sql += "WHERE d.source = %s "
        params.append(source)
    sql += "ORDER BY c.embedding <=> %s LIMIT %s"
    params += [vector, k]

    async with pool.connection() as conn:
        async with conn.transaction():
            await conn.execute("SELECT set_config('hnsw.ef_search', %s, true)", (str(ef_search),))
            if source is not None:
                await conn.execute("SELECT set_config('hnsw.iterative_scan', 'strict_order', true)")
            cursor = await conn.execute(sql, params)
            rows = await cursor.fetchall()
    return [Hit(source=r[0], title=r[1], content=r[2], distance=r[3]) for r in rows]

app/answer.py:

Python
from dataclasses import dataclass

import anthropic
from anthropic.types import SearchResultBlockParam

from app.retrieve import Hit

MODEL = "claude-sonnet-5-5"
MAX_TOKENS = 2048
SYSTEM_PROMPT = (
    "You answer questions about Acme Billing using only the search results provided. "
    "Answer in one to three sentences. If the search results do not contain the answer, "
    "say that the docs don't cover it. Never use outside knowledge."
)

SONNET_INPUT_USD_PER_MTOK = 2.0
SONNET_OUTPUT_USD_PER_MTOK = 10.0
VOYAGE_4_USD_PER_MTOK = 0.06


class AnswerError(Exception):
    pass


@dataclass
class Answer:
    text: str
    citations: list[dict[str, str]]
    input_tokens: int
    output_tokens: int


def estimate_cost_usd(embed_tokens: int, input_tokens: int, output_tokens: int) -> float:
    return (
        embed_tokens * VOYAGE_4_USD_PER_MTOK
        + input_tokens * SONNET_INPUT_USD_PER_MTOK
        + output_tokens * SONNET_OUTPUT_USD_PER_MTOK
    ) / 1_000_000


def to_search_results(hits: list[Hit]) -> list[SearchResultBlockParam]:
    return [
        {
            "type": "search_result",
            "source": hit.source,
            "title": hit.title,
            "content": [{"type": "text", "text": paragraph} for paragraph in hit.content.split("\n\n")],
            "citations": {"enabled": True},
        }
        for hit in hits
    ]


async def answer_question(client: anthropic.AsyncAnthropic, question: str, hits: list[Hit]) -> Answer:
    message = await client.messages.create(
        model=MODEL,
        max_tokens=MAX_TOKENS,
        system=SYSTEM_PROMPT,
        messages=[
            {
                "role": "user",
                "content": [*to_search_results(hits), {"type": "text", "text": question}],
            }
        ],
    )
    if message.stop_reason != "end_turn":
        raise AnswerError(f"stop_reason={message.stop_reason}, request {message._request_id}")

    parts: list[str] = []
    citations: dict[tuple[str, str], dict[str, str]] = {}
    for block in message.content:
        if block.type != "text":
            continue
        parts.append(block.text)
        for citation in block.citations or []:
            if citation.type == "search_result_location":
                citations.setdefault(
                    (citation.source, citation.cited_text),
                    {"source": citation.source, "title": citation.title or "", "cited_text": citation.cited_text},
                )
    return Answer("".join(parts), list(citations.values()), message.usage.input_tokens, message.usage.output_tokens)

app/main.py:

Python
import hmac
import logging
from contextlib import asynccontextmanager, contextmanager

import anthropic
import psycopg
from fastapi import Depends, FastAPI, Header, HTTPException, Request
from pydantic import BaseModel, Field
from voyageai.error import VoyageError

from app.answer import AnswerError, answer_question, estimate_cost_usd
from app.config import get_settings
from app.db import init_schema, make_pool
from app.embed import embed_query
from app.retrieve import search

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("app")


@asynccontextmanager
async def lifespan(app: FastAPI):
    settings = get_settings()
    await init_schema(settings.database_url)
    pool = make_pool(settings.database_url)
    await pool.open(wait=True)
    app.state.pool = pool
    app.state.claude = anthropic.AsyncAnthropic(api_key=settings.anthropic_api_key)
    try:
        yield
    finally:
        await app.state.claude.close()
        await pool.close()


app = FastAPI(lifespan=lifespan)


def require_api_key(x_api_key: str | None = Header(default=None)) -> None:
    expected = get_settings().api_key.encode("utf-8")
    supplied = (x_api_key or "").encode("utf-8")
    if not hmac.compare_digest(supplied, expected):
        raise HTTPException(status_code=401, detail="Invalid API key.")


class Question(BaseModel):
    question: str = Field(min_length=1, max_length=500)
    k: int = Field(default=4, ge=1, le=10)
    source: str | None = Field(default=None, max_length=200)


class Citation(BaseModel):
    source: str
    title: str
    cited_text: str


class AskResponse(BaseModel):
    answer: str
    citations: list[Citation]


class HitOut(BaseModel):
    source: str
    title: str
    content: str
    distance: float


@contextmanager
def translate_errors():
    try:
        yield
    except AnswerError as exc:
        logger.warning("unusable answer: %s", exc)
        raise HTTPException(status_code=502, detail="The model did not return a usable answer.") from exc
    except (anthropic.APIError, VoyageError) as exc:
        logger.exception("upstream failure")
        raise HTTPException(status_code=502, detail="An upstream service failed.") from exc
    except psycopg.Error as exc:
        logger.exception("database failure")
        raise HTTPException(status_code=503, detail="The database is unavailable.") from exc


@app.post("/search", response_model=list[HitOut], dependencies=[Depends(require_api_key)])
async def search_route(body: Question, request: Request):
    with translate_errors():
        query_vector, _ = await embed_query(body.question)
        hits = await search(request.app.state.pool, query_vector, k=body.k, source=body.source)
    return [HitOut(source=h.source, title=h.title, content=h.content, distance=h.distance) for h in hits]


@app.post("/ask", response_model=AskResponse, dependencies=[Depends(require_api_key)])
async def ask(body: Question, request: Request):
    with translate_errors():
        query_vector, embed_tokens = await embed_query(body.question)
        hits = await search(request.app.state.pool, query_vector, k=body.k, source=body.source)
        answer = await answer_question(request.app.state.claude, body.question, hits)
    logger.info(
        "ask: embed_tokens=%d input_tokens=%d output_tokens=%d est_usd=%.6f citations=%d",
        embed_tokens,
        answer.input_tokens,
        answer.output_tokens,
        estimate_cost_usd(embed_tokens, answer.input_tokens, answer.output_tokens),
        len(answer.citations),
    )
    return AskResponse(answer=answer.text, citations=[Citation(**c) for c in answer.citations])

eval/run_eval.py:

Python
import argparse
import asyncio
import json
from pathlib import Path

from app.config import get_settings
from app.db import make_pool
from app.embed import embed_query
from app.retrieve import search


async def run(k: int, ef_search: int) -> None:
    questions = json.loads(Path("eval/questions.json").read_text(encoding="utf-8"))
    pool = make_pool(get_settings().database_url)
    await pool.open(wait=True)
    hits_at_k = 0
    reciprocal_ranks: list[float] = []
    missed: list[dict] = []
    try:
        for item in questions:
            vector, _ = await embed_query(item["question"])
            hits = await search(pool, vector, k=k, ef_search=ef_search)
            rank = next(
                (i for i, hit in enumerate(hits, start=1) if hit.source == item["source"]), None
            )
            if rank is None:
                missed.append(item)
                reciprocal_ranks.append(0.0)
            else:
                hits_at_k += 1
                reciprocal_ranks.append(1 / rank)
    finally:
        await pool.close()

    total = len(questions)
    print(f"questions={total} k={k} ef_search={ef_search}")
    print(f"hit@{k}={hits_at_k}/{total}")
    print(f"MRR={sum(reciprocal_ranks) / total:.3f}")
    for item in missed:
        print(f"missed: {item['question']} (expected {item['source']})")


if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--k", type=int, default=4)
    parser.add_argument("--ef-search", type=int, default=40)
    args = parser.parse_args()
    asyncio.run(run(args.k, args.ef_search))

eval/questions.json:

JSON
[
  {"question": "How long do refunds take?", "source": "docs/refunds.md"},
  {"question": "Can I get my money back on an annual plan?", "source": "docs/refunds.md"},
  {"question": "Where do I download an invoice as a PDF?", "source": "docs/invoices.md"},
  {"question": "How do I change the company name on my invoices?", "source": "docs/invoices.md"},
  {"question": "Which cards and payment methods do you accept?", "source": "docs/payment-methods.md"},
  {"question": "What happens when my card payment fails?", "source": "docs/payment-methods.md"},
  {"question": "How do I upgrade from Starter to Team?", "source": "docs/subscriptions.md"},
  {"question": "Can I cancel in the middle of a billing period?", "source": "docs/subscriptions.md"},
  {"question": "Is VAT added to my invoice?", "source": "docs/taxes.md"},
  {"question": "How do I add my VAT number?", "source": "docs/taxes.md"},
  {"question": "When are overage charges billed?", "source": "docs/usage-billing.md"},
  {"question": "Can I get an alert before my usage gets expensive?", "source": "docs/usage-billing.md"}
]

The six help pages follow. Acme Billing is invented, and so is every number in them.

docs/refunds.md:

markdown
# Refunds

## Refund eligibility

You can request a refund within 30 days of any charge. Annual plans are refundable within 30 days of the first payment or of a renewal.

Overage charges from usage-based billing can't be refunded once the invoice is paid.

## Refund timing

Approved refunds reach the original payment method in 5 to 10 business days. Bank transfers can take up to 14 days.

You get an email when a refund is approved.

## How to request a refund

Open Billing, then Refunds, pick the charge and select Request refund. Support reviews every request within 2 business days.

docs/invoices.md:

markdown
# Invoices

## Finding your invoices

Open Billing, then Invoices. Every invoice downloads as a PDF.

## Changing invoice details

Edit the company name and address under Billing, then Company details. Changes apply to future invoices only.

Issued invoices can't be edited. Ask support for a credit note and a reissued invoice.

## Invoice recipients

Add finance contacts under Billing, then Invoice recipients. Each contact receives a copy of every new invoice by email.

docs/payment-methods.md:

markdown
# Payment methods

## Accepted payment methods

We accept Visa, Mastercard and American Express cards. Annual plans can also be paid by bank transfer.

## Failed payments

We retry a failed card payment after 3 days and again after 7 days. If both retries fail, the account becomes read-only until a payment succeeds.

## Updating your card

Open Billing, then Payment methods, and select Replace card. The new card is used for the next charge.

docs/subscriptions.md:

markdown
# Subscriptions

## Plans

Acme Billing has three plans: Starter, Team and Business.

## Changing plans

Upgrades take effect immediately, and the price difference is prorated for the rest of the billing period. Downgrades take effect at the next renewal.

To upgrade, open Billing, then Subscription, and choose a new plan.

## Cancelling

Cancel at any time from Billing, then Subscription. Access continues until the end of the paid period, and no further charges are made.

docs/taxes.md:

markdown
# Taxes

## VAT and sales tax

Acme Billing adds VAT to invoices for customers in EU countries. VAT isn't added when a business customer has a valid VAT number on file.

## Adding a VAT number

Open Billing, then Company details, and enter your VAT number. Validation takes up to 1 business day, and the number applies to future invoices.

## Tax on refunds

Refunds include the tax charged on the original payment.

docs/usage-billing.md:

markdown
# Usage-based billing

## Overage charges

Usage beyond your plan's included volume is billed as overage once a month, in arrears, on the first day of the next month.

## Usage alerts

Set an alert under Billing, then Usage alerts. We email you when usage reaches 80% and 100% of the amount you choose.

Reference

pgvector distance operators

The four operators, from the pgvector README. You used <=> to match vector_cosine_ops.

OperatorDistance
<->L2 (Euclidean)
<=>Cosine distance
<#>Negative inner product
<+>L1

HNSW parameters and limits

ParameterWhere it appliesDefault
mindex build16
ef_constructionindex build64
hnsw.ef_searcheach query40

The build parameters are the defaults, and step 3 wrote them out. The query-time knob is the cheap one to turn: SET hnsw.ef_search = 100;. The README caps an indexed vector column at 2,000 dimensions, halfvec at 4,000 and binary vectors at 64,000. That rules out Voyage's 2048-dimension option for a plain vector column, and it's why this tutorial stays at 1024.

Voyage embedding limits

Voyage's embeddings reference lists voyage-4-large, voyage-4, voyage-4-lite and voyage-4-nano, a 32,000-token context, output dimensions of 1024 (the default), 256, 512 and 2048, and a limit of 1,000 texts per request. The code batches 128 at a time, well under that limit.

What a query costs

A query has two bills: embedding the question (Voyage) and generating the answer (Claude). Both report their own token counts: Voyage returns total_tokens on each embedding result, and the Claude response carries a usage object with input_tokens and output_tokens. estimate_cost_usd in app/answer.py turns them into dollars, and /ask logs the figure on every request.

Voyage's pricing page lists per-million-token prices of $0.12 for voyage-4-large, $0.06 for voyage-4 and $0.02 for voyage-4-lite, with the first 200M tokens free per account. The Claude prices (checked 2026-10-01) are in the models overview:

ModelModel IDInput per MTokOutput per MTok
Sonnet 5.5claude-sonnet-5-5$2$10
Opus 5.5claude-opus-5-5$4$20
Haiku 4.5claude-haiku-4-5-20251001$1$5

A worked example with assumed sizes, not measured ones: a 20-token question, four retrieved chunks that make 1,500 input tokens with the system prompt, and a 200-token answer. Embedding is 20 x $0.06 / 1M, about $0.000001. Input is 1,500 x $2 / 1M = $0.003. Output is 200 x $10 / 1M = $0.002. Total: roughly $0.005 per question, before any thinking tokens, which bill as output.

The embedding side is almost free by comparison. Retrieval cost is dominated by how many chunks you pass to Claude, so k is also a cost control. Indexing the six sample docs is a few thousand tokens at $0.06 per million, and it falls inside the free 200M tokens anyway.

Common mistakes

  • Mismatching dimensions. vector(1024) only works with a model and output_dimension that produce 1024 numbers. A wider output (Voyage offers 2048) also breaks the index, because a vector column can't be indexed beyond 2,000 dimensions.
  • Omitting input_type. Anthropic's guide says not to. Embed stored text as document and questions as query, with the same model on both sides.
  • Pairing the wrong operator and operator class. An index built with vector_cosine_ops serves cosine distance, so query with <=>.
  • Registering the vector type on one connection instead of every pool connection. Put register_vector_async in the pool's configure callback, and create the extension before the pool opens.
  • Filtering without iterative scan. A WHERE clause on an approximate index can return fewer than k rows. Set hnsw.iterative_scan for filtered queries.
  • Treating an answer with no citations as grounded. An empty citations list means the model cited nothing, so show it as unverified or refuse it.

Next steps

  • Add hybrid search. Postgres full-text search catches exact terms (error codes, product names) that embeddings blur, and you can merge the two result lists before they reach Claude.
  • Add reranking. Voyage's rerank-2.5 (and rerank-2.5-lite) reorders the top candidates by relevance before you build the search results. Re-run step 8 to see whether it earns its latency.
  • Try halfvec. The README lets you index half-precision vectors up to 4,000 dimensions, which opens up Voyage's 2048-dimension output.
  • Stream the answer. The Claude API tutorial linked at the top shows the stream helper and a server-sent-events route, and the /ask handler is a drop-in place for it.

Hiring Python developers through HighCircl

If you'd rather add an engineer than build this alone, you can hire Python developers through HighCircl. Python is one of the stacks it covers, across seven European countries. Senior engineers run €45-105/hr ($50-115/hr), with a transparent 20% capped margin on top of what the engineer earns and no minimum hours. Matching takes 72 hours. Every engineer goes through four stages of vetting run by senior engineers, and about 1 in 10 applicants pass.

FAQ

Why use pgvector instead of a dedicated vector database?

If you already run Postgres, pgvector keeps your chunks, their metadata and your filters in one database you already back up, secure and understand. This tutorial's JOIN and WHERE d.source = ... are ordinary SQL. A dedicated vector store is worth evaluating when one Postgres instance stops being enough, and your own measurements (not a general claim) should decide that.

HNSW or IVFFlat?

This tutorial uses HNSW, which the pgvector README documents with its build parameters m and ef_construction and the query-time setting hnsw.ef_search. IVFFlat is the README's other index type. Compare them on your own corpus with the step 8 script, and pick the one that gives the recall you need at the latency you can accept.

Does Anthropic offer an embedding model?

No. Anthropic's embeddings guide says it doesn't offer its own embedding model and recommends Voyage AI. This tutorial uses voyage-4, with voyage-4-large and voyage-4-lite as the alternatives on the same pricing page.

How many dimensions should the vector column have?

It must equal what your embedding model returns. Voyage 4 defaults to 1024 and also offers 256, 512 and 2048. A vector column can be indexed up to 2,000 dimensions and halfvec up to 4,000, so 2048 only works with halfvec. Smaller dimensions mean smaller storage, and you can measure what they cost in recall with the eval script.

How do I get citations from Claude?

Send your retrieved chunks as search_result blocks with citations: {"enabled": true} and read the search_result_location citations on the response's text blocks. Each one gives the source, the title, the cited_text and the block indexes. Split chunks into smaller content blocks if you want finer-grained citations, because the block is the smallest unit Claude can cite.

How much does a RAG query cost?

Two meters run: the Voyage embedding of the question, and Claude's input and output tokens. On voyage-4 ($0.06 per million tokens) and claude-sonnet-5-5 ($2 input, $10 output per million), an assumed 1,500-token prompt with a 200-token answer comes to about $0.005. The input side grows with the number of chunks you send, so k is the main cost lever. Log usage and compute your own figure.

Share this article

Author Image

HighCircl Editorial Team

The HighCircl editorial team writes about hiring software engineers, nearshore development, and engineering team building. Our articles draw on direct experience sourcing and placing senior developers across Poland, Hungary, Slovakia, Serbia, Slovenia, Romania, and Spain — and on candid conversations with the CTOs and engineering leads who hire them.

HighCircl is a nearshore engineering network that delivers matched candidate shortlists in 72 hours. Every piece of content we publish is informed by real engagement data: actual developer rates, real hiring timelines, and what separates engineering teams that scale cleanly from those that stall.

Take Me to the Experts

Access our network of industry-leading software engineers.

Start Now