Write a custom wire codec
Write encode and decode for events AI Transport has never seen. Declare the descriptor table once and the client and agent transports carry your events with no other change.
What you build
The same chat app the other two quickstarts build, against a provider AI Transport ships no codec for:
- Tokens stream from your own event types to every connected client in realtime.
- A stop button cancels the run in flight, and every tab sees it close.
- A second tab shows the same stream, because every subscribed client receives the same messages.
- Reloading reads the conversation back from your own store.
Only the codec is new. The client and agent halves are the same code the Vercel and OpenAI quickstarts use, because the transports take any codec.
Prerequisites
- Node.js 22 or later.
- An Ably account with an API key.
- A model provider that gives you a stream of events. This guide invents one so you can follow it with no provider at all.
Install dependencies
Install the AI Transport SDK, the Ably client, and Next.js. No provider SDK is involved:
npm install @ably/ai-transport ably next react react-domSet up authentication
Create an auth endpoint at /api/auth/token that returns an Ably JWT to the browser, signed with the user's client id and the channel capabilities they need, as described in Set up authentication.
The client below fetches from it with authUrl: '/api/auth/token'. The agent route runs on your own infrastructure, so it uses the API key directly.
Configure the channel rule
AI Transport streams each response by appending tokens to a single channel message. That needs the Message annotations, updates, deletes, and appends rule (mutableMessages) on the namespace your conversations live on. Enable it once per Ably app through the dashboard, Control API, or CLI.
Declare your events
Create app/codec/events.ts. This example code provides a minimal set of example events that you might want to send over the transport.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
/** One message in your provider's world. */
export interface Turn {
id: string;
role: 'user' | 'assistant';
text: string;
}
/** Client to agent. Every input carries a `kind` the codec dispatches on. */
export type MyInput = { kind: 'user-message'; message: Turn };
/** Agent to client. Text arrives as a start / delta / end trio. */
export type MyOutput =
| { type: 'text-start'; id: string }
| { type: 'text-delta'; id: string; delta: string }
| { type: 'text-end'; id: string }
| { type: 'tool-call'; callId: string; toolName: string; args: unknown }
| { type: 'done'; reason: string };Write the codec
Create app/codec/index.ts. A codec describes how your events should be encoded to, and decoded from, Ably messages. The codec is defined as a table of input and output events, and how they map to an Ably message.
There are three kinds of mapping: stream, event, and drop. Use the stream group for text deltas, where the content of the message should be appended together to create the final output. A stream is made of a start, a delta, and an end. Use event for a single Ably message publish for that event. Use drop to keep that event out of the Ably channel:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
import { defineCodec, jsonField, strField } from '@ably/ai-transport';
import type { MyInput, MyOutput } from './events';
// Each field binds one header key to one property, so encode and decode
// can never disagree about where a value lives.
const idField = strField('id');
const callIdField = strField('callId', '');
const toolNameField = strField('toolName', '');
const reasonField = strField('reason', '');
const argsField = jsonField('args');
const messageIdField = strField('messageId', '');
export const myCodec = defineCodec<MyInput, MyOutput>()({
adapterTag: 'my-provider',
output: ({ event, stream }) => [
// One growing message per assistant reply. `streamId` is what ties the
// three phases together on the wire.
stream('text', {
streamId: (chunk) => chunk.id,
fields: [idField],
start: { type: 'text-start' },
delta: { type: 'text-delta', field: 'delta', decode: ({ rebuild }) => rebuild([idField]) },
end: { type: 'text-end' },
}),
// Discrete: one publish each.
event('tool-call', { fields: [callIdField, toolNameField, argsField] }),
event('done', { fields: [reasonField] }),
],
input: ({ batch }) => [
// A user message fans out into one wire event per part. This provider has
// one part kind, so the table has one entry.
batch('user-message', {
explode: (input) => [{ type: 'text', text: input.message.text }],
partTypeOf: (part) => part.type,
parts: (p) => [
p('text', {
data: { encode: (part) => part.text, decode: (d) => ({ text: String(d ?? '') }) },
}),
],
// Stamped on every part, so the decode side rebuilds the envelope from any one.
messageHeaders: (input) => ({
codecHeaders: { messageId: input.message.id },
transportHeaders: { role: input.message.role },
}),
assemble: (part, ctx) => ({
message: { id: messageIdField.read(ctx.codecHeaders), role: 'user', text: part.text },
}),
}),
],
});Every event is described or dropped
The output table should cover all your event types. An output type that is neither described nor dropped throws an error on encode rather than being skipped.
What a wire codec does not need
The transports accept a WireCodec, which is createEncoder, createDecoder, and the optional adapterTag you set above. Only the durable sessions act on adapterTag: they build their own channel and register the tag on it. The transports take a channel you built, so nothing registers the tag here. That is all defineCodec builds, so the descriptor tables above are the whole codec.
Merging events into messages is what durable sessions add, and the codecs the sessions consume are built by a different factory.
Create the agent route
Create app/api/chat/route.ts. This file shows the agent side of the application. The code creates the agent transport with your own codec, connects to the channel, locates the input event that was sent on the invocation, loads the previous conversation messages from your database, and calls your model provider. runModel stands in for that provider, yielding the events you declared. The output is piped directly into the run and onto the Ably channel, and the run carries a reason when it ends.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
// Agent-side. This runs on your own infrastructure, so it holds the key.
import { after } from 'next/server';
import * as Ably from 'ably';
import { createAgentTransport } from '@ably/ai-transport';
import { myCodec } from '../../codec';
import { appendTurn, loadConversation } from '../../store';
const ably = new Ably.Realtime({ key: process.env.ABLY_API_KEY });
// Your provider, reduced to what this guide needs: an async iterable of MyOutput.
async function* runModel(conversation, signal) {
const id = crypto.randomUUID();
yield { type: 'text-start', id };
for await (const fragment of yourProvider(conversation, { signal })) {
yield { type: 'text-delta', id, delta: fragment };
}
yield { type: 'text-end', id };
yield { type: 'done', reason: 'stop' };
}
export async function POST(req) {
const { conversationId, eventId } = await req.json();
const transport = createAgentTransport({
channel: ably.channels.get(conversationId),
codec: myCodec,
clientId: 'agent',
});
await transport.connect();
after(async () => {
try {
// locateInput scans channel history for the input carrying this event id.
// It is the one history read the SDK does for you.
const trigger = await transport.locateInput(eventId);
const prompt = trigger?.inputs.find((input) => input.kind === 'user-message');
// Passing inputCodecMessageId lets a cancel keyed on the input find this
// run, including one that arrived before the run opened.
const run = transport.openRun({ inputCodecMessageId: trigger?.meta.codecMessageId });
const conversation = [...loadConversation(conversationId), prompt.message];
appendTurn(conversationId, prompt.message);
const reply = { id: crypto.randomUUID(), role: 'assistant', text: '' };
const { reason } = await run.pipe(
(async function* () {
for await (const chunk of runModel(conversation, run.abortSignal)) {
if (chunk.type === 'text-delta') reply.text += chunk.delta;
yield chunk;
}
})(),
);
await run.end({ reason });
if (reason === 'complete') appendTurn(conversationId, reply);
} finally {
transport.close();
}
});
return Response.json({ ok: true });
}Create the chat component
Create app/chat.tsx. The client publishes the user's input, then calls the agent route with the eventId that publish returned:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
'use client';
// Client-side. The browser fetches a token; never put an API key here.
import { useEffect, useRef, useState } from 'react';
import * as Ably from 'ably';
import { createClientTransport } from '@ably/ai-transport';
import { myCodec } from './codec';
import { createFold } from './fold';
export function Chat({ conversationId, seed }) {
const [messages, setMessages] = useState(seed);
const [input, setInput] = useState('');
const [activeRunId, setActiveRunId] = useState(undefined);
const transportRef = useRef(null);
useEffect(() => {
const ably = new Ably.Realtime({ authUrl: '/api/auth/token' });
const fold = createFold();
const transport = createClientTransport({
channel: ably.channels.get(conversationId),
codec: myCodec,
});
transportRef.current = transport;
transport.subscribe((event) => {
if (event.kind === 'run-lifecycle') {
setActiveRunId(event.event.type === 'start' ? event.event.runId : undefined);
return;
}
fold.apply(event);
setMessages([...seed, ...fold.render()]);
});
void transport.connect();
return () => {
transport.close();
ably.close();
};
}, [conversationId]);
const send = async (text) => {
const sent = await transportRef.current.publishInput({
kind: 'user-message',
message: { id: crypto.randomUUID(), role: 'user', text },
});
// The SDK publishes; waking the agent is yours.
await fetch('/api/chat', {
method: 'POST',
body: JSON.stringify({ conversationId, eventId: sent.eventId }),
});
};
return (
<div>
{messages.map((m, i) => (
<div key={i}><strong>{m.role}:</strong> {m.text}</div>
))}
<form onSubmit={(e) => { e.preventDefault(); void send(input); setInput(''); }}>
<input value={input} onChange={(e) => setInput(e.target.value)} />
{activeRunId ? (
<button type="button" onClick={() => transportRef.current.cancel(activeRunId)}>Stop</button>
) : (
<button type="submit">Send</button>
)}
</form>
</div>
);
}The input you publish is one of your own MyInput values, and your input table decides how it reaches the wire. publishInput emits a local echo to your subscribe handler, so the user's own message appears without waiting for the round trip.
Merge the model output into conversation messages
Create app/fold.ts.
Most model output comes in chunks, parts, or events. These events are streamed directly over the Ably channel, and delta events like text-streaming are appended together into a single message. You merge these events into whatever your UI renders. AI Transport doesn't do this for you. See durable sessions for the APIs that will merge events into conversation messages.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
// Client-side. Fold the transport's event stream into a renderable list.
export function createFold() {
const messages = new Map(); // codecMessageId -> { role, text, stepId, stepStartSerial }
const canonical = new Map(); // stepId -> the newest attempt's start serial
function apply(event) {
if (event.kind !== 'message') return;
const { codecMessageId, stepId, stepStartSerial, role } = event.meta;
if (!codecMessageId) return;
// A retry publishes under the same step id with a higher start serial.
// Record the newest, so render drops the attempt it replaced.
if (stepId && stepStartSerial && (canonical.get(stepId) ?? '') < stepStartSerial) {
canonical.set(stepId, stepStartSerial);
}
const entry = messages.get(codecMessageId) ?? { role, text: '', stepId, stepStartSerial };
for (const input of event.inputs) {
if (input.kind === 'user-message') entry.text = input.message.text;
}
for (const output of event.outputs) {
if (output.type === 'text-delta') entry.text += output.delta;
}
messages.set(codecMessageId, entry);
}
function render() {
return [...messages.values()].filter(
(m) => !m.stepId || canonical.get(m.stepId) === m.stepStartSerial,
);
}
return { apply, render };
}The output types this function ignores are the ones your interface does not render yet. Adding tool-call rendering means one more branch here.
Integrate with your database
Create app/store.ts. This Map stands in for whatever database your application already uses. You own the conversation at this level, so the store is part of the code here rather than hidden behind a helper:
1
2
3
4
5
6
7
8
9
10
11
// Stand-in for your own database. Replace both functions with real queries.
const conversations = new Map();
export function loadConversation(conversationId) {
return conversations.get(conversationId) ?? [];
}
export function appendTurn(conversationId, turn) {
const existing = conversations.get(conversationId) ?? [];
conversations.set(conversationId, [...existing, turn]);
}Wire it together
Create app/page.tsx. The page reads the stored conversation and passes it to the component as the seed, which is hydration you write yourself rather than something AI Transport provides:
1
2
3
4
5
6
7
import { Chat } from './chat';
import { loadConversation } from './store';
export default function Page() {
const conversationId = 'conversations:demo';
return <Chat conversationId={conversationId} seed={loadConversation(conversationId)} />;
}Run the app
Start the dev server:
npm run devOpen http://localhost:3000 in two tabs and send a message from one. Tokens stream into both. Open the channel in your Ably dashboard to watch what your descriptors produced: one message per assistant reply growing by append, and a discrete message per tool-call and done.
What happens when you send a message
publishInputruns your input table'suser-messagebatch, publishes the resulting wire events, and returns thecodecMessageIdandeventId.- Your POST wakes the agent. AI Transport sends no HTTP of its own; the endpoint is yours.
- The agent opens a run and pipes your provider's events. The output table decides which become appends on one message and which become their own.
- Every subscribed client decodes them back into your own event types through the same table, so encode and decode can never drift.
Understand the architecture
Your codec owns the mapping between your provider's events and Ably messages. AI Transport owns everything around it: the run and step lifecycle, cancel routing, steering, and history. Streaming covers who implements what, row by row, and codec architecture covers what the descriptor drivers do underneath.
Explore next
- Wire codec reference: the contract the transports require.
- Codec architecture: the encoder and decoder cores that read your descriptors.
- Wire protocol: the message names and headers your table writes.
- Token streaming: why a stream group becomes one growing message.
- Durable sessions: what merges the events into messages for you.