Event Mapping
@threadplane/ag-ui reduces AG-UI protocol events into the runtime-neutral Agent contract from @threadplane/chat. Every event that reaches the browser lands on a Signal your template already reads, which is why the same <chat> composition renders a run from any backend that speaks the protocol. This page is the compatibility map: first the streaming example as the worked case, then every event the adapter handles and the state it changes.
What the demo does
The Run tab shows the prebuilt <chat> composition in front of a one-node LangGraph graph served over AG-UI. The example configures nothing beyond the connection, so the demo opens on the composition's welcome screen — the heading "How can I help?" above a single input, with no suggestions projected into it.
Ask anything and the answer arrives a fragment at a time rather than in one block. The end-to-end test for this example asks "Tell me one quick fact about Angular signals in two sentences." Any prompt does the same work: the graph has one node, the node calls the model once, and the tokens reach the transcript while the call is still running.
That is the whole feature. It is also the smallest complete AG-UI run, which makes it the right place to read the event sequence off the wire.
How it is built
Four files carry the demo: a graph with one node, a FastAPI server that fronts it with the AG-UI protocol, an application config that registers the agent, and a component that hands the agent to <chat>. Open the Code tab to read them in place.
The graph that streams tokens
The graph uses LangGraph's MessagesState, so its state is the message list and nothing else. One node reads the system prompt from disk, prepends it to the conversation, and awaits a single model call. The model client is constructed with streaming=True, which is what makes the call emit token chunks the adapter can forward while the node is still awaiting its result.
"""
LangGraph Streaming Graph
A minimal StateGraph that demonstrates real-time token streaming from an LLM.
Uses LangGraph's MessagesState for compatibility with the LangGraph SDK client.
The graph uses LangSmith for observability — every invocation is traced
automatically when LANGCHAIN_TRACING_V2=true is set.
"""
from pathlib import Path
from langgraph.graph import StateGraph, MessagesState, END
from langgraph.checkpoint.memory import MemorySaver
from langchain_openai import ChatOpenAI
from langchain_core.messages import SystemMessage
PROMPTS_DIR = Path(__file__).parent.parent / "prompts"
def build_streaming_graph():
"""
Constructs the LangGraph StateGraph for streaming.
The graph has a single node that calls the LLM with the system prompt
and user message. Uses MessagesState so the LangGraph SDK can send
and receive messages directly.
Returns:
A compiled StateGraph ready for invocation
"""
llm = ChatOpenAI(model="gpt-5-mini", streaming=True)
async def generate(state: MessagesState) -> dict:
"""
Generate a streaming response from the LLM.
Reads the system prompt and prepends it to the conversation,
then invokes the LLM with streaming enabled.
"""
system_prompt = (PROMPTS_DIR / "streaming.md").read_text()
messages = [SystemMessage(content=system_prompt)] + state["messages"]
response = await llm.ainvoke(messages)
return {"messages": [response]}
graph = StateGraph(MessagesState)
graph.add_node("generate", generate)
graph.set_entry_point("generate")
graph.add_edge("generate", END)
return graph.compile(checkpointer=MemorySaver())
# The graph instance — referenced by server.py
graph = build_streaming_graph()The graph compiles with MemorySaver(). The ag-ui-langgraph wrapper calls graph.aget_state(config) after the stream drains to build its closing snapshots, and a graph compiled without a checkpointer cannot answer that call. An in-memory saver is a development choice: a deployment that must survive a restart needs a durable one behind the same checkpointer= argument.
Serving the graph over AG-UI
The backend is a FastAPI application. LangGraphAgent wraps the compiled graph, add_langgraph_fastapi_endpoint mounts it at /agent as an AG-UI event stream, and /ok is a plain health check.
from fastapi import FastAPI
from ag_ui_langgraph import LangGraphAgent, add_langgraph_fastapi_endpoint
from .graph import graph
agent = LangGraphAgent(name="streaming", graph=graph)
app = FastAPI(title="cockpit-ag-ui-streaming")
add_langgraph_fastapi_endpoint(app, agent, path="/agent")
@app.get("/ok")
def ok() -> dict:
return {"ok": True}That wrapper is the piece that translates LangGraph's own stream into the protocol events listed below.
Providing the agent
provideAgent() from @threadplane/ag-ui registers the agent once for the whole application, and it is the only provider the <chat> composition requires. Your own application passes { url: 'https://your-backend.example.com/agent' } directly; this example passes a factory because it resolves its endpoint at runtime from the host that serves the demo.
import { injectCockpitRuntimeConnection } from '@threadplane/cockpit-telemetry';
import { ApplicationConfig } from '@angular/core';
import { provideAgent } from '@threadplane/ag-ui';
export const appConfig: ApplicationConfig = {
providers: [
provideAgent(() => {
const connection = injectCockpitRuntimeConnection();
if (connection.adapter !== 'ag-ui') {
throw new Error('incompatible runtime');
}
return {
url: connection.url,
};
}),
],
};The component that reads the Signals
The component is three lines of behavior. injectAgent() returns the agent the provider configured, and <chat> takes it as an input. Message rendering, the input, the typing indicator, and error display all live inside the composition, which reads the same Signals this page maps.
import { Component } from '@angular/core';
import { ChatComponent } from '@threadplane/chat';
import { injectAgent } from '@threadplane/ag-ui';
import { ExampleChatLayoutComponent } from '@threadplane/example-layouts';
/**
* Streaming demo — simplest possible @threadplane/chat integration with AG-UI.
*
* Retrieves the agent with injectAgent() (provided by provideAgent /
* provideFakeAgent) and passes it to the prebuilt <chat> composition. The
* composition handles message rendering, input, typing indicator, and error
* display internally.
*
* Demonstrates the chat-runtime decoupling: same <chat> composition as the
* LangGraph cockpit, AG-UI runtime instead of LangGraph.
*/
@Component({
selector: 'app-streaming',
standalone: true,
imports: [ChatComponent, ExampleChatLayoutComponent],
template: `
<example-chat-layout>
<chat main [agent]="agent" class="flex-1 min-w-0" />
</example-chat-layout>
`,
})
export class StreamingComponent {
protected readonly agent = injectAgent();
}injectAgent() must run inside an Angular injection context: a field initializer, as it is here, or a constructor body.
The events one turn produces
One prompt against this graph produces the following sequence. The order is the wrapper's; the effects are the reducer's.
RUN_STARTED.status()becomesrunning,isLoading()becomestrue, anderror(),interrupt(),customEvents(), and the subagent map are cleared for the new run. The user message is already in the transcript at this point, becausesubmit()appends it optimistically before the run opens.STEP_STARTEDnaming the node,generate. The reducer has no case for step events, so this one changes nothing. It is on the wire for consumers that want node boundaries.TEXT_MESSAGE_STARTon the first non-empty token chunk. The reducer creates the assistant message slot and marks it as streaming. This is the moment the empty bubble appears under your prompt.TEXT_MESSAGE_CONTENT, once per chunk, each carrying adelta. The reducer appends the delta to that message's content, and the composition renders the growing text. This is the streaming the demo exists to show.TEXT_MESSAGE_ENDwhen the model call ends. The reducer treats it as a no-op: the content is already accumulated in the Signal.STATE_SNAPSHOTwhen the node exits, carrying the graph state. The reducer replacesstate()with it and runs the citations bridge over the transcript. This graph publishes nostate.citations, so nothing is attached.STEP_FINISHED, then a closingSTATE_SNAPSHOTandMESSAGES_SNAPSHOTread back from the checkpoint. The snapshot message carries the final message id rather than the streaming chunk id used in step 3, so the reducer maps the in-flight assistant message onto its snapshot counterpart instead of leaving two bubbles, and re-applies citations fromstate()against the new ids.RUN_FINISHED. The run is marked complete,status()returns toidle, andisLoading()becomesfalse.
RAW events are interleaved through the stream as well, one per upstream LangGraph event. The adapter ignores them, as it ignores any event type it does not know.
Event reference
The adapter exposes messages, status, isLoading, error, toolCalls, state, interrupt, interruptSession, subagents, and customEvents as Signals, events$ as an Observable, and clientTools as a capability. Alongside submit, retry, stop, and regenerate, it provides the hydration Promise ready, reconcileInterrupt(), and dispose(). The tables below name the surface each event writes to.
Run lifecycle
| AG-UI event | Surface | Behavior |
|---|---|---|
RUN_STARTED | status, isLoading, error, interrupt, interruptSession, customEvents, subagents | Sets status to running and isLoading to true; clears the error, visible interrupt, custom events, subagents, and unfinished tool-argument buffer. A matching resume becomes acknowledged; the claim is retained until a terminal outcome, rather than treated as completed. Ignored when the run id does not match this delivery. |
RUN_FINISHED (no outcome, or a success outcome) | status, isLoading, messages, interruptSession | Settles the run and returns to idle. A collected compatibility interrupt becomes pending at this boundary; otherwise the run completes. Ignored on the same run-id gate as RUN_STARTED. |
RUN_FINISHED (outcome { type: 'interrupt' }) | interrupt, interruptSession, status, isLoading | Records the full native batch as pending and returns to idle. In auto or protocol, projects { id, value: { interrupts, runId }, resumable: true }; explicit command profiles project the compatibility interrupt. In default auto mode the native batch takes precedence over compatibility events in either arrival order. |
RUN_FINISHED (outcome { type: 'cancelled' }) | status, isLoading, messages, state, interruptSession | Settles delivery as aborted, returns to idle, and sets no error. The producer stopped the run on request. state rolls back to the last committed value. A resume attempt in flight lands in a recoverable interrupt phase: interruptSession moves to recovery-required, or to uncertain if no RUN_STARTED arrived. |
RUN_FINISHED / RUN_ERROR with usage | usage | Replaces the optional usage signal with { entries } for that run. Recorded only when the run finalizes (any outcome) and on RUN_ERROR; cleared at the next RUN_STARTED. |
RUN_FINISHED (success with pendingToolCallIds) | clientTools.pending | The declared list is the authoritative pending set for client tools, and an empty list means nothing is pending. Absent means derive from the stream. The set is recorded only when the run finalizes and is cleared at the next RUN_STARTED. |
RUN_ERROR | status, isLoading, error | Ends the run as failed, sets status to error, stops loading, and stores the event's message as an error (the raw event when no message is present). Ignored on the same run-id gate as RUN_STARTED. |
STEP_STARTED, STEP_FINISHED | none | Not reduced. Node boundaries pass through untouched. |
Messages
| AG-UI event | Surface | Behavior |
|---|---|---|
TEXT_MESSAGE_START | messages | Creates or reuses the assistant message slot for messageId and marks it streaming for this run. |
TEXT_MESSAGE_CONTENT | messages | Appends delta to that message's content. |
TEXT_MESSAGE_END | none | A no-op. The content is already accumulated from the deltas. |
REASONING_MESSAGE_START | messages | Creates or reuses the slot with an empty reasoning field and records a start time. |
REASONING_MESSAGE_CONTENT, REASONING_MESSAGE_CHUNK | messages | Appends delta to message.reasoning. Both event names are handled identically. |
REASONING_MESSAGE_END | messages | Sets reasoningDurationMs from the browser-measured start and end times. |
MESSAGES_SNAPSHOT | messages, toolCalls | Replaces the message list with the snapshot, bridges each assistant message's toolCalls onto toolCallIds, merges snapshot-only tool calls into toolCalls by id, keeps the run's in-flight assistant message aligned across the id change, and re-applies citations from the current state. |
MESSAGES_SNAPSHOT (tool messages with part-list content) | messages | Text parts and URL-sourced image parts become chat content blocks. Every other part is preserved under extra['ag-ui'].parts. Snapshots do not populate ToolCall.parts; only TOOL_CALL_RESULT does. |
Reasoning and text write to the same message when they share a messageId, so one assistant bubble can carry both reasoning and content. The duration is measured in the browser between the start and end events: it is useful for display, not for billing or tracing.
Tool calls
| AG-UI event | Surface | Behavior |
|---|---|---|
TOOL_CALL_START | toolCalls, messages | Adds a tool call with status running; when parentMessageId is present, links the call to that assistant message, creating the message slot if the turn produced no text. |
TOOL_CALL_ARGS | toolCalls | Appends delta to a per-call buffer, parses the buffer whenever it becomes valid JSON, and keeps the last good args while more fragments arrive. |
TOOL_CALL_END | toolCalls | Parses the accumulated buffer one final time, discards the buffer, and marks the call complete. |
TOOL_CALL_RESULT | toolCalls | Stores the result on the matching call. A string is parsed as JSON when possible. For a content-part list, result is the text of the parts, parsed as JSON only when exactly one text part holds valid JSON. Every part is kept under parts, and an all-media list yields an empty-string result. |
Argument deltas are fragments of one JSON document, not standalone JSON, which is why the reducer buffers rather than parses each delta. If the buffer never parses, args stays at the last good value, or {} when there was none.
{ type: 'TOOL_CALL_START', toolCallId: 'search-1', toolCallName: 'search' }
{ type: 'TOOL_CALL_ARGS', toolCallId: 'search-1', delta: '{"q":"Angular"}' }
{ type: 'TOOL_CALL_RESULT', toolCallId: 'search-1', content: { hits: 3 } }
{ type: 'TOOL_CALL_END', toolCallId: 'search-1' }The resulting entry in agent.toolCalls():
{
id: 'search-1',
name: 'search',
args: { q: 'Angular' },
result: { hits: 3 },
status: 'complete',
}State
| AG-UI event | Surface | Behavior |
|---|---|---|
STATE_SNAPSHOT | state, messages | Replaces state with snapshot and runs the citations bridge over the transcript. |
STATE_DELTA | state, messages | Applies the JSON Patch operations in delta to a clone of the current state, then runs the same bridge. |
STATE_DELTA carries standard JSON Patch operations:
{
type: 'STATE_DELTA',
delta: [
{ op: 'replace', path: '/topic', value: 'support' },
{ op: 'add', path: '/caseId', value: 'case-123' },
],
}The citations bridge reads state.citations, an object keyed by message id, and copies the entries onto the matching messages. See Citations for the accepted shapes.
Custom events and interrupts
| AG-UI event | Surface | Behavior |
|---|---|---|
CUSTOM named on_interrupt | interrupt, interruptSession, status, isLoading | Collects a compatibility interrupt, parsing a string value as JSON. The batch becomes pending at a terminal event or clean transport close. A native batch retains precedence in auto mode. |
CUSTOM named state_update with an object value | customEvents, events$ | Appends { name, data } to customEvents and emits { type: 'state_update', data } on events$. |
CUSTOM (every other name) | customEvents, events$ | Appends { name, data } to customEvents and emits { type: 'custom', name, data } on events$. |
Every non-interrupt custom event is fanned out to two surfaces, not one. Reach for the customEvents() Signal when you want an accumulated per-run snapshot for reactive rendering, and for events$ when you want a transient stream for side effects or telemetry. The Signal is documented on injectAgent() and toAgent().
Subagents and activity
| AG-UI event | Surface | Behavior |
|---|---|---|
SUBAGENT_STARTED | subagents | Creates the entry for subagentRunId with status running, or fills in the identity of an entry a child content event created first. toolCallId becomes parentToolCallId, else the id an earlier activity entry recorded, else subagentRunId. |
SUBAGENT_FINISHED | subagents | Marks the entry complete and records result. An outcome of { type: 'suspended' } keeps it running. |
SUBAGENT_ERROR | subagents | Marks the entry error and records message. |
ACTIVITY_SNAPSHOT | subagents | Creates or merges an activity entry; replace: true overwrites its content rather than merging. Entries with activityType: 'subagent' project into subagents(). |
ACTIVITY_DELTA | subagents | Applies a JSON Patch to an existing entry's content. An unknown id or a malformed patch is dropped rather than thrown. |
Text and tool events may also carry a subagentRunId. An attributed content event feeds that child's card and never the parent transcript, and one that arrives before SUBAGENT_STARTED still gets a card, which SUBAGENT_STARTED then fills in with identity. A suspended outcome keeps the card running: the run resumes under the same id, and the pause itself surfaces through interrupt(), not through the card.
Everything else
| AG-UI event | Surface | Behavior |
|---|---|---|
RAW and any unrecognized type | none | Ignored. A newer protocol version does not crash the adapter. |
Submit, stop, and failures
agent.submit({ message }) builds a local user message, appends it to messages and the source agent's list, and starts the run. Omitting message starts a run without appending a user message, provided no interrupt blocks new input. submit({ resume }) requires a pending batch and claims it before dispatch. Native responses must cover every pending ID exactly once; a single unwrapped value is accepted only for a single interrupt. Ordinary input and regeneration cannot abandon a pending approval.
agent.retry() replays captured input without appending the user message again. For a resume, an error carrying requestNotDispatched: true retains the pending batch and exact decision for retry. Uncertain delivery or failure after acknowledgement requires reconcileInterrupt() with an application-provided authoritative reconciler before retrying. A failed resume restores committed protocol messages and state; provisional stream changes do not become the next request's baseline, though failed messages may remain visible in the local transcript. Native snapshots are authoritative, while compatibility pause state commits at a terminal event or clean close.
interruptTransport selects auto, protocol, legacy-command, or mastra-command. In auto, native outcomes use top-level resume; compatibility payloads use forwardedProps.command.resume, with command.interruptEvent for Mastra payloads carrying toolCallId. Select a command profile explicitly if the backend requires it even when native events are present.
agent.stop() ends the current run locally as aborted, clears the error, and calls source.abortRun(). Whether the backend stops producing depends on the AG-UI source implementation.
onRunFailed comes from the AG-UI subscriber API rather than from a protocol event. The adapter treats it as a run failure: status becomes error, isLoading becomes false, and error holds the thrown value. An abort that the component requested is not reported as an error.
Unsupported protocol areas
The adapter carries interrupts, subagents, activity, state, and client tools. It does not implement history or time travel. Unknown protocol events are ignored rather than treated as errors, so a backend ahead of this adapter degrades to the events it does know.