Get 2,500 events tracked for freeSign up now

All tutorials
Syed Raza15 min

Cost per graded chunk in a LangChain RAG agent

A rewriter, an LLM relevance grader over retrieved chunks, an answerer and a grounding check in LangChain 1.x on OpenAI, then a budget on the grader.

Framework
LangChain 1.x
Provider
openai
Python
3.11+
Capsera
pip install capsera

This builds a grounded question-answering agent over a document set in LangChain 1.x against OpenAI: a rewriter turns one question into retrieval queries, a relevance grader reads every candidate chunk and votes it in or out, an answerer composes a cited answer from what survived, and a grounding check reads the answer back against those chunks. One question is four model calls at the floor and fifteen at the cap, and the role that decides where in that range a run lands is the cheapest one. The last third of this page adds Capsera: one init() call to capture every request, four decorator lines to say which role made it, and a blocking budget on the role that multiplies.

What you will build

answer.py, a single-file agent with four roles and a supervisor:

  • rewriter — one call. Reads a cheap index of the knowledge base and the question, writes at most MAX_QUERIES retrieval queries in the index's own vocabulary.
  • grader — one call per unique candidate chunk, run in parallel. Reads one chunk and replies KEEP or DROP. The multiplier.
  • answerer — one call, on the stronger model. Writes the answer from the surviving chunks with a chunk id on every factual sentence.
  • grounding checker — one call. Reads the answer back against those chunks and names the sentences they do not support, which is what buys a revision.

Retrieval is not a model call. A keyword retriever in Python produces the candidate set, and the citation check at the end is a regex. Both are deliberate: a RAG agent that spends model calls on jobs a function can do has a cost curve that nothing in the prompt explains.

Each role is its own function. That matters later — it is the seam the instrumentation attaches to — and it is the better shape regardless, because a pipeline that inlines four model calls has nothing you can name, budget, or move to a different model.

Written against LangChain 1.x with langchain-openai, calling the OpenAI API directly rather than through a gateway. Models: gpt-4.1-mini for the rewriter, the graders and the grounding checker; gpt-4.1 for the answerer.

Prerequisites

Python 3.11 or newer.

pip install "langchain>=1.0,<2" langchain-openai

The major is pinned so this page stays true as LangChain moves. langchain-openai is left unpinned on purpose: install the one that matches your langchain major.

You need an OpenAI API key. ChatOpenAI reads OPENAI_API_KEY from the environment, so export it and no code has to touch it:

export OPENAI_API_KEY=...

Step 1 — Fix the knowledge base, and retrieve from it in code

Start with the chunks, not the agent. The fixture below is seven short chunks standing in for whatever your vector store returns, and it is built around the case that makes RAG hard rather than the case that makes it look good: two chunks state opposite rules and one of them governs, and three more match the question's vocabulary while settling nothing.

"""answer.py — the knowledge base a grounded Q&A agent is allowed to read."""
import re

CHUNKS = {
    "billing-policy-s3-cancellations": (
        "Annual plans are billed in advance and are non-refundable once the "
        "term has started. A customer who cancels mid-term keeps access until "
        "the term ends. Monthly plans are not refunded for a partial month."
    ),
    "billing-policy-s4-price-changes": (
        "When a price increase is applied to an in-flight annual term, the "
        "customer may cancel within sixty days of the notice and receive a "
        "pro-rated refund of the unused months. This overrides the "
        "non-refundable rule in section 3."
    ),
    "terms-of-service-s7": (
        "Notice of a price change is sent to the account billing contact at "
        "least thirty days before it takes effect. The notice date, not the "
        "effective date, starts the sixty-day cancellation window."
    ),
    "release-notes-2026-02": (
        "Billing: refunds issued from the admin console now appear on the "
        "invoice timeline within a minute instead of at the nightly job. No "
        "change to refund eligibility or to how amounts are calculated."
    ),
    "support-macro-refund-decline": (
        "Macro text for declining a refund: Thanks for writing in. Your annual "
        "plan is non-refundable once the term has started, so I am not able to "
        "issue a refund for the remaining months."
    ),
    "runbook-refund-failures": (
        "If a refund fails with provider error 402 the original charge is more "
        "than one hundred and eighty days old and the gateway will not reverse "
        "it. Issue account credit instead and file a finance ticket."
    ),
    "pricing-page-copy": (
        "Switch to annual and pay for ten months instead of twelve. Cancel any "
        "time and you keep access until the end of your term."
    ),
}

QUESTION = (
    "A customer on an annual plan cancelled in month three, five weeks after we "
    "told them the price was going up. Are they owed a refund, and for how much?"
)

_WORD = re.compile(r"[a-z0-9]+")
_STOP = frozenset(
    "the a an and or of to in on for is are was were it its this that with "
    "from at by we our you your not no".split()
)


def _terms(text: str) -> set[str]:
    return {w for w in _WORD.findall(text.lower()) if len(w) > 2 and w not in _STOP}


def chunk_index() -> str:
    return "\n".join(
        f"{cid}: {' '.join(body.split()[:10])}..." for cid, body in CHUNKS.items()
    )


def chunk_text(chunk_id: str) -> str:
    return CHUNKS[chunk_id]


def retrieve(query: str, k: int) -> list[str]:
    wanted = _terms(query)
    if not wanted:
        return []
    scored = [
        (len(wanted & _terms(body)) / len(wanted), cid)
        for cid, body in CHUNKS.items()
    ]
    scored = sorted(
        (pair for pair in scored if pair[0]), key=lambda pair: (-pair[0], pair[1])
    )
    return [cid for _, cid in scored[:k]]


def candidates(queries: list[str], per_query: int, cap: int) -> list[str]:
    pool: list[str] = []
    for query in queries:
        for cid in retrieve(query, per_query):
            if cid not in pool:
                pool.append(cid)
    return pool[:cap]


if __name__ == "__main__":
    print(chunk_index())

retrieve scores on the share of the query's terms a chunk covers, not on the raw overlap count. Raw counts reward length, and the longest chunk in a real knowledge base is usually the release notes. Normalising by query length is one division and it stops the agent paying a grading call to find that out.

It also drops anything that scored zero rather than padding the list to k. That matters for cost in a way it does not matter for quality: k is a ceiling on how many chunks get graded, not a quota to fill, and every chunk the retriever adds to be polite is a model call you pay for in full.

candidates is the function that saves the most money on this page, and it is seven lines. Four retrieval queries at k=3 is twelve retrievals but rarely twelve distinct chunks — related queries pull the same governing policy — so deduplicating before grading rather than after is the difference between twelve grading calls and five. It preserves order, so the chunk the first query ranked highest is still first when the cap bites.

chunk_index() is the cheap artefact that makes the expensive one targeted. The rewriter sees seven one-line stubs, not seven chunks, so the call that decides what to search for costs a fraction of the calls that read the results.

Step 2 — Rewrite the question into queries the index can answer

A support question is written in the customer's words. The knowledge base is written in the company's. Retrieval across that gap is the failure that gets blamed on the model, and one cheap call closes most of it.

# fragment
from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI

FAST_MODEL = "gpt-4.1-mini"
MAX_QUERIES = 4

_BULLET = re.compile(r"^\s*(?:[-*•]|\d+[.)])\s*")


def _lines(text: str, limit: int) -> list[str]:
    cleaned = (_BULLET.sub("", line).strip() for line in text.splitlines())
    return [line for line in cleaned if line][:limit]


REWRITE_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You turn one support question into search queries for a keyword "
            "index. One query per line, no numbering and no commentary. Each "
            "query targets a different fact the answer needs, and uses the "
            "vocabulary of the index rather than the customer's. Never write a "
            "query about something the index does not cover. At most "
            "{max_queries} queries.",
        ),
        ("human", "Question:\n{question}\n\nChunk index:\n{index}"),
    ]
)

rewriter_llm = ChatOpenAI(model=FAST_MODEL, max_tokens=192, temperature=0)
rewrite_chain = REWRITE_PROMPT | rewriter_llm | StrOutputParser()


def rewrite(question: str) -> list[str]:
    raw = rewrite_chain.invoke(
        {
            "question": question,
            "index": chunk_index(),
            "max_queries": MAX_QUERIES,
        }
    )
    return _lines(raw, MAX_QUERIES) or [question]

Lines rather than JSON. A one-query-per-line contract degrades into a shorter list when the model ignores you; a malformed JSON object degrades into an exception, and a retry on a parse failure is a call you pay for twice.

or [question] is the fallback that keeps a bad rewrite from becoming a failed run. If the model returns nothing usable, the original question goes to the retriever unmodified — worse retrieval, still an answer. Agents that treat an empty list as a list tend to produce a confident "I could not find anything", which is the most expensive kind of wrong because nobody re-runs it.

The cap lives in Python. A model asked for "at most four" will occasionally write six, and the slice in _lines is what makes that harmless — not because six queries is untidy, but because each extra query is up to k more chunks and therefore up to k more grading calls.

Step 3 — Grade every candidate chunk, in parallel

This is the role that makes the pipeline a RAG agent rather than a prompt with a search in front of it. The retriever ranks on word overlap, which is why support-macro-refund-decline and pricing-page-copy will be near the top of the candidate list for this question: both are about annual plans and refunds, and neither is policy. One of them is a canned reply that states the rule the exception overrides.

# fragment
DOCS_PER_QUERY = 3
MAX_CANDIDATES = 10
GRADE_CONCURRENCY = 4

GRADE_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You decide whether one knowledge-base chunk carries a fact that "
            "helps answer a question. Reply with the single word KEEP or the "
            "single word DROP and nothing else. DROP marketing copy, canned "
            "reply templates and release notes that only mention the topic: "
            "they share the question's vocabulary and settle nothing. KEEP a "
            "chunk that states a rule, a condition or a date, even if another "
            "chunk contradicts it.",
        ),
        ("human", "Question: {question}\n\nChunk [{chunk_id}]:\n{chunk}"),
    ]
)

grader_llm = ChatOpenAI(model=FAST_MODEL, max_tokens=4, temperature=0)
grade_chain = GRADE_PROMPT | grader_llm | StrOutputParser()


def grade(question: str, chunk_ids: list[str]) -> list[str]:
    verdicts = grade_chain.batch(
        [
            {"question": question, "chunk_id": cid, "chunk": chunk_text(cid)}
            for cid in chunk_ids
        ],
        config={"max_concurrency": GRADE_CONCURRENCY},
    )
    return [
        cid
        for cid, verdict in zip(chunk_ids, verdicts, strict=True)
        if verdict.strip().upper().startswith("KEEP")
    ]

One chunk per call, not the whole candidate list per call. Grading the list in one call is cheaper and it is the wrong trade: the verdicts become positional, a model that skips an item silently shifts every chunk after it, and you cannot tell a careful DROP from a lost line. One chunk per call also means the prompt is a fixed instruction plus one chunk, so raising k multiplies a small number instead of growing a large one.

"KEEP a chunk that states a rule even if another chunk contradicts it" is the instruction that makes the answer correct here. A grader asked only for relevance will DROP the section-3 non-refundable rule once it has seen the section-4 exception, because the exception is the better answer — and then the answerer never learns there was a rule to override, so it states the refund as though it were the default. Resolving a conflict is the answerer's job. The grader's job is to make sure both sides of it arrive.

max_tokens=4 is the whole economic argument for a grader. Output tokens are the expensive half of a call and this role needs one word, so each grading call is close to the floor price of a request. That is what makes it affordable to run ten of them — and it is also why the grader row surprises people later: a cheap call made many times is still the row that grows.

max_concurrency bounds how many are in flight. Ten sequential grading calls is ten round trips of latency for a decision no chunk depends on.

zip(..., strict=True) because pairing chunks to verdicts by position is only safe if the lengths agree, and a silent truncation here would keep a chunk the grader rejected and drop one it wanted.

Step 4 — Answer from what survived

The answerer is the only role on the stronger model, because it is the only one whose output a person reads. The other three produce input for another model or for a Python branch.

# fragment
ANSWER_MODEL = "gpt-4.1"

ANSWER_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You answer a support question from the chunks you are given and "
            "from nothing else. Every sentence that states a fact ends with "
            "the id of the chunk it came from, in square brackets. If two "
            "chunks conflict, say which one governs and why, citing both. If "
            "the chunks do not settle the question, say so in one sentence and "
            "stop. Never compute a figure the chunks do not support. Rewrite "
            "any sentence listed as unsupported, or drop it.",
        ),
        (
            "human",
            "Question:\n{question}\n\nChunks:\n{sources}\n\n"
            "Sentences from your previous draft that were not supported:\n"
            "{flagged}",
        ),
    ]
)

answerer_llm = ChatOpenAI(model=ANSWER_MODEL, max_tokens=700, temperature=0)
answer_chain = ANSWER_PROMPT | answerer_llm | StrOutputParser()


def _sources(chunk_ids: list[str]) -> str:
    return "\n\n".join(f"[{cid}]\n{chunk_text(cid)}" for cid in chunk_ids)


def compose(question: str, chunk_ids: list[str], flagged: list[str]) -> str:
    return answer_chain.invoke(
        {
            "question": question,
            "sources": _sources(chunk_ids),
            "flagged": "\n".join(flagged) or "none",
        }
    )

One prompt for both the first draft and the revision, with flagged reading none on the first pass. A separate revision prompt is a second thing to keep in step with the first, and the only difference between the two passes is a list that is usually empty.

"Never compute a figure the chunks do not support" is aimed at the specific failure this question invites. The chunks say pro-rated refund of the unused months; they do not say what the plan cost. A model that has been asked how much is owed will reach for a number anyway, and nine months of an unstated price is a figure that looks like arithmetic and is actually invention. The grounding check in the next step catches it, but the cheaper fix is not producing it.

Step 5 — Read the answer back against the chunks

Do not ask the answerer to verify itself in the same call. A model that has just written a sentence will defend it, and the check costs the same whether it is honest or not. A separate call on the cheaper model, which sees the chunks and the answer but not the reasoning that connected them, is the whole mechanism here.

# fragment
MAX_FLAGGED = 10

_CITATION = re.compile(r"\[([a-z0-9\-.]+)\]")

GROUNDING_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You check an answer against the chunks it was written from. Reply "
            "with the single word GROUNDED if every factual sentence is "
            "supported by the chunk it cites. Otherwise copy out, one per "
            "line, only the sentences that are not supported. Add nothing "
            "else: no preamble, no corrections, no explanation.",
        ),
        ("human", "Chunks:\n{sources}\n\nAnswer:\n{answer}"),
    ]
)

checker_llm = ChatOpenAI(model=FAST_MODEL, max_tokens=384, temperature=0)
grounding_chain = GROUNDING_PROMPT | checker_llm | StrOutputParser()


def check_grounding(answer: str, chunk_ids: list[str]) -> list[str]:
    raw = grounding_chain.invoke(
        {"sources": _sources(chunk_ids), "answer": answer}
    )
    if raw.strip().upper().startswith("GROUNDED"):
        return []
    return _lines(raw, MAX_FLAGGED)


def miscited(answer: str, chunk_ids: list[str]) -> list[str]:
    allowed = set(chunk_ids)
    return sorted({cid for cid in _CITATION.findall(answer) if cid not in allowed})

The checker is asked to copy sentences, not to describe problems. A description has to be interpreted before it can be used, which means either a third call or a human; a copied sentence goes straight back into the next prompt as the flagged list. It is also the cheaper output — quoting is shorter than explaining.

miscited is free and it catches a different failure. The checker reads whether a claim is supported; the regex reads whether the cited id was even in the set. An answerer that has seen billing-policy-s4-price-changes in a prompt will sometimes cite billing-policy-s5-..., and that citation looks correct to everything except a set membership test. It returns a list rather than raising, because an answer with one bad reference is worth reading with that reference flagged.

Step 6 — The supervisor

One function, and every branch the agent can take is visible in it.

# fragment
from typing import Any

MAX_REVISIONS = 1


def answer_question(question: str) -> dict[str, Any]:
    queries = rewrite(question)
    pool = candidates(queries, DOCS_PER_QUERY, MAX_CANDIDATES)
    if not pool:
        return {"answer": None, "reason": "nothing in the index matched"}

    kept = grade(question, pool)
    if not kept:
        return {"answer": None, "reason": "no chunk survived grading"}

    draft = compose(question, kept, [])
    unsupported = check_grounding(draft, kept)
    revisions = 0
    while unsupported and revisions < MAX_REVISIONS:
        draft = compose(question, kept, unsupported)
        unsupported = check_grounding(draft, kept)
        revisions += 1

    return {
        "answer": draft,
        "reason": None,
        "queries": queries,
        "graded": len(pool),
        "kept": kept,
        "revisions": revisions,
        "unsupported": unsupported,
        "miscited": miscited(draft, kept),
    }

MAX_REVISIONS is the bound that stops the obvious loop. Answer, check, revise, check is a cycle with no natural end — a checker can always find one more sentence to doubt, and every lap costs an answerer call on the expensive model plus a checker call. One revision is a defensible default because the second one rarely changes the verdict: if a sentence is still unsupported after being flagged and rewritten, the chunks do not contain the fact.

The last unsupported is returned rather than swallowed. A run that exhausted its revisions still produces an answer, and the sentences the checker would not pass travel with it. That is the honest shape: the caller decides whether to show it, not the agent that wrote it.

Count the calls: one rewriter, one to ten graders, one or two answerers, one or two checkers. Four at the floor and fifteen at the cap. Both of those numbers come from your constants; nothing in your file decides where between them a given question lands.

The whole file

"""answer.py — a grounded RAG agent: rewrite, grade, answer, check."""
import re
from typing import Any

from langchain_core.output_parsers import StrOutputParser
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI

FAST_MODEL = "gpt-4.1-mini"
ANSWER_MODEL = "gpt-4.1"

MAX_QUERIES = 4
DOCS_PER_QUERY = 3
MAX_CANDIDATES = 10
GRADE_CONCURRENCY = 4
MAX_REVISIONS = 1
MAX_FLAGGED = 10

CHUNKS = {
    "billing-policy-s3-cancellations": (
        "Annual plans are billed in advance and are non-refundable once the "
        "term has started. A customer who cancels mid-term keeps access until "
        "the term ends. Monthly plans are not refunded for a partial month."
    ),
    "billing-policy-s4-price-changes": (
        "When a price increase is applied to an in-flight annual term, the "
        "customer may cancel within sixty days of the notice and receive a "
        "pro-rated refund of the unused months. This overrides the "
        "non-refundable rule in section 3."
    ),
    "terms-of-service-s7": (
        "Notice of a price change is sent to the account billing contact at "
        "least thirty days before it takes effect. The notice date, not the "
        "effective date, starts the sixty-day cancellation window."
    ),
    "release-notes-2026-02": (
        "Billing: refunds issued from the admin console now appear on the "
        "invoice timeline within a minute instead of at the nightly job. No "
        "change to refund eligibility or to how amounts are calculated."
    ),
    "support-macro-refund-decline": (
        "Macro text for declining a refund: Thanks for writing in. Your annual "
        "plan is non-refundable once the term has started, so I am not able to "
        "issue a refund for the remaining months."
    ),
    "runbook-refund-failures": (
        "If a refund fails with provider error 402 the original charge is more "
        "than one hundred and eighty days old and the gateway will not reverse "
        "it. Issue account credit instead and file a finance ticket."
    ),
    "pricing-page-copy": (
        "Switch to annual and pay for ten months instead of twelve. Cancel any "
        "time and you keep access until the end of your term."
    ),
}

QUESTION = (
    "A customer on an annual plan cancelled in month three, five weeks after we "
    "told them the price was going up. Are they owed a refund, and for how much?"
)

_WORD = re.compile(r"[a-z0-9]+")
_BULLET = re.compile(r"^\s*(?:[-*•]|\d+[.)])\s*")
_CITATION = re.compile(r"\[([a-z0-9\-.]+)\]")
_STOP = frozenset(
    "the a an and or of to in on for is are was were it its this that with "
    "from at by we our you your not no".split()
)

REWRITE_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You turn one support question into search queries for a keyword "
            "index. One query per line, no numbering and no commentary. Each "
            "query targets a different fact the answer needs, and uses the "
            "vocabulary of the index rather than the customer's. Never write a "
            "query about something the index does not cover. At most "
            "{max_queries} queries.",
        ),
        ("human", "Question:\n{question}\n\nChunk index:\n{index}"),
    ]
)

GRADE_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You decide whether one knowledge-base chunk carries a fact that "
            "helps answer a question. Reply with the single word KEEP or the "
            "single word DROP and nothing else. DROP marketing copy, canned "
            "reply templates and release notes that only mention the topic: "
            "they share the question's vocabulary and settle nothing. KEEP a "
            "chunk that states a rule, a condition or a date, even if another "
            "chunk contradicts it.",
        ),
        ("human", "Question: {question}\n\nChunk [{chunk_id}]:\n{chunk}"),
    ]
)

ANSWER_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You answer a support question from the chunks you are given and "
            "from nothing else. Every sentence that states a fact ends with "
            "the id of the chunk it came from, in square brackets. If two "
            "chunks conflict, say which one governs and why, citing both. If "
            "the chunks do not settle the question, say so in one sentence and "
            "stop. Never compute a figure the chunks do not support. Rewrite "
            "any sentence listed as unsupported, or drop it.",
        ),
        (
            "human",
            "Question:\n{question}\n\nChunks:\n{sources}\n\n"
            "Sentences from your previous draft that were not supported:\n"
            "{flagged}",
        ),
    ]
)

GROUNDING_PROMPT = ChatPromptTemplate.from_messages(
    [
        (
            "system",
            "You check an answer against the chunks it was written from. Reply "
            "with the single word GROUNDED if every factual sentence is "
            "supported by the chunk it cites. Otherwise copy out, one per "
            "line, only the sentences that are not supported. Add nothing "
            "else: no preamble, no corrections, no explanation.",
        ),
        ("human", "Chunks:\n{sources}\n\nAnswer:\n{answer}"),
    ]
)

rewriter_llm = ChatOpenAI(model=FAST_MODEL, max_tokens=192, temperature=0)
grader_llm = ChatOpenAI(model=FAST_MODEL, max_tokens=4, temperature=0)
answerer_llm = ChatOpenAI(model=ANSWER_MODEL, max_tokens=700, temperature=0)
checker_llm = ChatOpenAI(model=FAST_MODEL, max_tokens=384, temperature=0)

rewrite_chain = REWRITE_PROMPT | rewriter_llm | StrOutputParser()
grade_chain = GRADE_PROMPT | grader_llm | StrOutputParser()
answer_chain = ANSWER_PROMPT | answerer_llm | StrOutputParser()
grounding_chain = GROUNDING_PROMPT | checker_llm | StrOutputParser()


def _terms(text: str) -> set[str]:
    return {w for w in _WORD.findall(text.lower()) if len(w) > 2 and w not in _STOP}


def _lines(text: str, limit: int) -> list[str]:
    cleaned = (_BULLET.sub("", line).strip() for line in text.splitlines())
    return [line for line in cleaned if line][:limit]


def chunk_index() -> str:
    return "\n".join(
        f"{cid}: {' '.join(body.split()[:10])}..." for cid, body in CHUNKS.items()
    )


def chunk_text(chunk_id: str) -> str:
    return CHUNKS[chunk_id]


def retrieve(query: str, k: int) -> list[str]:
    wanted = _terms(query)
    if not wanted:
        return []
    scored = [
        (len(wanted & _terms(body)) / len(wanted), cid)
        for cid, body in CHUNKS.items()
    ]
    scored = sorted(
        (pair for pair in scored if pair[0]), key=lambda pair: (-pair[0], pair[1])
    )
    return [cid for _, cid in scored[:k]]


def candidates(queries: list[str], per_query: int, cap: int) -> list[str]:
    pool: list[str] = []
    for query in queries:
        for cid in retrieve(query, per_query):
            if cid not in pool:
                pool.append(cid)
    return pool[:cap]


def _sources(chunk_ids: list[str]) -> str:
    return "\n\n".join(f"[{cid}]\n{chunk_text(cid)}" for cid in chunk_ids)


def rewrite(question: str) -> list[str]:
    raw = rewrite_chain.invoke(
        {
            "question": question,
            "index": chunk_index(),
            "max_queries": MAX_QUERIES,
        }
    )
    return _lines(raw, MAX_QUERIES) or [question]


def grade(question: str, chunk_ids: list[str]) -> list[str]:
    verdicts = grade_chain.batch(
        [
            {"question": question, "chunk_id": cid, "chunk": chunk_text(cid)}
            for cid in chunk_ids
        ],
        config={"max_concurrency": GRADE_CONCURRENCY},
    )
    return [
        cid
        for cid, verdict in zip(chunk_ids, verdicts, strict=True)
        if verdict.strip().upper().startswith("KEEP")
    ]


def compose(question: str, chunk_ids: list[str], flagged: list[str]) -> str:
    return answer_chain.invoke(
        {
            "question": question,
            "sources": _sources(chunk_ids),
            "flagged": "\n".join(flagged) or "none",
        }
    )


def check_grounding(answer: str, chunk_ids: list[str]) -> list[str]:
    raw = grounding_chain.invoke(
        {"sources": _sources(chunk_ids), "answer": answer}
    )
    if raw.strip().upper().startswith("GROUNDED"):
        return []
    return _lines(raw, MAX_FLAGGED)


def miscited(answer: str, chunk_ids: list[str]) -> list[str]:
    allowed = set(chunk_ids)
    return sorted({cid for cid in _CITATION.findall(answer) if cid not in allowed})


def answer_question(question: str) -> dict[str, Any]:
    queries = rewrite(question)
    pool = candidates(queries, DOCS_PER_QUERY, MAX_CANDIDATES)
    if not pool:
        return {"answer": None, "reason": "nothing in the index matched"}

    kept = grade(question, pool)
    if not kept:
        return {"answer": None, "reason": "no chunk survived grading"}

    draft = compose(question, kept, [])
    unsupported = check_grounding(draft, kept)
    revisions = 0
    while unsupported and revisions < MAX_REVISIONS:
        draft = compose(question, kept, unsupported)
        unsupported = check_grounding(draft, kept)
        revisions += 1

    return {
        "answer": draft,
        "reason": None,
        "queries": queries,
        "graded": len(pool),
        "kept": kept,
        "revisions": revisions,
        "unsupported": unsupported,
        "miscited": miscited(draft, kept),
    }


if __name__ == "__main__":
    result = answer_question(QUESTION)
    if result["answer"] is None:
        print(f"no answer: {result['reason']}")
    else:
        print(result["answer"])
        print(
            f"\n{result['graded']} chunk(s) graded, {len(result['kept'])} kept, "
            f"{result['revisions']} revision(s)"
        )
        if result["unsupported"]:
            print(f"still unsupported: {len(result['unsupported'])} sentence(s)")
        if result["miscited"]:
            print(f"miscited: {', '.join(result['miscited'])}")
python answer.py

You get an answer with a chunk id on each factual sentence, a count of chunks graded and kept, and whatever the checker would not pass. The interesting part is which chunks survive: an answer that cites section 3 and section 4 together has found the conflict and resolved it, and an answer that cites only section 3 has been handed the canned decline macro and agreed with it. Both are real outcomes of an agent shaped like this, and the difference between them is the grader.

The application is finished. Nothing below changes what it does.

Seeing what it cost

What you do not have is any idea which of the four roles spent the money. The provider's invoice has one line for the API key, and the four to fifteen calls behind it are indistinguishable — same key, same account, three of the four roles on the same model.

pip install capsera

At the top of answer.py, above everything else:

# fragment
import os

import capsera

capsera.init(api_key=os.environ["CAPSERA_API_KEY"], endpoint="https://api.capsera.ai")

That is the capture step, and it is already complete. capsera.init() patches the provider client that ChatOpenAI constructs underneath, so from that call on every request the agent makes is recorded: the rewriter's, each grader's inside batch(), the answerer's, the checker's. The prompts are untouched, the chains are untouched, batch() is untouched, and no traffic changed route — the SDK runs inside your process rather than in front of it, so there is no proxy, no base_url to swap and nothing new in the request path.

One addition for a script: a background thread ships events on an interval, so a process that exits immediately can exit before the queue drains.

# fragment
if __name__ == "__main__":
    result = answer_question(QUESTION)
    print(result["answer"])
    capsera.shutdown()

capsera.shutdown() flushes pending events and stops the worker. A long-running service does not need it; a script, a cron job or a notebook cell does, and omitting it is the usual reason a first run reports nothing.

Run it now and you have a total, a per-model split and token counts per call. You also have every one of those events attributed to unknown, because nothing has said which role made it. unknown is a real row, not a dropped one — an unattributed cost is still a cost — but it is one row where you wanted four.

Name the roles

The calls are already captured. The decorator only says who:

# fragment
@capsera.agent("rewriter", team="answers-desk", task_type="query-expansion")
def rewrite(question: str) -> list[str]:
    ...  # body unchanged


@capsera.agent("grader", team="answers-desk", task_type="relevance")
def grade(question: str, chunk_ids: list[str]) -> list[str]:
    ...  # body unchanged


@capsera.agent("answerer", team="answers-desk", task_type="composition")
def compose(question: str, chunk_ids: list[str], flagged: list[str]) -> str:
    ...  # body unchanged


@capsera.agent("grounding-checker", team="answers-desk", task_type="verification")
def check_grounding(answer: str, chunk_ids: list[str]) -> list[str]:
    ...  # body unchanged

Four added lines, and that is the entire diff. No function body changed, no call site changed, no model construction changed, no argument threaded through the supervisor. Every LLM call made while a decorated function is on the stack is attributed to that agent, including calls made by the functions it calls and by the framework it invokes.

The fan-out keeps its identity. grade does its work in grade_chain.batch(...), which dispatches to worker threads rather than calling on the thread you were on, and the grader scope is still what gets recorded. Capsera holds the attribution stack in a ContextVar and LangChain's batch() dispatches through a context-copying executor, so the scope the caller was in is the scope the worker sees — verified empirically on LangChain 1.3 with capsera 0.3.0. batch() also produces one event per underlying call rather than one per batch, so ten candidate chunks are ten grader events (coverage table).

The credit there belongs to LangChain, and it is worth knowing where the guarantee stops. A ContextVar does not cross a thread boundary on its own — it crosses when the code starting the thread copies the context, which batch() does and a hand-rolled pool does not. Grade the candidates yourself with a bare ThreadPoolExecutor and the same decorator records unknown instead of grader, with no error to tell you. Measured on the same versions:

# fragment
import contextvars
from concurrent.futures import ThreadPoolExecutor

# Loses the scope: workers start with a fresh context.
with ThreadPoolExecutor(max_workers=GRADE_CONCURRENCY) as pool:
    verdicts = list(pool.map(grade_one, pool_ids))

# Keeps it: hand each worker a copy of the caller's context.
ctx = contextvars.copy_context()
with ThreadPoolExecutor(max_workers=GRADE_CONCURRENCY) as pool:
    verdicts = list(pool.map(lambda cid: ctx.run(grade_one, cid), pool_ids))

Async needs nothing extra — a ContextVar is snapshotted per task, so roles awaited concurrently under asyncio.gather keep their own attribution.

Leave the supervisor undecorated. answer_question makes no LLM call of its own; every call inside it happens while one of the four roles is on the stack, and the innermost scope is the one recorded. Decorating it would add a row that can never have spend on it. For a per-question total, use the team field — which is what team="answers-desk" on all four decorators is for.

Cost the revision separately

The revision pass is the part of this agent whose cost you cannot see from the role rows alone. A run that passed the grounding check first time and a run that had to rewrite both report as answerer, on the expensive model, and the row says nothing about which. capsera.tag() opens a sub-scope inside a function you have already decorated, which is exactly this case:

# fragment
import capsera


@capsera.agent("answerer", team="answers-desk", task_type="composition")
def compose(question: str, chunk_ids: list[str], flagged: list[str]) -> str:
    if not flagged:
        return _compose(question, chunk_ids, flagged)
    # A revision is the same agent rewriting its own draft against the same
    # chunks. Cost the second pass separately from the first.
    with capsera.tag(
        "answerer", team="answers-desk", task_type="answer-revision"
    ):
        return _compose(question, chunk_ids, flagged)

_compose is the original body, moved down one level and otherwise unchanged. The revision keeps the agent name, so a budget scoped to answerer still counts both passes, and it changes the task type, so the share of answer spend that came from failed grounding checks is a number you can read instead of a suspicion.

Use tag() for a genuine sub-scope like this one. It is the wrong tool for naming an iteration: tagging each grading call with its chunk id would give you one row per chunk, most of which never appear on the next question, and an agent budget cannot accumulate against a name that lived for one answer. The per-invocation axis has its own fields on the same event — customer_id and cost_center on @capsera.agent, session_id on capsera.tag — so you can slice by which customer or which knowledge base a question was for without fragmenting the agent dimension that budgets and cost attribution depend on.

Put a budget on the role that multiplies

The rewriter makes one call per question. The answerer makes one plus MAX_REVISIONS, and the checker matches it. All three of those counts are in your file. The grader's is not: it is the number of queries a model chose to write, multiplied by your k, minus whatever the deduplication in candidates happened to remove. Point this agent at a knowledge base of four thousand chunks instead of seven and the other three rows do not move at all — the answerer still makes one call, because grading is what stands between a bigger corpus and a bigger prompt. The grader is the only row that scales with the thing you are most likely to grow. That is multi-agent amplification in one sentence, and it is why the cap belongs on the grader rather than on the run.

Create the budget in the dashboard with scope agent and the agent set to grader — the same string as the decorator. That string is the whole join between the two halves of this page: a budget can only act on a boundary the events carry, and spend attributed to unknown cannot be governed by a per-agent budget. Pick the amount from a week of your own answering spend rather than a figure from a tutorial, and set the action to block.

Enforcement is on by default (enable_budget_enforcement=True in init()) and the check runs before the provider call. When the grader is over its budget, BudgetExceededError is raised inside your process and the request never leaves it, so no tokens are spent on it. This is pre-call enforcement: the check happens where your code is, not in front of the provider, and a blocked call costs you the check rather than the completion.

Then decide what a blocked grading call means for the answer. batch() raises on the first failure and abandons the remaining inputs unless you tell it otherwise, so this is a decision you make rather than one you inherit:

# fragment
import capsera


def grade(question: str, chunk_ids: list[str]) -> tuple[list[str], int]:
    """Returns the chunks that passed, and how many were never graded."""
    verdicts = grade_chain.batch(
        [
            {"question": question, "chunk_id": cid, "chunk": chunk_text(cid)}
            for cid in chunk_ids
        ],
        config={"max_concurrency": GRADE_CONCURRENCY},
        return_exceptions=True,
    )

    kept: list[str] = []
    ungraded = 0
    for cid, verdict in zip(chunk_ids, verdicts, strict=True):
        if isinstance(verdict, capsera.BudgetExceededError):
            ungraded += 1
            continue
        if isinstance(verdict, BaseException):
            raise RuntimeError(f"grading {cid} failed: {verdict}")
        if verdict.strip().upper().startswith("KEEP"):
            kept.append(cid)
    return kept, ungraded

return_exceptions=True asks batch() to hand failures back as items in the result list instead of raising on the first one, which is what lets you tell a budget block from a flaky call and count the chunks that went unread. Without it, one blocked call ends the grading pass and the remaining chunks are never even attempted.

Counting the blocks rather than ignoring them is the part that matters, and it is where a RAG agent degrades worse than a research agent does. An ungraded chunk is not an irrelevant chunk. Narrower evidence in a research brief makes the brief vaguer; narrower evidence here can make the answer wrong, because the chunk the budget stopped may be the one that governs. On this question that is not hypothetical: drop billing-policy-s4-price-changes and every surviving chunk agrees that annual plans are non-refundable, and the answer is confident, well cited and the opposite of the policy.

So the count travels with the result, and the supervisor reuses the branch it already had:

# fragment
from typing import Any


def answer_question(question: str) -> dict[str, Any]:
    queries = rewrite(question)
    pool = candidates(queries, DOCS_PER_QUERY, MAX_CANDIDATES)
    if not pool:
        return {"answer": None, "reason": "nothing in the index matched"}

    kept, ungraded = grade(question, pool)
    if not kept:
        reason = (
            "the grader budget stopped every chunk"
            if ungraded
            else "no chunk survived grading"
        )
        return {"answer": None, "reason": reason}

    draft = compose(question, kept, [])
    unsupported = check_grounding(draft, kept)
    revisions = 0
    while unsupported and revisions < MAX_REVISIONS:
        draft = compose(question, kept, unsupported)
        unsupported = check_grounding(draft, kept)
        revisions += 1

    return {
        "answer": draft,
        "reason": None,
        # Not a complete review of the retrieved set. Do not present this
        # answer as one.
        "partial": ungraded > 0,
        "ungraded": ungraded,
        "queries": queries,
        "graded": len(pool) - ungraded,
        "kept": kept,
        "revisions": revisions,
        "unsupported": unsupported,
        "miscited": miscited(draft, kept),
    }

There is no second code path for "we ran out of money". A grading pass that returns nothing lands in the branch the agent already had for "no chunk survived grading", with a different reason string, and a pass that returns something answers from it and says so. That is the only reason this degrade is cheap enough to be worth doing — and partial is the field that has to reach whatever renders the answer, because an answer from a partially reviewed set is exactly as confident-looking as one from a complete one.

The single-call roles still need somewhere for a refusal to land. compose and check_grounding call invoke() directly, so there is no return_exceptions to catch a block for them — if you also put a budget on answerer or on the answers-desk team, it surfaces as a raise:

# fragment
import capsera


def run_once(question: str) -> dict[str, Any] | None:
    """Returns None when a budget stopped a role that makes a single call."""
    try:
        return answer_question(question)
    except capsera.BudgetExceededError as exc:
        print(f"budget reached before an answer existed: {exc}")
        return None

BudgetExceededError is the only exception the SDK raises deliberately; everything else fails open. Do not catch it and retry immediately — the budget will still be exceeded. Lowering DOCS_PER_QUERY and rerunning, queueing the question until the period rolls over, and failing are all defensible; for an agent answering from a queue, failing usually is.

Three limits to know before you rely on this as a hard ceiling.

The blast radius is the role, not the run. A budget on grader stops grading calls. The rewriter, the answerer and the checker are different agents with no budget on them, so a question that gets two usable chunks still reaches the answerer on the stronger model and still pays for it — twice, if the grounding check fails. That is the right shape here, because the grader is the row that grows with the knowledge base, but it is worth deciding on purpose rather than discovering: if what you want capped is the answerer's model spend, that is a second budget on answerer, not a side effect of this one.

A parallel grading pass can overshoot. The check reads budget state computed from delivered events, and events are delivered in batches on an interval (flush_interval_ms, 500 ms by default), so a burst of concurrent calls can each pass a check taken before any of them was recorded. Here the burst is GRADE_CONCURRENCY — four as shipped, and whatever you raise it to when the candidate list gets longer. Set the budget below a figure you cannot exceed, or lower the concurrency; it is the concurrency that decides the overshoot, not MAX_CANDIDATES.

It fails open, on purpose. If the pre-call check does not complete within budget_check_timeout (1 second by default), the call is allowed. A network problem between your service and Capsera should not halt production traffic. Of the three actions, block and margin downgrade act at runtime; throttle records the decision but does not delay the call. The full behaviour is in budgets and enforcement.

What this run actually cost

No dollar figures on this page: this agent has not been run against a real key, and a number invented here would not be yours anyway — it depends on your chunks, how many queries your rewriter writes, and today's prices. What is worth predicting is the shape of the result.

Four agent rows — rewriter, grader, answerer, grounding-checker — under one team, answers-desk, with a cost per run for the question as a whole, and the answerer row split by task type into composition and answer-revision.

Read each row as call count and cost per call, not as a total. Three things move on this agent and each has a different fix. The grader's call count moving means the rewriter is issuing more queries or your queries are hitting more distinct chunks, which is a retrieval change rather than a prompt one. The answerer's answer-revision share moving means the grounding check is failing more often, which is the expensive kind of movement because each revision pays for the stronger model and a second checker call. And the checker's cost per call moving means more chunks are surviving grading, because its input is every kept chunk plus the answer.

Which row is largest is not predictable in advance, and that is the point of measuring it rather than reasoning about it. The answerer makes one or two calls on the more expensive model with every kept chunk in the prompt; the grader makes up to ten tiny ones on the cheaper model. Where the crossover falls depends on how long your chunks are and on how many candidates you let through, and it decides whether the next thing to change is the answerer's model or the width of retrieval.

The grader is also the row with the most repetition in it, and repetition is the one place a cheap call can get cheaper. Every grading call in a pass carries the same instruction and differs only in the chunk and the question, so a long, identical prefix is sitting in front of ten requests. The cache read share column is what tells you whether any of it was billed ten times instead of once; it is a number to read rather than one to predict, because a shared prefix has to be long enough and identical enough to be worth anything, and putting the instruction before the chunk is what earns it. Turning on enable_prompt_analysis=True in init() puts a figure on how often one system-prompt hash repeats, which is the evidence for whether that restructuring is worth doing: prompt caching that never engaged.

Each event also carries the file that made the call, which points at answer.py rather than into langchain_core, so a row you did not expect traces back to code rather than to a model name.

Questions this page answers

How do I get per-agent cost for a LangChain RAG agent? Call capsera.init() once, which patches the provider client ChatOpenAI constructs underneath, so every call the agent makes is already recorded. Then put @capsera.agent("<role>") on each role function — the rewriter, the grader, the answerer, the grounding checker. The decorator says who made the call; it does not change whether the call is captured, and no function body, call site or model construction changes.

Which agent in a RAG pipeline should carry the budget? The relevance grader, because it is the only role whose call count is not a constant in your file. The rewriter chooses how many retrieval queries to issue, and the grader runs once per unique chunk those queries return, so grader calls are a model's choice multiplied by your k. The rewriter, the answerer and the grounding checker are one call each per pass.

Why does an LLM relevance grader cost more than the answer it filters for? Because it is priced per chunk, not per run. Each grading call is tiny — one chunk in, one word out — but there is one for every candidate the retrieval step returned, and the answerer makes a single call no matter how many chunks survived. Raising k or letting the rewriter issue more queries moves the grader row and nothing else, so the two roles cross over at a candidate count you have to measure rather than guess.

What happens to a RAG answer when a budget stops part of the grading? An ungraded chunk is not an irrelevant chunk, so it cannot be treated as one. Count the chunks the budget stopped, pass the count out with the answer, and keep the result marked as a partial review — the chunk that governs the question may be among the ones that were never read. This is the one place a RAG agent degrades worse than a research agent: narrower evidence does not make the answer vaguer, it can make it wrong.

Questions this page answers

How do I get per-agent cost for a LangChain RAG agent?
Call capsera.init() once, which patches the provider client ChatOpenAI constructs underneath, so every call the agent makes is already recorded. Then put @capsera.agent("<role>") on each role function — the rewriter, the grader, the answerer, the grounding checker. The decorator says who made the call; it does not change whether the call is captured, and no function body, call site or model construction changes.
Which agent in a RAG pipeline should carry the budget?
The relevance grader, because it is the only role whose call count is not a constant in your file. The rewriter chooses how many retrieval queries to issue, and the grader runs once per unique chunk those queries return, so grader calls are a model's choice multiplied by your k. The rewriter, the answerer and the grounding checker are one call each per pass.
Why does an LLM relevance grader cost more than the answer it filters for?
Because it is priced per chunk, not per run. Each grading call is tiny — one chunk in, one word out — but there is one for every candidate the retrieval step returned, and the answerer makes a single call no matter how many chunks survived. Raising k or letting the rewriter issue more queries moves the grader row and nothing else, so the two roles cross over at a candidate count you have to measure rather than guess.
What happens to a RAG answer when a budget stops part of the grading?
An ungraded chunk is not an irrelevant chunk, so it cannot be treated as one. Count the chunks the budget stopped, pass the count out with the answer, and keep the result marked as a partial review — the chunk that governs the question may be among the ones that were never read. This is the one place a RAG agent degrades worse than a research agent: narrower evidence does not make the answer vaguer, it can make it wrong.

Give every agent an identity, a budget, and hard limits.

One line of code. Anthropic, OpenAI, and Google Gemini.

See pricing