# ClientTransport The `ClientTransport` publishes to and receives from one [Ably channel](https://ably.com/docs/channels.md), 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()`](#connect) and never detaches the channel. #### Javascript ``` 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(options: ClientTransportOptions): ClientTransport` 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()`](#connect) resolves. ### Parameters | Parameter | Required | Description | Type | | --- | --- | --- | --- | | options | required | The transport's channel, codec, and identity. |
|
| Property | Description | Type | | --- | --- | --- | | channel | The Ably channel to publish on and receive from. The transport subscribes its own listener on `connect()`; the channel stays caller-owned and is never detached. | `Ably.RealtimeChannel` | | codec | The wire tier of the codec: its encoder serializes inputs and its decoder classifies inbound messages. Any full `Codec` satisfies it. |
| | clientId | Optional. The publishing client's Ably `clientId`, set as the `run-client-id` header on inputs. Omit for an anonymous connection; the header is then absent and the local echo's `clientId` is `undefined`. | String | | historyPageSize | Optional. Wire messages fetched per `channel.history()` round trip in [`history()`](#history). Defaults to `100`. | Number | | logger | Optional. Logger for diagnostics. | `Logger` |
| Property | Description | Type | | --- | --- | --- | | adapterTag | Optional. Ably-Agent identifier registered on the channel. Omit to opt out of registration. | String | | createEncoder | `(channel, options?) => Encoder`. Create a stateful encoder bound to the channel. | Function | | createDecoder | `() => Decoder`. Create a stateful decoder converting inbound Ably messages into typed inputs and outputs. | Function |
### Returns `ClientTransport`. The transport, not yet connected. ## Connect the transport `connect(): Promise` 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`. Resolves once the channel is attached. A failure both rejects this call and is emitted on [`on('error')`](#on), as an [`ErrorInfo`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo). ## Subscribe to transport events `subscribe(handler: (event: TransportEvent) => void): () => void` Subscribe to classified transport events. Shorthand for [`on('event', handler)`](#on). Fires for live wire events and for the optimistic local echo that [`publishInput`](#publish-input) emits. History batches do not pass through here; they are returned from [`history()`](#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 ``` 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 | Parameter | Required | Description | Type | | --- | --- | --- | --- | | handler | required | Called with each event in wire order. |
|
| Property | Description | Type | | --- | --- | --- | | kind | Discriminator. `'message'` for a codec-decoded message, `'run-lifecycle'` for a run start, suspend, resume, or end, `'step-lifecycle'` for a step start or end. | `'message' \| 'run-lifecycle' \| 'step-lifecycle'` | | meta | Present when `kind` is `'message'`. The transport-tier metadata for the carrying wire message. |
| | inputs | Present when `kind` is `'message'`. The decoded client inputs from this message, in wire order. May be empty. | `TInput[]` | | outputs | Present when `kind` is `'message'`. The decoded agent outputs from this message, in wire order. May be empty. | `TOutput[]` | | event | Present when `kind` is `'run-lifecycle'` or `'step-lifecycle'`. The parsed lifecycle event. |
|
| Property | Description | Type | | --- | --- | --- | | codecMessageId | The logical message this event belongs to. Key accumulated content on it: many events across many wire messages share one id. `undefined` when the wire carried none. | String | | runId | The run this message was published under, or `undefined` for a run-less user input. | String | | stepId | The step attempt that published this output, or `undefined` when the message belonged to no step. Stable across retries of the same step. | String | | stepStartSerial | The identity of the step attempt, which is the serial of its `ai-step-start`. For a given `stepId` the highest value is the canonical attempt; earlier attempts are superseded. | String | | serial | Ably channel serial of the message, or `undefined` for an optimistic local echo. | String | | versionSerial | The append version serial, which advances per delivery on an appending stream. Use it to dedup a replayed wire message. `undefined` for an optimistic local echo. | String | | versionTimestamp | The append version timestamp, epoch milliseconds. `undefined` for an optimistic local event. | Number | | timestamp | Ably server timestamp in epoch milliseconds, or `undefined` for an optimistic local echo. | Number | | clientId | The publisher's Ably `clientId`, or `undefined` for an anonymous connection. | String | | role | The `role` header, for example `'user'` or `'assistant'`. `undefined` when the wire carried none. | String | | messageName | The Ably message name, for example `ai-input` or `ai-output`. `undefined` for an optimistic local echo. | String | | parent | Structure header that identifies the preceding message on this branch. Carried verbatim; the transport does not interpret it. | String | | forkOf | Structure header that identifies the message this one replaces. Carried verbatim. | String | | regenerates | Structure header that identifies the message this run regenerates. Carried verbatim. | String | | inputCodecMessageId | Structure header that identifies the input that triggered this run. Carried verbatim. | String | | inputCodecMessageIds | On a run end or suspend, the codec-message-ids of every input the run's output considered, the trigger and every steering message carried on the run's output. `undefined` elsewhere, when the run produced no output, or when the header is malformed. | `string[]` | | steerCodecMessageIds | The steering-message codec-message-ids the agent drained before the step attempt that produced this output. `undefined` when the wire carried no such header. | `string[]` | | transport | The complete `extras.ai.transport` header tier, verbatim. Empty object when the wire carried no transport tier. | `Record` | | codec | The complete `extras.ai.codec` header tier, verbatim. Empty object when the wire carried no codec tier. | `Record` | | headers | The application's own headers from Ably's `extras.headers` slot, outside the SDK's envelope. Empty object when the wire carried none. | `Record` |
| Property | Description | Type | | --- | --- | --- | | type | Which lifecycle transition this is. | `'start' \| 'suspend' \| 'resume' \| 'end'` | | runId | The run this event concerns. | String | | clientId | The owning client's Ably `clientId`. | String | | invocationId | The invocation this event was published under. Empty string when the wire carried none. | String | | serial | Ably channel serial of the lifecycle message, or `undefined` for an optimistic local event. | String | | timestamp | Ably server timestamp in epoch milliseconds. Absent for an optimistic local event. | Number | | reason | Present when `type` is `'end'`. How the run ended. | `'complete' \| 'cancelled' \| 'error'` | | error | Present when `type` is `'end'` and `reason` is `'error'`. The terminal error. | [`ErrorInfo`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo) | | parent | Present when `type` is `'start'`. The codec-message-id of the parent message. Omitted for a root run. | String | | forkOf | Present when `type` is `'start'`. The codec-message-id of the user prompt being forked, when the run is an edit. | String | | regenerates | Present when `type` is `'start'`. The codec-message-id of the assistant message this run regenerates. | String | | inputCodecMessageId | Present when `type` is `'start'`, and only when the publishing agent set `input-codec-message-id` on the opening event. A durable [`AgentSession`](https://ably.com/docs/ai-transport/api/javascript/core/agent-session.md) sets it from the triggering input it located. A run opened through the standalone `AgentTransport` never carries it on `ai-run-start`, so the field is absent there whether or not that agent knew the trigger. | String |
### 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. | Value | Description | | --- | --- | | `event` | Classified [`TransportEvent`](#subscribe) 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`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo) with a distinguishing `code`. A single decode failure drops that one message and emits the error rather than tearing down the stream. |
### Parameters | Parameter | Required | Description | Type | | --- | --- | --- | --- | | event | required | Which stream to subscribe to. | `'event' \| 'ably-message' \| 'error'` | | handler | required | Called with each payload on that stream. | Function |
### Returns `() => void`. An unsubscribe function. ## Publish a client input `publishInput(event: TInput, opts?: PublishInputOptions): Promise` Publish one codec input event to the channel. Requires [`connect()`](#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 ``` 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 | Parameter | Required | Description | Type | | --- | --- | --- | --- | | event | required | The codec input event to publish. | `TInput` | | opts | optional | Per-publish overrides. |
|
| Property | Description | Type | | --- | --- | --- | | codecMessageId | The codec-message-id to publish under. Defaults to a fresh id. Supplying one makes the input an amendment of an existing message, so the transport emits no local echo for it. | String | | parent | Structure: the codec-message-id of the preceding message on this branch. Omit for linear chat. | String | | forkOf | Structure: the codec-message-id this input replaces, which is an edit fork. Omit for linear chat. | String | | regenerates | Structure: the codec-message-id this input regenerates. Omit for linear chat. | String | | runId | Reuse a known run id, which makes this a continuation of an existing run. Omit for a fresh send; the agent generates the run id at run start. | String | | headers | Arbitrary application headers, published in Ably's own `extras.headers` slot outside the SDK's envelope and surfaced back on `meta.headers`, on the wire echo and the optimistic local echo alike. | `Record` |
### Returns `Promise`. Resolves once the input is published, or rejects with an [`ErrorInfo`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo). | Property | Description | Type | | --- | --- | --- | | codecMessageId | The codec-message-id the input was published under, either your option value or a freshly generated id. This keys optimistic-echo reconciliation. | String | | eventId | 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. | String | | runId | 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()`](#close) and on channel continuity loss; a rejection handler is pre-attached, so ignoring it raises nothing. | `Promise` |
## Cancel a run `cancel(runId: string): Promise` Publish a cancel envelope targeting the given run. Requires [`connect()`](#connect). Stateless: the transport holds no run registry, so supply a run id you captured from a `run-lifecycle` start event or from [`publishInput`](#publish-input)'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 | Parameter | Required | Description | Type | | --- | --- | --- | --- | | runId | required | The run to cancel. | String |
### Returns `Promise`. Resolves once the cancel envelope is published, or rejects with an [`ErrorInfo`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo). ## Steer an open run `steer(runId: string | Promise, 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()`](#connect). The first argument accepts a run id or a promise of one. Pass [`publishInput`](#publish-input)'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 ``` const { published, outcome } = transport.steer(sent.runId, { kind: 'message', payload: followUp }); const { serial } = await published; const { consumed, runTerminalReason } = await outcome; ``` ### Parameters | Parameter | Required | Description | Type | | --- | --- | --- | --- | | runId | required | The open run to steer, or a promise resolving to its id. | `string \| Promise` | | event | required | The codec input event to publish as the steering message. | `TInput` |
### Returns | Property | Description | Type | | --- | --- | --- | | published | 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()`](#close), on channel continuity loss, or when the run's end arrives first. Rejects if the publish fails, or if the `runId` promise rejects. | `Promise<{ serial: string \| undefined }>` | | outcome | 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()`](#close) and on channel continuity loss. |
|
| Property | Description | Type | | --- | --- | --- | | consumed | Whether the steering message's codec-message-id appeared in the `steer-codec-message-ids` headers observed on the run's output. `true` resolves at an end or a suspend alike; `false` resolves only at a run end, because a suspend leaves it pending for a later resume to consume. | Boolean | | runTerminalReason | How the run ended, present when a run end settled the outcome. Absent when a suspend settled it. | `'complete' \| 'cancelled' \| 'error'` |
## Page channel history `history(opts?: TransportHistoryOptions): Promise>` Page the channel's history backwards from the attach point and return the classified events as a batch. Requires [`connect()`](#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 | Parameter | Required | Description | Type | | --- | --- | --- | --- | | opts | optional | Batch bounds. |
|
| Property | Description | Type | | --- | --- | --- | | limit | Stop paging once at least this many events have been collected. Page granular: the call finishes the page it is on, so the batch may exceed the limit. Omit to page to channel exhaustion. | Number | | signal | Abort signal, checked between page fetches. When it fires the call rejects with `OperationCancelled`, and the cursor stays resumable so a later call continues from where this one stopped. | `AbortSignal` | | onPage | `() => void`. Called after each page fetch completes and before the next begins. A durable consumer uses it as a heartbeat while a long scan pages the channel. It observes progress only: the return value is ignored, and a throw is logged and the walk continues, so the callback cannot fail the call it is reporting on. | Function |
### Returns `Promise>`. Rejects with `SessionHistoryFetchFailed` when a page fetch fails after retries. | Property | Description | Type | | --- | --- | --- | | events | The classified events in this batch, oldest-first. Each batch is older than the one the previous call returned, so a consumer prepends. |
| | exhausted | 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. | Boolean |
## 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 ``` 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(); ``` ## Related Topics - [Agent transport](https://ably.com/docs/ai-transport/api/javascript/transport/agent-transport.md): API reference for the AI Transport AgentTransport: the factory, connect, openRun, locateInput, history, and the run and step handles that publish a run's output and lifecycle. - [Wire codec](https://ably.com/docs/ai-transport/api/javascript/transport/wire-codec.md): API reference for the AI Transport WireCodec: the encode and decode contract the client and agent transports require, and how defineCodec builds one from a descriptor table. ## Documentation Index To discover additional Ably documentation: 1. Fetch [llms.txt](https://ably.com/llms.txt) for the canonical list of available pages. 2. Identify relevant URLs from that index. 3. Fetch target pages as needed. Avoid using assumed or outdated documentation paths.