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.

JavaScript

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

optionsrequiredClientTransportOptions
The transport's channel, codec, and identity.

Returns

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): () => void

Subscribe 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.

JavaScript

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

handlerrequiredTransportEvent
Called with each event in wire order.

Returns

() => void. An unsubscribe function.

Subscribe to a specific stream

on(event: 'event' | 'ably-message' | 'error', handler: (payload) => void): () => void

Subscribe to one of the transport's three receive streams. subscribe is shorthand for the 'event' form.

event
Classified TransportEvent values, fired once per inbound wire message that produces one.
ably-message
The raw inbound Ably.InboundMessage, emitted after the matching typed event so a handler sees state an earlier subscriber merged.
error
Channel and subscription failures, and codec decode failures, each an ErrorInfo 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'
Which stream to subscribe to.
handlerrequiredFunction
Called with each payload on that stream.

Returns

() => 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.

JavaScript

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

eventrequiredTInput
The codec input event to publish.
optsoptionalPublishInputOptions
Per-publish overrides.

Returns

Promise<PublishInputResult>. Resolves once the input is published, or rejects with an ErrorInfo.

codecMessageIdString
The codec-message-id the input was published under, either your option value or a freshly generated id. This keys optimistic-echo reconciliation.
eventIdString
The per-publish event-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>
Resolves with the run id when the transport observes the first run start whose 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

runIdrequiredString
The run to cancel.

Returns

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): SteerResult

Publish 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.

JavaScript

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>
The open run to steer, or a promise resolving to its id.
eventrequiredTInput
The codec input event to publish as the steering message.

Returns

publishedPromise<{ serial: string | undefined }>
Resolves with { 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.
outcomeSteerOutcome
Resolves at the run's next suspend when the run's output consumed the steering message, and at the run's end otherwise. Rejects on close() 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

optsoptionalTransportHistoryOptions
Batch bounds.

Returns

Promise<TransportHistoryResult<TInput, TOutput>>. Rejects with SessionHistoryFetchFailed when a page fetch fails after retries.

eventsTransportEvent
The classified events in this batch, oldest-first. Each batch is older than the one the previous call returned, so a consumer prepends.
exhaustedBoolean
True when the cursor reached the channel's attach point, meaning the whole history has been returned and further calls resolve with an empty batch.

Close the transport

close(): void

Unsubscribe 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

JavaScript

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();