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_QUERIESretrieval queries in the index's own vocabulary. - grader — one call per unique candidate chunk, run in parallel. Reads one
chunk and replies
KEEPorDROP. 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