Agentic RAG API with FastAPI and LangGraph
An agentic RAG API is a web service whose routes hand a question to an agent that decides whether to search, grades what the search returned, rewrites the question when the results are poor, and then answers from the documents.
Last updated: 09 Oct, 2026 · FastAPI 0.143 · LangGraph 1.2
AgentOps named the project of this module. This lesson opens it up: how the stack starts, which routes the API offers, how the agent behind one of them is wired, and which job fills the search index. Plain RAG runs one fixed line: search, then answer. The agent adds decisions between the steps, and two of its nodes are guardrails.
Starting the stack with Docker Compose
This part of the video starts at 6:03:50. It opens the project's Makefile, reads what its start, stop and restart targets do, and then runs make start: the terminal shows the images being built and the containers starting. These are the targets:
start: ## Start all services
docker compose up --build -d
stop: ## Stop all services
docker compose down
status: ## Show service status
docker compose ps
logs: ## Show service logs
docker compose logs -fShown as it ran in the video, not run here: it needs Docker and the project's .env file, which holds the accounts for Neon, Upstash, Langfuse, Jina and the model provider.
A Makefile gives short names to long commands. make start runs docker compose up --build -d: build the images, start the containers, and return the terminal (-d, detached). The terminal in the video ends with [+] up 7/7: two images built, one network created and four containers started.
| Container | What runs in it | Port | Starts after |
|---|---|---|---|
rag-api | The FastAPI application, built from the repository's Dockerfile | 8000 | OpenSearch reports healthy |
rag-opensearch | OpenSearch 2.19.5, one node: the search index | 9200, 9600 | |
rag-dashboards | OpenSearch Dashboards, a web page for looking into the index | 5601 | OpenSearch |
rag-airflow | Airflow, which runs the daily ingestion job | 8080 | OpenSearch reports healthy |
Each service has a health check in compose.yml, and depends_on with condition: service_healthy makes the API and Airflow wait until OpenSearch answers. The four containers share one bridge network, rag-network, so the API reaches the index at http://opensearch:9200 by service name. After the start, the first thing to read is the API's log, docker logs -f rag-api, which lists every service it connected to.
Reading the API routes
This part of the video starts at 6:08:15. FastAPI serves an interactive page of its routes at localhost:8000/docs. The video opens the health route, runs it, and then scrolls through the search and question routes before sending the first question to the agent. In the project's code the health route checks the database, OpenSearch and the model API; /api/v1/ask runs the project's own prompt builder and model client, only the agent route uses LangGraph, and that route is phase 7 of the repository.
| Route | What it does | Runs the agent | Reads the cache |
|---|---|---|---|
GET /api/v1/health | Checks the database, OpenSearch and the model API | no | no |
POST /api/v1/hybrid-search/ | Returns ranked chunks, no answer | no | no |
POST /api/v1/ask | Search, build a prompt, one model call, return the answer | no | yes |
POST /api/v1/stream | The same, streamed as the model writes | no | yes |
POST /api/v1/ask-agentic | Runs the LangGraph agent | yes | no |
POST /api/v1/feedback | Attaches a score to a trace | no | no |
The health reply on screen in the video lists the database, OpenSearch with Index 'arxiv-papers-chunks' with 317 documents, and the OpenAI API.
The first agent request in the video sends this body: {"categories": ["cs.AI", "cs.LG"], "query": "what is Vector Policy?", "top_k": 3, "use_hybrid": true}. The container log for it reads Generated answer of length: 1830 characters, Sources found: 1, Retrieval attempts: 1 and Execution time: 25.16s.
Exercising the routes with TestClient
The request model
FastAPI reads the request body into a pydantic model. The limits on the fields are the first input check of the service: a question longer than 1000 characters or a top_k above 10 never reaches the search engine.
class AskRequest(BaseModel):
query: str = Field(..., description="User's question", min_length=1, max_length=1000)
top_k: int = Field(3, description="Number of top chunks to retrieve", ge=1, le=10)
use_hybrid: bool = Field(True, description="Use hybrid search (BM25 + vector)")
model: Optional[str] = Field(None, description="Model ID for generation")
categories: Optional[List[str]] = Field(None, description="Filter by arXiv categories")Calling the app in-process
The video calls the running container through Swagger UI and curl. The example below builds a small app with the same request models and four of the routes and calls it with FastAPI's TestClient, which sends requests straight to the app object, so it needs no server and no Docker. The answers are stubs and say so; the routing, the validation and the cache check are real.
import hashlib
import json
import uuid
from typing import List, Optional
from fastapi import FastAPI
from fastapi.testclient import TestClient
from pydantic import BaseModel, Field
class AskRequest(BaseModel): # same fields and limits as the project's AskRequest
query: str = Field(..., min_length=1, max_length=1000)
top_k: int = Field(3, ge=1, le=10)
use_hybrid: bool = True
model: Optional[str] = None
categories: Optional[List[str]] = None
class FeedbackRequest(BaseModel):
trace_id: str
score: float = Field(..., ge=-1, le=1)
comment: Optional[str] = Field(None, max_length=1000)
app = FastAPI(title="arXiv Paper Curator API", version="0.1.0")
CACHE, SCORES, PIPELINE_RUNS = {}, [], []
def cache_key(r: AskRequest) -> str: # the project's key recipe
data = {"query": r.query, "model": r.model, "top_k": r.top_k, "use_hybrid": r.use_hybrid,
"categories": sorted(r.categories) if r.categories else []}
return "exact_cache:" + hashlib.sha256(json.dumps(data, sort_keys=True).encode()).hexdigest()[:16]
@app.get("/api/v1/health")
def health():
return {"status": "ok", "version": app.version,
"services": {"cache": {"status": "healthy", "message": f"{len(CACHE)} keys"}}}
@app.post("/api/v1/ask")
def ask(r: AskRequest):
key = cache_key(r)
if key not in CACHE: # miss: run the pipeline (a stub here) and store
PIPELINE_RUNS.append(r.query)
CACHE[key] = {"query": r.query, "answer": f"(stub answer for: {r.query})", "sources": [],
"chunks_used": r.top_k, "search_mode": "hybrid" if r.use_hybrid else "bm25"}
return CACHE[key]
@app.post("/api/v1/ask-agentic")
def ask_agentic(r: AskRequest): # no cache on this route, as in the project
blocked = "pasta" in r.query.lower() # a stub rule where the project calls the guardrail
return {"query": r.query, "answer": "(stub refusal)" if blocked else "(stub answer)",
"retrieval_attempts": 0 if blocked else 1, "trace_id": uuid.uuid4().hex,
"guardrail_filter": "topic_blocked: off-topic-queries" if blocked
else "Content passed all guardrail checks"}
@app.post("/api/v1/feedback")
def feedback(r: FeedbackRequest):
SCORES.append(r.model_dump())
return {"success": True, "message": "Feedback recorded successfully"}
client = TestClient(app)
print("GET /health ", client.get("/api/v1/health").json())
for attempt in ("first", "second"):
r = client.post("/api/v1/ask", json={"query": "What is Vector Policy Optimization?"})
print(f"POST /ask ({attempt:6}) ", r.status_code, "pipeline runs so far:", len(PIPELINE_RUNS))
a = client.post("/api/v1/ask-agentic", json={"query": "What is the best pasta recipe?"}).json()
print("POST /ask-agentic ", a["guardrail_filter"], "| retrieval_attempts", a["retrieval_attempts"],
"| trace_id length", len(a["trace_id"]))
fb = {"trace_id": a["trace_id"], "score": 0.9, "comment": "This answer was very helpful and accurate!"}
print("POST /feedback ", client.post("/api/v1/feedback", json=fb).json())
bad = client.post("/api/v1/feedback", json={"trace_id": a["trace_id"], "score": 1.5})
print("POST /feedback 1.5 ", bad.status_code, bad.json()["detail"][0]["msg"])
bad = client.post("/api/v1/ask", json={"query": "x", "top_k": 11})
print("POST /ask top_k=11 ", bad.status_code, bad.json()["detail"][0]["msg"])
print("routes ", sorted(client.get("/openapi.json").json()["paths"]))GET /health {'status': 'ok', 'version': '0.1.0', 'services': {'cache': {'status': 'healthy', 'message': '0 keys'}}}
POST /ask (first ) 200 pipeline runs so far: 1
POST /ask (second) 200 pipeline runs so far: 1
POST /ask-agentic topic_blocked: off-topic-queries | retrieval_attempts 0 | trace_id length 32
POST /feedback {'success': True, 'message': 'Feedback recorded successfully'}
POST /feedback 1.5 422 Input should be less than or equal to 1
POST /ask top_k=11 422 Input should be less than or equal to 10
routes ['/api/v1/ask', '/api/v1/ask-agentic', '/api/v1/feedback', '/api/v1/health']What the eight calls returned
- The second
/askdid not run the pipeline. Both calls returned 200, and the pipeline counter stayed at 1: the second answer came from the stored entry. - The agent route answered with a filter name and
retrieval_attempts 0. A blocked question never reaches the search step. The reply carries a 32-character trace id. - The feedback route accepted 0.9 for that trace id and returned the same message the project returns.
- A score of 1.5 and a
top_kof 11 were rejected with status 422 before any handler code ran:Input should be less than or equal to 1andInput should be less than or equal to 10. That is pydantic enforcingle=1andle=10.
The agent graph
This part of the video starts at 6:50:37. It reads the project's graph document: a rewrite node that runs at its own temperature, the retrieve node, and OpenSearch wired in as a tool. The diagram on screen is the project's earlier design, in which a model scored each question from 0 to 100 against a threshold of 60; the code on the branch calls the guardrail service there and adds an output guardrail after the answer.
Behind /api/v1/ask-agentic is a LangGraph graph. A graph is a set of nodes (functions that read and update one shared state) and edges that say which node runs next. A conditional edge picks the next node from the state.
The eight nodes
workflow = StateGraph(AgentState, context_schema=Context)
workflow.add_node("guardrail", ainvoke_guardrail_step)
workflow.add_node("out_of_scope", ainvoke_out_of_scope_step)
workflow.add_node("retrieve", ainvoke_retrieve_step)
workflow.add_node("tool_retrieve", ToolNode(tools))
workflow.add_node("grade_documents", ainvoke_grade_documents_step)
workflow.add_node("rewrite_query", ainvoke_rewrite_query_step)
workflow.add_node("generate_answer", ainvoke_generate_answer_step)
workflow.add_node("output_guardrail", ainvoke_output_guardrail_step)Shown as it ran in the video, not run here: it needs the project's services: OpenSearch, the embedding API, the model and the guardrail.
guardrailsends the question to Amazon Bedrock Guardrails and stores a score: 100 when the question is allowed, 0 when it is blocked.out_of_scopewrites a polite refusal. Nothing is searched.retrieveasks for the search tool by writing a tool call;tool_retrieveis a prebuiltToolNodethat runs it. The search is a tool, not a node of its own.grade_documentsasks the model whether the returned chunks are relevant to the question.rewrite_queryasks the model for a better question, at temperature 0.3.generate_answerwrites the answer from the chunks;output_guardrailchecks that answer before it leaves.
The edges and the routing rule
workflow.add_edge(START, "guardrail")
workflow.add_conditional_edges(
"guardrail",
continue_after_guardrail,
{"continue": "retrieve", "out_of_scope": "out_of_scope"},
)
workflow.add_edge("out_of_scope", END)
workflow.add_edge("tool_retrieve", "grade_documents")
workflow.add_edge("rewrite_query", "retrieve")
workflow.add_edge("generate_answer", "output_guardrail")
workflow.add_edge("output_guardrail", END)The routing function, with its logging lines left out:
def continue_after_guardrail(state: AgentState, runtime: Runtime[Context]) -> Literal["continue", "out_of_scope"]:
score = state["guardrail_result"].score
threshold = runtime.context.guardrail_threshold
return "continue" if score >= threshold else "out_of_scope"Two more conditional edges complete the picture. After retrieve, LangGraph's tools_condition sends the run to tool_retrieve when the last message holds a tool call. After grade_documents, the run goes to generate_answer when the chunks are relevant or two retrieval attempts have been used, and to rewrite_query otherwise. A blocked question ends at out_of_scope; it is never rewritten.
The threshold comes from a Context object passed in at run time (context_schema=Context), not from the state. The running app in the video logs Guardrail threshold: 40. Since the guardrail node only ever writes 0 or 100, any threshold from 1 to 100 gives the same route.
Running the same wiring with stub nodes
The example keeps the project's node names, edges and routing rule and replaces every service with a few lines of Python: a keyword rule for the guardrail, a one-entry dictionary for the index, a fixed rewrite. It prints the path each question takes, at both thresholds.
from dataclasses import dataclass
from typing import Annotated, Literal, Optional, TypedDict
from langchain_core.messages import AIMessage, AnyMessage, HumanMessage
from langchain_core.tools import tool
from langgraph.graph import END, START, StateGraph
from langgraph.graph.message import add_messages
from langgraph.prebuilt import ToolNode, tools_condition
from langgraph.runtime import Runtime
@dataclass
class Context: # per-run settings, read through runtime.context
guardrail_threshold: int = 40
max_retrieval_attempts: int = 2
class State(TypedDict):
messages: Annotated[list[AnyMessage], add_messages]
guardrail_score: Optional[int]
retrieval_attempts: int
routing_decision: Optional[str]
PAPERS = {"vector policy optimization": "VPO trains policies that produce diverse solutions."}
@tool
def retrieve_papers(query: str) -> str:
"""Search the paper index."""
return next((text for key, text in PAPERS.items() if key in query.lower()), "no matching chunk")
def question(state): # the latest human message
return [m for m in state["messages"] if isinstance(m, HumanMessage)][-1].content
def guardrail(state: State, runtime: Runtime[Context]):
allowed = any(word in question(state).lower() for word in ("paper", "policy"))
return {"guardrail_score": 100 if allowed else 0} # allowed is 100, blocked is 0
def continue_after_guardrail(state: State, runtime: Runtime[Context]) -> Literal["continue", "out_of_scope"]:
return "continue" if state["guardrail_score"] >= runtime.context.guardrail_threshold else "out_of_scope"
def out_of_scope(state: State):
return {"messages": [AIMessage("I can only help with questions about research papers.")]}
def retrieve(state: State): # asks for the search tool by writing a tool call
attempt = (state.get("retrieval_attempts") or 0) + 1
call = {"name": "retrieve_papers", "args": {"query": question(state)}, "id": f"call_{attempt}"}
return {"messages": [AIMessage("", tool_calls=[call])], "retrieval_attempts": attempt}
def grade_documents(state: State, runtime: Runtime[Context]):
relevant = state["messages"][-1].content != "no matching chunk"
give_up = state["retrieval_attempts"] >= runtime.context.max_retrieval_attempts
return {"routing_decision": "generate_answer" if relevant or give_up else "rewrite_query"}
def rewrite_query(state: State): # a fixed rewrite where the project asks an LLM
return {"messages": [HumanMessage("What is Vector Policy Optimization?")]}
def generate_answer(state: State):
chunk = [m for m in state["messages"] if m.type == "tool"][-1].content
return {"messages": [AIMessage(f"(answer written from: {chunk})")]}
def output_guardrail(state: State): # the project checks the answer against the chunks here
return {}
g = StateGraph(State, context_schema=Context)
for name, node in [("guardrail", guardrail), ("out_of_scope", out_of_scope), ("retrieve", retrieve),
("tool_retrieve", ToolNode([retrieve_papers])), ("grade_documents", grade_documents),
("rewrite_query", rewrite_query), ("generate_answer", generate_answer),
("output_guardrail", output_guardrail)]:
g.add_node(name, node)
g.add_edge(START, "guardrail")
g.add_conditional_edges("guardrail", continue_after_guardrail,
{"continue": "retrieve", "out_of_scope": "out_of_scope"})
g.add_edge("out_of_scope", END)
g.add_conditional_edges("retrieve", tools_condition, {"tools": "tool_retrieve", END: END})
g.add_edge("tool_retrieve", "grade_documents")
g.add_conditional_edges("grade_documents", lambda s: s["routing_decision"],
{"generate_answer": "generate_answer", "rewrite_query": "rewrite_query"})
g.add_edge("rewrite_query", "retrieve")
g.add_edge("generate_answer", "output_guardrail")
g.add_edge("output_guardrail", END)
graph = g.compile()
for q in ["What is the best pasta recipe?",
"What is Vector Policy Optimization?",
"Find the paper on VPO"]:
for threshold in (40, 60):
steps = graph.stream({"messages": [HumanMessage(q)]}, stream_mode="updates",
context=Context(guardrail_threshold=threshold))
path = " > ".join(name for step in steps for name in step)
print(f"threshold {threshold}: {q}\n {path}")threshold 40: What is the best pasta recipe? guardrail > out_of_scope threshold 60: What is the best pasta recipe? guardrail > out_of_scope threshold 40: What is Vector Policy Optimization? guardrail > retrieve > tool_retrieve > grade_documents > generate_answer > output_guardrail threshold 60: What is Vector Policy Optimization? guardrail > retrieve > tool_retrieve > grade_documents > generate_answer > output_guardrail threshold 40: Find the paper on VPO guardrail > retrieve > tool_retrieve > grade_documents > rewrite_query > retrieve > tool_retrieve > grade_documents > generate_answer > output_guardrail threshold 60: Find the paper on VPO guardrail > retrieve > tool_retrieve > grade_documents > rewrite_query > retrieve > tool_retrieve > grade_documents > generate_answer > output_guardrail
What the three paths show
- The pasta question took two nodes:
guardrail > out_of_scope. No retrieval, no model call. - The in-scope question took six nodes and ended at
output_guardrail, so the answer was checked on the way out. - "Find the paper on VPO" took ten. The first search returned nothing useful,
grade_documentschoserewrite_query, and the second search succeeded. This is the loop in the diagram, run once. - Thresholds 40 and 60 gave identical paths for all three questions, because the score is either 0 or 100.
The ingestion DAG that fills the index
This part of the video starts at 6:31:21. The agent can only answer from papers that are in the index. A scheduled Airflow job puts them there, and the video opens it in the Airflow UI and reads out its five steps.
Airflow runs workflows written as a DAG, a directed acyclic graph: tasks with arrows between them and no way back to an earlier task. The DAG here is named arxiv_paper_ingestion.
dag = DAG(
"arxiv_paper_ingestion",
default_args=default_args, # owner, retries 2, retry_delay 30 minutes
schedule="0 6 * * 1-5", # Monday-Friday at 6 AM UTC
max_active_runs=1,
catchup=False,
)
fetch_task = PythonOperator(task_id="fetch_daily_papers", python_callable=fetch_daily_papers, dag=dag)
# setup_task, index_hybrid_task and report_task are PythonOperator tasks too; cleanup_task is a BashOperator
setup_task >> fetch_task >> index_hybrid_task >> report_task >> cleanup_taskShown as it ran in the video, not run here: it needs Airflow with its metadata database. The project's image installs Airflow 2.10.3, and PythonOperator and BashOperator are imported from airflow.operators.python and airflow.operators.bash, the Airflow 2 paths.
- The cron line
0 6 * * 1-5means 06:00 UTC on Monday to Friday. fetch_daily_paperscalls the arXiv API, downloads and parses the PDFs and writes one row per paper to thepaperstable in Postgres.index_papers_hybridsplits each paper into chunks, asks the embedding API for one vector per chunk and writes text and vector to the OpenSearch indexarxiv-papers-chunks.- Tasks pass small results to each other through XCom, Airflow's store for values a later task needs: the report task reads
fetch_resultsandhybrid_index_stats. - A failed task is retried twice, 30 minutes apart, which suits a source such as arXiv that limits how fast it can be called.
/ask vs /ask-agentic
/api/v1/ask | /api/v1/ask-agentic | |
|---|---|---|
| Flow | Fixed: search, prompt, one model call | A graph that branches and can loop once |
| Model calls for an answered question | One | Two in the video's trace: grading and generation |
| Guardrails | None | One node on the question, one on the answer |
| Off-topic question | Searched and answered from whatever was found | Refused before any search |
| Poor search results | Answered anyway | Question rewritten, searched again |
| Answer cache | Yes | No |
Where you use an agentic RAG API
- Questions that vary in quality. A rewrite step rescues short or vague questions that a fixed pipeline would answer badly.
- A narrow domain with open access. The guardrail node keeps a paper assistant from answering about recipes or elections.
- Several clients for one agent. A web page, a chat bot and another agent can all call the same route, and every call is traced the same way.
Related
- Previous: AgentOps
- Next: Tracing agents with Langfuse
- See also: Amazon Bedrock Guardrails, Redis caching for RAG
- Reference: FastAPI: Testing, LangGraph: Graph API, the LangGraph tutorial
- In the graph example, change
max_retrieval_attemptsinContextto 1 and run it: "Find the paper on VPO" now goes straight fromgrade_documentstogenerate_answer, with no rewrite. - In the API example, send
{"query": ""}to/api/v1/ask: the reply is 422 withString should have at least 1 character, frommin_length=1. - In the API example, post the same question to
/api/v1/askwith"top_k": 5: the pipeline counter goes to 2, because the key changed.
You understood something today that you didn't yesterday.