Streaming

Stream your agent's responses directly to clients in realtime over any Ably channel. The channel provides bi-directional control over your agent's operation.

Streaming allows your AI agents to stream their output directly to subscribed clients over an Ably channel. The channel provides bi-directional control over the agent's operation, allowing you to easily cancel, interrupt, or steer the output in realtime. The output reaches every subscribing client as it is produced, and any client that disconnects can automatically reconnect and resume from where it left off. Those properties come from the channel rather than from anything AI Transport adds on top: Ably orders and delivers every message on a fault-tolerant global network, at low latency from any region. Message ordering covers what order means when publishing from different regions.

Why an HTTP response is not enough

A typical SSE streaming response over a single HTTP connection can only serve one client, and doesn't allow for realtime control and steering over your agent's operation. A second device the user picks up sees nothing at all, because the response belongs to the request that opened it. There's also no built-in reconnect, resume or replay; Ably channels do all three for you, so a client that disconnects gets the rest of the response when it comes back.

Understand the model

One conversation is one Ably channel. There are two transports for connecting to the channel: a client transport that publishes the user's input and reads the agent's output, and an agent transport that opens a run and pipes the model's output into it. Everything published by either the client or agent is an ordinary channel message, so ordering, storage, and replay come from the channel.

Diagram showing a run as a group of channel messages, with lifecycle events marking its start and end, and the steps inside it carrying the agent's output

There are two lifecycle concepts: runs and steps. A run is one turn of agent work with a start, an end, and a reason it ended. A step is a single retryable unit of work within a run. Steps have an identifier, and any later step published with the same identifier supersedes an earlier one. This allows you to supersede failed output that has already been written to the channel by republishing to the same step.

The streaming layer adds cancellation, the run and step lifecycle, and steering on top of the channel. A client can cancel an inflight run to stop an agent, and every subscribed client can see the run end. A client can steer a run while the agent is streaming. The agent can see that steering message between model calls, and reports back on the run whether the steering message was considered in the output.

Client and agent example

Both transports attach to the same channel. The client publishes an input, reads the agent's output off its subscription, and calls its own endpoint to start the agent:

JavaScript

1

2

3

4

5

6

7

8

9

10

11

12

13

14

15

16

17

18

// Client-side. The browser fetches a token; never put an API key here.
import { createUIMessageCodec } from '@ably/ai-transport/vercel';

const codec = createUIMessageCodec();
const ably = new Ably.Realtime({ authUrl: '/auth' });
const transport = createClientTransport({
  channel: ably.channels.get('conversation-42'),
  codec,
  clientId: 'alice',
});
await transport.connect();

transport.subscribe((e) => {
  if (e.kind === 'message') render(e.meta.codecMessageId, e.outputs);
});

const sent = await transport.publishInput({ kind: 'message', payload: prompt });
await fetch('/api/chat', { method: 'POST', body: JSON.stringify({ eventId: sent.eventId }) });

The agent finds the input that woke it, opens a new run, and pipes the model's output back:

JavaScript

1

2

3

4

5

6

7

8

9

10

11

12

13

14

// Agent-side. This runs on your own infrastructure, so it holds the key.
import { createUIMessageCodec } from '@ably/ai-transport/vercel';

const codec = createUIMessageCodec();
const ably = new Ably.Realtime({ key: process.env.ABLY_API_KEY });
const transport = createAgentTransport({ channel: ably.channels.get(conversationId), codec });
await transport.connect();

const trigger = await transport.locateInput(eventId);
const run = transport.openRun({ inputCodecMessageId: trigger?.meta.codecMessageId });

await run.pipe(llmStream([...conversationFromYourDatabase, ...(trigger?.inputs ?? [])]));
await run.end({ reason: 'complete' });
transport.close();

openRun publishes the run start, pipe publishes the output inside a step, and end closes the run with a reason.

Know who implements what

Streaming splits the work of a conversation three ways:

Part of the systemWho implements it
Token streamingSDK
Run and step lifecycleSDK
Cancelling a run in flightSDK
Steering a run while it streamsSDK
Fan-out to several clientsAbly channel
Ordering, storage, replayAbly channel
Resume after a disconnectAbly channel
Where the conversation livesYour application
Turning decoded codec events into the messages your UI rendersYour application
Waking the agentYour application

What you still own

The conversation lives in whatever database your application already uses. You are responsible for providing the conversation history to clients, usually over a REST endpoint you own. You merge the model output chunks from the event stream into whatever your UI renders. You should also subscribe to the run and step lifecycle events, so you can detect superseded output on the client and update your UI.

You are also responsible for waking an agent, typically using an HTTP request that carries an invocation. The client should publish the input message to the transport, then call your own HTTP endpoint with the invocation to run the agent. Ably delivers the message and nothing more.

How to get started

To get started, see one of the quickstarts with the Vercel AI SDK or OpenAI.

Codecs are how you encode and decode your events into Ably messages on the channel. The SDK ships with codecs for the Vercel AI UI SDK format, and for the OpenAI Responses format. For any other format see the custom wire codec quickstart.