Workflows: a custom RAG pipeline with @step
A Workflow is a class whose steps pass events to each other, so you build a retrieve-then-answer pipeline you control step by step.
Last updated: 28 Sep, 2026 · LlamaIndex 0.14
A query engine hides retrieve and answer inside one call. A workflow splits them into named steps joined by events, so you can add a step, store something between steps, or change the order.
Declaring an event and steps
An Event is a small object one step returns and the next receives. A method marked @step becomes a step; its argument type says which event starts it.
from llama_index.core.workflow import Workflow, step, StartEvent, StopEvent, Event, Context
class Retrieved(Event):
question: str
context: strThe retrieve step
The first step takes the StartEvent, retrieves nodes, saves the source names in ctx.store, and returns a Retrieved event with the joined text.
@step
async def retrieve(self, ctx: Context, ev: StartEvent) -> Retrieved:
nodes = self.retriever.retrieve(ev.question)
await ctx.store.set("sources", sorted({n.metadata["file_name"] for n in nodes}))
context = "\n".join(n.get_content() for n in nodes)
return Retrieved(question=ev.question, context=context)The answer step
The second step receives the Retrieved event, builds the prompt, asks the stand-in model, reads the sources back out of the context store, and returns a StopEvent, which ends the run.
@step
async def answer(self, ctx: Context, ev: Retrieved) -> StopEvent:
prompt = f"---------------------{ev.context}---------------------\nQuery:{ev.question}\nAnswer:"
reply = self.llm.complete(prompt).text
sources = await ctx.store.get("sources")
return StopEvent(result=f"{reply} (from {', '.join(sources)})")View the code here
# Lamps
The LMP-204 desk lamp has a known cable fault. Stop using a lamp with a damaged cable and we will replace it free of charge.
All lamps come with a two year guarantee against electrical faults.
Bulbs are not covered by the refund policy once they have been used.
The LMP-310 floor lamp needs a bulb with an E27 fitting, which is sold separately.
# Refunds
You can get a full refund within 30 days of delivery. The money goes back to the card you paid with within 5 working days of us receiving the item.
Items bought in a sale can be refunded too, but the delivery charge is not returned.
To start a refund, open the order in your account and choose Return an item. Print the label and drop the parcel at any post office.
Personalised items cannot be refunded unless they arrive damaged.
# Delivery
Standard delivery takes 3 to 5 working days and is free on orders over 40.
Express delivery arrives the next working day if you order before 2pm. It costs 6.
We deliver to the mainland only. Parcels to islands take 2 extra working days.
If a parcel has not arrived after 10 working days, contact us and we will send a replacement.
import re
from llama_index.core.llms import CompletionResponse, CustomLLM, LLMMetadata
from llama_index.core.llms.callbacks import llm_completion_callback
def stems(text):
"""Words longer than three letters, cut to five letters, so refund and refunds match."""
return {w[:5] for w in re.findall(r"[a-z0-9-]+", text.lower()) if len(w) > 3}
class ExtractiveLLM(CustomLLM):
"""Answers with the context sentence that shares most words with the question."""
@property
def metadata(self):
return LLMMetadata(model_name="extractive")
@llm_completion_callback()
def complete(self, prompt, formatted=False, **kwargs):
context = prompt.split("---------------------")[1]
question = prompt.split("Query:")[1].split("Answer:")[0]
asked = stems(question)
sentences = [s.strip() for s in re.split(r"(?<=[.!?])\s+|\n+", context)]
sentences = [s for s in sentences if s and not s.startswith("#") and ": " not in s[:20]]
best = max(sentences, key=lambda s: len(asked & stems(s)), default="")
if len(asked & stems(best)) < 2:
return CompletionResponse(text="I could not find that in the documents.")
return CompletionResponse(text=best)
@llm_completion_callback()
def stream_complete(self, prompt, formatted=False, **kwargs):
yield self.complete(prompt)
The workflow answering two questions
The whole program. The two steps run in order, joined by the Retrieved event, and the context store carries the source names from the first step to the second.
import asyncio
from llama_index.core import Settings, SimpleDirectoryReader, VectorStoreIndex
from llama_index.core.workflow import Context, Event, StartEvent, StopEvent, Workflow, step
from llama_index.embeddings.huggingface import HuggingFaceEmbedding
from extractive_llm import ExtractiveLLM
Settings.embed_model = HuggingFaceEmbedding(model_name="sentence-transformers/all-MiniLM-L6-v2")
class Retrieved(Event):
question: str
context: str
class RAGWorkflow(Workflow):
def __init__(self, index, **kwargs):
super().__init__(**kwargs)
self.retriever = index.as_retriever(similarity_top_k=2)
self.llm = ExtractiveLLM()
@step
async def retrieve(self, ctx: Context, ev: StartEvent) -> Retrieved:
nodes = self.retriever.retrieve(ev.question)
await ctx.store.set("sources", sorted({n.metadata["file_name"] for n in nodes}))
context = "\n".join(n.get_content() for n in nodes)
return Retrieved(question=ev.question, context=context)
@step
async def answer(self, ctx: Context, ev: Retrieved) -> StopEvent:
prompt = f"---------------------{ev.context}---------------------\nQuery:{ev.question}\nAnswer:"
reply = self.llm.complete(prompt).text
sources = await ctx.store.get("sources")
return StopEvent(result=f"{reply} (from {', '.join(sources)})")
async def main():
docs = SimpleDirectoryReader("help").load_data()
index = VectorStoreIndex.from_documents(docs)
workflow = RAGWorkflow(index, timeout=60)
for question in ["How long until my refund money reaches my card?", "What are your opening hours?"]:
result = await workflow.run(question=question)
print(question)
print(" ", result)
asyncio.run(main())How long until my refund money reaches my card? The money goes back to the card you paid with within 5 working days of us receiving the item. (from delivery.md, refunds.md) What are your opening hours? I could not find that in the documents. (from delivery.md, refunds.md)
How the events moved through the run
- The start event carried the question into
retrieve, which returned aRetrievedevent. - The Retrieved event triggered
answer, because its argument type isRetrieved. - The context store passed the source names between steps without putting them in the event.
- The absent question still ran both steps; the stand-in model refused because no sentence matched.
Query engine vs a workflow
| Approach | Steps | Control |
|---|---|---|
| Query engine | Fixed: retrieve then answer | Little; one call |
| Workflow | Whatever steps you write | Full: add steps, branch, store state |
When a workflow is worth the extra code
- Adding a step, such as a check or a rewrite, between retrieval and answering.
- Branching on the question before deciding how to answer.
- Passing state between steps that should not travel inside the events.
Related
- Previous: Long-term memory: recalling past sessions
- Next: Multi-agent: routing between two agents
- See also: Query engines: the prompt the model receives
- Reference: Understanding: workflows
- Add a middle step that prints the question before answering.
- Store the number of retrieved nodes in
ctx.storeand print it in the answer. - Change
similarity_top_kto 1 and see whether the refund question still finds its sentence.
Slow is fine. Stopping is the only problem.