ClientTransport
The ClientTransport publishes to and receives from one Ably channel, through one codec. It publishes client inputs, classifies inbound messages into typed events, cancels and steers runs, and pages channel history. It builds no conversation tree and holds no run registry, so your own code merges the event stream into what your application renders.
Construct one with createClientTransport from the core entry point. The channel stays caller-owned: the transport subscribes its own listener on connect() and never detaches the channel.
1
2
3
4
5
6
7
8
9
10
11
12
13
import * as Ably from 'ably';
import { createClientTransport } from '@ably/ai-transport';
import { createUIMessageCodec } from '@ably/ai-transport/vercel';
const ably = new Ably.Realtime({ authUrl: '/api/auth/token' });
const transport = createClientTransport({
channel: ably.channels.get('conversation-42'),
codec: createUIMessageCodec(),
clientId: 'alice',
});
await transport.connect();Create a client transport
function createClientTransport<TInput, TOutput>(options: ClientTransportOptions<TInput, TOutput>): ClientTransport<TInput, TOutput>Construct a transport bound to a channel and a codec. Construction is synchronous. The constructor registers a channel state listener, which the transport uses to drain in-flight steering messages on continuity loss; the transport subscribes its message listener and publishes only after connect() resolves.
Parameters
optionsrequiredClientTransportOptionsReturns
ClientTransport<TInput, TOutput>. The transport, not yet connected.
Connect the transport
connect(): Promise<void>Subscribe the transport's listener to the channel and attach it, starting live event delivery. Every other method requires a successful connect() first.
Single-flight and idempotent: concurrent and repeat calls share one attempt, and a failed attempt is retried by the next call.
Returns
Promise<void>. Resolves once the channel is attached. A failure both rejects this call and is emitted on on('error'), as an ErrorInfo.
Subscribe to transport events
subscribe(handler: (event: TransportEvent<TInput, TOutput>) => void): () => voidSubscribe to classified transport events. Shorthand for on('event', handler).
Fires for live wire events and for the optimistic local echo that publishInput emits. History batches do not pass through here; they are returned from history() instead.
Delivery is synchronous and in registration order. Each event reaches every subscriber before the next is processed, and a throw from a handler is caught and logged, so the remaining subscribers still receive the event.
1
2
3
4
const off = transport.subscribe((event) => {
if (event.kind === 'message') fold(event.meta, event.outputs);
if (event.kind === 'run-lifecycle' && event.event.type === 'end') done(event.event.reason);
});Parameters
handlerrequiredTransportEventReturns
() => void. An unsubscribe function.
Subscribe to a specific stream
on(event: 'event' | 'ably-message' | 'error', handler: (payload) => void): () => voidSubscribe to one of the transport's three receive streams. subscribe is shorthand for the 'event' form.
eventTransportEvent values, fired once per inbound wire message that produces one.ably-messageAbly.InboundMessage, emitted after the matching typed event so a handler sees state an earlier subscriber merged.errorErrorInfo with a distinguishing code. A single decode failure drops that one message and emits the error rather than tearing down the stream.Parameters
eventrequired'event' | 'ably-message' | 'error'handlerrequiredFunctionReturns
() => void. An unsubscribe function.
Publish a client input
publishInput(event: TInput, opts?: PublishInputOptions): Promise<PublishInputResult>Publish one codec input event to the channel. Requires connect().
An input published under a fresh codec-message-id reaches subscribe handlers immediately as a local message event, with serial and versionSerial undefined, so the sender sees its own input before the channel echoes it back. When you supply opts.codecMessageId the input amends something already on the channel, and the transport emits no local event for it. The real echo later carries the same codecMessageId, so your own code can key on it and combine the local event and the echo into one entry.
1
2
3
4
const sent = await transport.publishInput({ kind: 'message', payload: message });
// The SDK publishes; waking the agent is yours.
await fetch('/api/chat', { method: 'POST', body: JSON.stringify({ eventId: sent.eventId }) });Parameters
eventrequiredTInputoptsoptionalPublishInputOptionsReturns
Promise<PublishInputResult>. Resolves once the input is published, or rejects with an ErrorInfo.
codecMessageIdStringeventIdStringevent-id set on the wire message, distinct from codecMessageId. This is what an agent's locateInput matches to find the input that woke an invocation.runIdPromise<string>input-codec-message-id matches this publish's codecMessageId. Never resolves for an input that triggers no run. Rejects on close() and on channel continuity loss; a rejection handler is pre-attached, so ignoring it raises nothing.Cancel a run
cancel(runId: string): Promise<void>Publish a cancel envelope targeting the given run. Requires connect().
Stateless: the transport holds no run registry, so supply a run id you captured from a run-lifecycle start event or from publishInput's runId promise. The envelope carries a per-cancel event id so a channel rewind can redeliver it, and the agent transport ignores that id, so a redelivered cancel has no extra effect.
Parameters
runIdrequiredStringReturns
Promise<void>. Resolves once the cancel envelope is published, or rejects with an ErrorInfo.
Steer an open run
steer(runId: string | Promise<string>, event: TInput): SteerResultPublish a steering input into an open run and observe whether the run's output considered it. Returns synchronously with two promises. Requires connect().
The first argument accepts a run id or a promise of one. Pass publishInput's runId promise and the transport publishes the steering message as soon as the agent's run start arrives, so a follow-up message typed before the run existed still reaches it. Both returned promises reject if the runId promise rejects.
Steering a run whose end this transport has already received rejects both promises without publishing.
1
2
3
4
const { published, outcome } = transport.steer(sent.runId, { kind: 'message', payload: followUp });
const { serial } = await published;
const { consumed, runTerminalReason } = await outcome;Parameters
runIdrequiredstring | Promise<string>eventrequiredTInputReturns
publishedPromise<{ serial: string | undefined }>{ serial } once the transport observes the steering message's own channel echo, and with serial undefined when that echo will not arrive: on close(), on channel continuity loss, or when the run's end arrives first. Rejects if the publish fails, or if the runId promise rejects.outcomeSteerOutcomeclose() and on channel continuity loss.Page channel history
history(opts?: TransportHistoryOptions): Promise<TransportHistoryResult<TInput, TOutput>>Page the channel's history backwards from the attach point and return the classified events as a batch. Requires connect().
Each call returns the next older slice and leaves the cursor paused, so repeated calls walk toward the start of the channel. History events are returned here only; they never reach subscribe handlers.
Decoding shares the live stream's decoder, so a stream that spans history and the live subscription is not decoded twice, and a single undecodable message is skipped and emitted on error. Concurrent calls serialise.
Parameters
optsoptionalTransportHistoryOptionsReturns
Promise<TransportHistoryResult<TInput, TOutput>>. Rejects with SessionHistoryFetchFailed when a page fetch fails after retries.
eventsTransportEventexhaustedBooleanClose the transport
close(): voidUnsubscribe the transport's listener from the channel and stop all delivery.
Terminal. Every subsequent connect, publishInput, cancel, steer, or history call rejects, in-flight steer outcomes reject, and existing subscriptions receive nothing further. The channel is caller-owned and is not detached.
Example
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
import * as Ably from 'ably';
import { createClientTransport } from '@ably/ai-transport';
import { createUIMessageCodec } from '@ably/ai-transport/vercel';
const codec = createUIMessageCodec();
const ably = new Ably.Realtime({ authUrl: '/api/auth/token' });
const transport = createClientTransport({
channel: ably.channels.get('conversation-42'),
codec,
clientId: 'alice',
});
await transport.connect();
let activeRunId;
transport.subscribe((event) => {
if (event.kind === 'message') fold(event.meta, event.inputs, event.outputs);
if (event.kind === 'run-lifecycle') {
activeRunId = event.event.type === 'start' ? event.event.runId : undefined;
}
});
transport.on('error', (error) => console.error(error.message));
// Read what happened before this client attached.
const { events } = await transport.history({ limit: 50 });
for (const event of events) fold(event.meta, event.inputs, event.outputs);
// Publish, then wake the agent yourself.
const sent = await transport.publishInput({
kind: 'message',
payload: { id: crypto.randomUUID(), role: 'user', parts: [{ type: 'text', text: 'Plan a trip to Lisbon.' }] },
});
await fetch('/api/chat', { method: 'POST', body: JSON.stringify({ eventId: sent.eventId }) });
// Steer it, without waiting to learn the run id first.
const { outcome } = transport.steer(sent.runId, {
kind: 'message',
payload: { id: crypto.randomUUID(), role: 'user', parts: [{ type: 'text', text: 'Make it five days.' }] },
});
const { consumed } = await outcome;
if (activeRunId) await transport.cancel(activeRunId);
transport.close();