How does a LlamaIndex Workflow decide which @step runs next?
answer
- No edge list is ever declared
- Annotations declare consume and emit
- StartEvent in, StopEvent out
- Union return type means a branch
- Validation catches orphaned event types
basics
~20 sNothing declares the order. Each @step method's parameter type says which event it consumes and its return type says which events it emits, so LlamaIndex wires steps by matching event types. A run starts with StartEvent and ends when a step returns StopEvent.
solid answer
~50 sA LlamaIndex `Workflow` is a subclass whose async methods are decorated with `@step`. There is no edge list: the type annotation on a step's event parameter declares what that step consumes, and its return annotation declares what it may emit. At construction the framework builds the graph from those annotations and validates it — an event that nothing consumes, or a step consuming an event nothing produces, raises a validation error before anything runs. Execution begins by delivering a `StartEvent` and finishes when some step returns a `StopEvent`, whose `result` becomes the run's value. Branching is expressed by a union return type: return `EventA | EventB` and whichever you actually return decides the next step. Because dispatch is by event type, adding a branch means adding an event class and a step, not editing a router. Steps share data through `Context` (`await ctx.store.set(...)` / `get(...)`) and can push UI events with `ctx.write_event_to_stream(...)`.
code
python · 28 linesimport asyncio
from llama_index.core.workflow import (
Context,
Event,
StartEvent,
StopEvent,
Workflow,
step,
)
class DraftEvent(Event):
draft: str
class MyFlow(Workflow):
@step
async def write(self, ctx: Context, ev: StartEvent) -> DraftEvent:
await ctx.store.set("topic", ev.topic)
return DraftEvent(draft=f"draft about {ev.topic}")
@step
async def review(self, ctx: Context, ev: DraftEvent) -> StopEvent:
topic = await ctx.store.get("topic")
return StopEvent(result=f"{ev.draft} (reviewed, topic={topic})")
asyncio.run(MyFlow(timeout=60).run(topic="otters"))go deeper
Recall the three moving parts: methods marked @step, events as typed classes, and a run that begins with StartEvent and ends when a step returns StopEvent.
Explain that dispatch comes from type annotations — parameter type consumes, return type emits — and show branching via a union return rather than any conditional-edge API.
Demonstrate debugging instincts: a stalled workflow means an unconsumed event or an unsatisfied collect_events, and the default 45-second timeout hides as an unexplained failure until you raise it deliberately.
Argue the design tradeoff: event typing makes new capabilities additive and refactors cheap, at the price of control flow being implicit, which raises the bar for naming discipline and graph visualisation in a large team.
## Events are the wiring The Workflow API in llama-index-core 0.14 is event-driven rather than graph-declared. You subclass `Workflow`, define `async def` methods decorated with `@step`, and annotate types: ``` async def my_step(self, ctx: Context, ev: RetrievedEvent) -> AnsweredEvent | StopEvent ``` The framework reads those annotations: this step is subscribed to `RetrievedEvent`, and it may emit `AnsweredEvent` or `StopEvent`. When a step returns an event, the runtime looks up which step consumes that type and schedules it. Steps that are not waiting on the emitted type simply do not run. The consequence worth internalising is that **control flow lives in your event classes**, not in a wiring block — the type system is the router. Events are subclasses of `Event` (a Pydantic-style model), so they carry typed payload fields. Two steps that would consume the same payload still need distinct event classes if they must be scheduled independently; reusing one event type for two meanings is the most common design mistake. ## Start and stop Every run is kicked off by a `StartEvent`, which carries whatever keyword arguments you passed to `run()`. Exactly one step should consume it. The run terminates when any step returns a `StopEvent`; the value in its `result` field is what `await handler` yields. A workflow with no reachable `StopEvent` will run until the timeout rather than returning. ## Validation Because the graph is derived, LlamaIndex can check it. If a step emits an event type no step consumes, or a step waits on an event no step produces, you get a workflow validation error at run time rather than a silent hang. This catches the classic refactoring bug where you rename an event class and update only one side. ## Branching, fan-out and fan-in - **Branching**: return a union type and pick at run time. `-> SuccessEvent | RetryEvent` is a two-way branch with no conditional-edge API. - **Fan-out**: call `ctx.send_event(...)` several times to dispatch multiple events from one step. Unlike returning, `send_event` can emit many events, so it is the mechanism for map-style parallelism. - **Fan-in**: a step that must wait for several results calls `ctx.collect_events(ev, [TypeA, TypeB, TypeC])`. On each arriving event it returns `None` — meaning "not ready, do nothing" — until all the expected types have arrived, at which point it returns the collected list and the step proceeds. - **Concurrency**: `@step(num_workers=4)` lets that step process up to four events in parallel, which is what makes a fan-out actually concurrent instead of serialised. ## Context `Context` is the per-run shared object. Steps read and write it through the async store — `await ctx.store.set("key", value)` and `await ctx.store.get("key", default=None)` — which is how data that should not be threaded through every event travels. `ctx.write_event_to_stream(ev)` pushes an event into the run handler's stream so a UI can render progress without waiting for `StopEvent`. ## Running and observing `handler = workflow.run(topic="x")` returns immediately with a handler; `await handler` gives the final result, and `async for ev in handler.stream_events()` yields streamed events as they occur. `Workflow(timeout=..., verbose=True)` sets a wall-clock budget for the whole run — the default is 45 seconds, which is short for anything with several LLM calls, so raise it deliberately rather than discovering it as a mysterious timeout. `verbose=True` prints step transitions, and the utility drawing helpers can render the derived graph, which is the fastest way to check that the wiring you meant is the wiring you got. ## Why this design Compared with a declared node/edge graph, event typing makes composition cheap: a new capability is a new event class plus a step, and existing steps do not change. It also makes the failure surface unusual — a workflow that "does nothing" is almost always an event nobody consumes, and a workflow that hangs is almost always a fan-in that never received its full set. Both are diagnosed by listing which event types are produced and which are consumed, then comparing. Agents in LlamaIndex are themselves Workflows, so the same machinery — steps, events, `Context`, streaming — underlies the prebuilt agent classes and any custom orchestration you write.
- How do you make one step wait for several parallel results before continuing?Fan out with repeated `ctx.send_event(...)` calls, then have the collecting step call `ctx.collect_events(ev, [TypeA, TypeB, TypeC])`. That call returns `None` for every arrival until all listed types have been seen, so the step effectively no-ops until the set is complete, then returns the collected events and proceeds. Pair it with `@step(num_workers=n)` so the fanned-out work actually runs concurrently.
- A workflow you just refactored hangs and then times out. What is the first thing you check?Compare produced event types against consumed ones. A hang is almost always either an event emitted with no subscriber, or a `collect_events` waiting on a type that is no longer produced — often after renaming an event class and updating only one side. Turn on `verbose=True` to see which steps fired, and render the derived graph to confirm the wiring matches your intent before touching prompts or models.
- Why would you push events with ctx.write_event_to_stream instead of just returning them?Returning an event schedules the next step; it is control flow. `ctx.write_event_to_stream(ev)` publishes to the run handler's stream without scheduling anything, which is how you surface progress — status lines, partial tokens, intermediate findings — to a UI while the run continues. Consumers read them with `async for ev in handler.stream_events()` rather than waiting for the final `StopEvent`.
saying these in an interview costs you the question
- Looks for an add_edge or wiring call that does not exist
- Thinks steps execute in the order they are defined
- Assumes returning an event type nobody consumes is harmless
- Believes fan-out is parallel without num_workers
- Forgets the run ends only when a step returns StopEvent