AI SecurityNeMo Guardrails 0.24 · RAGAS 0.4 · OpenAI SDK 3.3 · Python 3.12 or 3.13
Dashboard
0%
1
Curious builder0 XP earned · 300 to level 2
0 daysFinish a lesson to begin
Badge collection0 of 6 unlocked
52 small wins to finish your pathNext lesson →

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

The Makefile and make start · from the Complete AI Security Course in 8 Hours video · 6:03:50 to 6:04:56

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:

makefile
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 -f

Shown 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.

ContainerWhat runs in itPortStarts after
rag-apiThe FastAPI application, built from the repository's Dockerfile8000OpenSearch reports healthy
rag-opensearchOpenSearch 2.19.5, one node: the search index9200, 9600
rag-dashboardsOpenSearch Dashboards, a web page for looking into the index5601OpenSearch
rag-airflowAirflow, which runs the daily ingestion job8080OpenSearch 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

The API routes in Swagger UI · from the Complete AI Security Course in 8 Hours video · 6:08:15 to 6:09:29

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.

RouteWhat it doesRuns the agentReads the cache
GET /api/v1/healthChecks the database, OpenSearch and the model APInono
POST /api/v1/hybrid-search/Returns ranked chunks, no answernono
POST /api/v1/askSearch, build a prompt, one model call, return the answernoyes
POST /api/v1/streamThe same, streamed as the model writesnoyes
POST /api/v1/ask-agenticRuns the LangGraph agentyesno
POST /api/v1/feedbackAttaches a score to a tracenono

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.

python
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.

ExampleRun on FastAPI 0.143
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"]))

What the eight calls returned

  • The second /ask did 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_k of 11 were rejected with status 422 before any handler code ran: Input should be less than or equal to 1 and Input should be less than or equal to 10. That is pydantic enforcing le=1 and le=10.

The agent graph

The rewrite node, the retrieve node and the search tool · from the Complete AI Security Course in 8 Hours video · 6:50:37 to 6:51:18

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 agent graph runs from START to the guardrail node, where a blocked question goes to out_of_scope and ends while an allowed question goes to retrieve, tool_retrieve and grade_documents, after which rewrite_query sends a new question back to retrieve when the chunks are not relevant and attempts are left, and otherwise generate_answer runs, then output_guardrail, then END.

The eight nodes

python
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.

  • guardrail sends the question to Amazon Bedrock Guardrails and stores a score: 100 when the question is allowed, 0 when it is blocked.
  • out_of_scope writes a polite refusal. Nothing is searched.
  • retrieve asks for the search tool by writing a tool call; tool_retrieve is a prebuilt ToolNode that runs it. The search is a tool, not a node of its own.
  • grade_documents asks the model whether the returned chunks are relevant to the question.
  • rewrite_query asks the model for a better question, at temperature 0.3.
  • generate_answer writes the answer from the chunks; output_guardrail checks that answer before it leaves.

The edges and the routing rule

python
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:

python
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.

ExampleThe project's graph wiring with stub nodes, run on LangGraph 1.2
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}")

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_documents chose rewrite_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

The Airflow ingestion DAG · from the Complete AI Security Course in 8 Hours video · 6:31:21 to 6:32:08

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.

Five Airflow tasks in a chain, setup_environment, fetch_daily_papers, index_papers_hybrid, generate_daily_report and cleanup_temp_files, of which the first four are PythonOperator tasks and the last is a BashOperator; fetch_daily_papers writes to the papers table in Neon Postgres, index_papers_hybrid writes to the OpenSearch index arxiv-papers-chunks, the report task reads the XCom values of both, and the DAG runs at 06:00 UTC, Monday to Friday.

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.

python
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_task

Shown 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-5 means 06:00 UTC on Monday to Friday.
  • fetch_daily_papers calls the arXiv API, downloads and parses the PDFs and writes one row per paper to the papers table in Postgres.
  • index_papers_hybrid splits each paper into chunks, asks the embedding API for one vector per chunk and writes text and vector to the OpenSearch index arxiv-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_results and hybrid_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
FlowFixed: search, prompt, one model callA graph that branches and can loop once
Model calls for an answered questionOneTwo in the video's trace: grading and generation
GuardrailsNoneOne node on the question, one on the answer
Off-topic questionSearched and answered from whatever was foundRefused before any search
Poor search resultsAnswered anywayQuestion rewritten, searched again
Answer cacheYesNo

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.
Watch out. Both guardrail nodes in the project fail open: when the call to the guardrail service raises an error, or no guardrail id is configured, the node logs a warning and lets the request through with a score of 100. That keeps the API answering during an outage, and it means an outage of the guardrail is also an outage of the protection. Decide which you want, and alert on the fallback.
Try it yourself
  • In the graph example, change max_retrieval_attempts in Context to 1 and run it: "Find the paper on VPO" now goes straight from grade_documents to generate_answer, with no rewrite.
  • In the API example, send {"query": ""} to /api/v1/ask: the reply is 422 with String should have at least 1 character, from min_length=1.
  • In the API example, post the same question to /api/v1/ask with "top_k": 5: the pipeline counter goes to 2, because the key changed.
PreviousAgentOps

You understood something today that you didn't yesterday.