# AgentTransport
The `AgentTransport` is the agent side of one [Ably channel](https://ably.com/docs/channels.md) and one codec. It opens runs, finds the input that woke an invocation, pages history for model context, and observes the channel so a cancel signal routes to the matching run handle and a steering message makes that run's `hasInput()` return `true`.
Construct one with `createAgentTransport` from the core entry point. It exposes no conversation-tree accessors. Everything it reads comes back as classified events, which the agent merges into whatever state it keeps.
#### Javascript
```
import * as Ably from 'ably';
import { createAgentTransport } from '@ably/ai-transport';
import { createUIMessageCodec } from '@ably/ai-transport/vercel';
// Agent-side. This runs on your own infrastructure, so it holds the key.
const ably = new Ably.Realtime({ key: process.env.ABLY_API_KEY });
const transport = createAgentTransport({
channel: ably.channels.get('conversation-42'),
codec: createUIMessageCodec(),
clientId: 'agent',
});
await transport.connect();
```
## Create an agent transport
`function createAgentTransport(options: AgentTransportOptions): AgentTransport`
Construct a transport bound to a channel and a codec. Construction is synchronous and passive: nothing is subscribed, no cancel signal is routed, and no run can open until [`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 run and step lifecycle and output on, and to receive cancel and steering signals 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 output and its decoder classifies the receive stream, `locateInput`, and `history`. Any full `Codec` satisfies it. | |
| clientId | Optional. The agent's Ably `clientId`, set as the `run-client-id` header on the run's lifecycle and output. The header carries an empty string when omitted. | String |
| historyPageSize | Optional. Wire messages fetched per channel-history page in [`locateInput`](#locate-input) and [`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
`AgentTransport`. 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 and cancel routing.
`history` and `locateInput` require a successful `connect()` first. `openRun` and `adoptRun` need only that you called it, and a failed connect surfaces when you await the run's opening publish or its first output. A run opened without it could silently miss the cancel signals addressed to it, which is why the guard is unconditional rather than advisory.
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 step-lifecycle events a run's output methods emit. History batches do not pass through here.
Steering messages arrive here like any other input, so an agent that assembles its own context reads them from this stream. Their effect on the run's loop is separate, on [`hasInput()`](#open-run).
### 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. | `RunLifecycleEvent \| StepLifecycleEvent` |
### 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, codec decode failures, a throw from a run's `onCancel` hook when that run set no [`onError`](#open-run-params), and a throw from a run's `onSteer` hook whether or not that run set one. Each is an [`ErrorInfo`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo) with a distinguishing `code`. |
### 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.
## Locate the input that woke this invocation
`locateInput(eventId: string, opts?: TransportHistoryOptions): Promise | undefined>`
Scan channel history for the input event whose `event-id` header matches `eventId`, so a fresh process can resume durably without the conversation tree. Requires [`connect()`](#connect).
This is the one history read the SDK performs on your behalf. It runs on a throwaway decoder, so it never perturbs the live receive stream's dedup state, and the triggering input is the newest message in that history, so it reads about one page whatever the conversation's length.
### Javascript
```
const trigger = await transport.locateInput(eventId);
const run = transport.openRun({ input: trigger });
```
### Parameters
### Returns
`Promise | undefined>`. Resolves `undefined` when no matching input is found in history.
## Open a run
`openRun(opts?: OpenRunOptions, hooks?: OpenRunHooks): AgentRunTransport`
Open a run and return a handle that publishes its output and lifecycle. Requires [`connect()`](#connect).
Pass the [located input](#locate-input) as `opts.input` and it determines how the run opens. The input's own `meta.runId`, the run-id header a client sets on a continuation, selects the opening event: present means this call re-enters that run and publishes `ai-run-resume`, absent means a fresh run and `ai-run-start`. The input also defaults `inputCodecMessageId`, and for a fresh run it defaults the structure options, so the agent never decides start-versus-resume itself.
The run is registered for cancel routing until it ends, and a cancel already buffered against the input's codec-message-id is honoured immediately.
### Javascript
```
const run = transport.openRun(
{ input: trigger },
{ onCancel: async (request) => request.message.clientId === userId },
);
```
### Parameters
| Parameter | Required | Description | Type |
| --- | --- | --- | --- |
| opts | optional | Run identity and structure. | |
| hooks | optional | Per-run callbacks and an external abort signal. | |
| Property | Description | Type |
| --- | --- | --- |
| input | The located input that woke this invocation, from [`locateInput`](#locate-input). Its `meta.runId` selects the opening event, so the caller never decides start-versus-resume. Its `meta.codecMessageId` defaults `inputCodecMessageId`, and for a fresh run its `meta.parent`, `meta.forkOf` and `meta.regenerates` default the structure options below. An option supplied explicitly wins over the input's value. | |
| runId | Without an `input`, supplying it marks this open as a continuation of that run. With an `input`, it pins the id of a fresh run only, which is how a durable agent makes a fresh-process retry re-enter the same run; a continuation's id comes from the input. | String |
| invocationId | Reuse a fixed invocation id. Omit to generate a fresh one, which is one per HTTP request. | String |
| parent | Structure: the codec-message-id of the parent message. Defaulted from `input` for a fresh run. Omit for a root run. | String |
| forkOf | Structure: the codec-message-id being forked, which is an edit. Defaulted from `input` for a fresh run. Omit unless forking. | String |
| regenerates | Structure: the codec-message-id this run regenerates. Defaulted from `input` for a fresh run. Omit unless regenerating. | String |
| inputCodecMessageId | The triggering input's codec-message-id. Defaulted from `input`, so supply it only when opening without one. It lets a cancel message keyed by input, which the client sends before it learns the run id, route to this run, including a cancel that arrived before this call, and it sets the `input-codec-message-id` header on the run's outputs and on its `ai-run-start`, and the client reads the `ai-run-start` header to resolve its `runId` promise. Without it, only a cancel message that carries the run id reaches this run. | String |
| Property | Description | Type |
| --- | --- | --- |
| signal | An external abort signal composed with the run's own. Aborting it fires `abortSignal` exactly as a cancel would. | `AbortSignal` |
| onCancel | `(request: CancelRequest) => Promise`. Authorise an incoming cancel. Return `false` to reject it and the run continues. Every cancel is accepted when omitted. | Function |
| onCancelled | `(write) => void \| Promise`. Runs when the abort fires, before in-flight streams close, so the agent can publish a final output on the way out. | Function |
| onSteer | `() => void`. Fires once per live steering message tracked for this run, and once for the whole batch of steering messages that arrived before this `openRun`, so the agent can race the arrival against an in-flight model call. A hint rather than the authority: `hasInput()` decides whether the loop runs again. Keep it synchronous and cheap, because it runs on the message-handling path. | Function |
| onError | `(error: Ably.ErrorInfo) => void`. Non-fatal run-scoped errors that have no other delivery path: a stream failure in `pipe` (`RunResponseStreamFailed`, also returned on `StreamResult.error`), a throw from `onCancel` (`RunCancelHandlerFailed`, and the run is not cancelled), and a throw from `onSteer` (`RunSteerHandlerFailed`, and the run is unaffected). An `onCancel` throw goes here when this is set, and to the transport's [`error`](#on) stream when it is not, never to both. A `pipe` stream failure calls this hook when you set one, and with no hook the failure is available only on the returned `StreamResult.error`. An `onSteer` throw always goes to the transport's `error` stream, even when this is set. Publish failures in the opening publish and in `end` reject their own promise instead. | Function |
| onAblyMessage | `(message: Ably.Message) => void`. Mutate the run's output messages before they are published. Run and step lifecycle messages publish straight to the channel, so they do not pass through this hook. | Function |
### Returns
#### AgentRunTransport
A handle for publishing an open run's output and lifecycle.
| Property | Description | Type |
| --- | --- | --- |
| runId | This run's id, generated at `openRun` or the reused continuation id. Available synchronously. | String |
| abortSignal | Fires when an accepted cancel routes to this run, or when the external `signal` aborts. An in-flight `pipe` call ends `'cancelled'` on its own; your agent watches `abortSignal` to abort its own work and then publishes `ai-run-end` itself. | `AbortSignal` |
##### Check for pending input
`hasInput(): boolean`
Whether the run has input awaiting a response. The agent's loop calls it to decide whether to iterate again.
Returns `true` until the run's first `pipe` or `send` call, which answers the triggering input, then `true` if a steering message has been tracked since the previous `hasInput()` call.
Reading it **drains** pending steering messages into the set the next step attempt publishes as its `steer-codec-message-ids` header, so there is no observe-only check. Call it once per loop iteration, immediately before assembling that iteration's context. It returns `false` once `abortSignal` has fired.
##### Pipe an output stream
`pipe(source: PipeSource): Promise`
Pipe an output stream through the encoder to the channel, as one implicit step whose start and end the SDK publishes for you. The SDK opens that step at the first output chunk, so a stream that produces nothing publishes no `ai-step-start` and no `ai-step-end`. Accepts a `ReadableStream` or any async iterable of the codec's output events. Resolves with `{ reason, error? }` when the stream completes, is cancelled, or errors.
##### Create a step
`createStep(opts?: StepOptions): RunStepTransport`
Create a re-attemptable unit of work within this run, whose retries supersede the prior attempt's output. Creating the handle is synchronous.
Pass an explicit `stepId` when the same logical step re-attempts in a separate process, which is a durable-execution retry, using the framework's own stable step identity. Omit it and the SDK generates an id scoped to this invocation, so an in-process retry reuses the same step id and a cross-process retry appends beside the failed attempt rather than superseding it.
`stepClientId` attributes the step to a participant, and you usually omit it. The SDK then inherits the `stepClientId` of the prior step in this invocation, and for a run's first step sets it to the empty string. Supply it when a steering message brings in a fresh input mid-run and the step should be attributed to that input's publisher rather than inheriting the prior step's value.
| Property | Description | Type |
| --- | --- | --- |
| stepId | This step's id, stable across retry attempts of the same step. | String |
| pipe | `(source) => Promise`. Pipe an output stream; every output carries this step's identity. | Function |
| send | `(event) => Promise`. Publish a single discrete output as one assistant message carrying this step's identity. The step must be active. | Function |
| end | `(params) => Promise`. Publish `ai-step-end`, closing the step. Idempotent. The reason is derived from piped output when omitted. | Function |
##### Suspend the run
`suspend(): Promise`
Publish `ai-run-suspend`, pausing the run without ending it. The handle blocks output while suspended, `resume()` re-opens it, and a later invocation can continue the run under the same id through a fresh `openRun`.
Throws while a step is active, so end the step first. No-op once suspended or ended. When the run has produced output, the `ai-run-suspend` event carries the `input-codec-message-ids` header, which lists every input the run considered so far.
##### Resume the run
`resume(): Promise`
Publish `ai-run-resume`, re-entering the run. It is a pure re-entry signal carrying no structure headers. A suspended handle accepts output again once the publish succeeds. Throws once the run has ended.
##### End the run
`end(params: RunEndParams): Promise`
Publish `ai-run-end`, ending the run. Terminal, and it auto-closes a still-open step first, so every client reading that step sees it end.
When the run has produced output, the `ai-run-end` event carries the `input-codec-message-ids` header, which lists every input the run considered, accumulated across suspend and resume. A client can check whether its steering message was processed by looking for its id in that header.
| Parameter | Required | Description | Type |
| --- | --- | --- | --- |
| reason | required | How the run ended. | `'complete' \| 'cancelled' \| 'error'` |
| error | optional | The terminal error to surface to clients, allowed only when `reason` is `'error'`. Omit to end in error without detail. | [`ErrorInfo`](https://ably.com/docs/ai-transport/api/errors.md#errorinfo) |
## Adopt an open run
`adoptRun(runId: string, opts?: AdoptRunOptions, hooks?: OpenRunHooks): AgentRunTransport`
Adopt a run that is already open, without publishing anything. The handle registers for cancel and steer routing exactly as [`openRun`](#open-run) does, and nothing reaches the channel until you publish output or a terminal event. Requires [`connect()`](#connect).
This is how a durable activity re-enters a run it has to gate or clean up. The transport holds no run state, so deciding whether there is anything left to do is yours: read [`history`](#history), then act.
An adopted run publishes no opening event, which is why the structure options stay on `openRun`; nothing here would publish them. To re-enter a run and continue it, call `openRun` with that run's id rather than `adoptRun`.
### Javascript
```
// Agent code, in a durable cleanup activity.
const run = transport.adoptRun(runId, { invocationId });
// The transport holds no run state, so the gate is yours: read history.
const { events } = await transport.history({ limit: 200 });
const ended = events.some(
(event) => event.kind === 'run-lifecycle' && event.event.type === 'end' && event.event.runId === runId,
);
// Nothing has reached the channel yet. Publish only if there is something to say.
if (!ended) {
await run.end({ reason: 'error' });
}
```
### Parameters
| Parameter | Required | Description | Type |
| --- | --- | --- | --- |
| runId | required | The id of the run to adopt. Must be non-empty. | String |
| opts | optional | Invocation-id pin. | |
| hooks | optional | Per-run callbacks and an external abort signal, the same shape `openRun` takes. | |
| Property | Description | Type |
| --- | --- | --- |
| invocationId | Reuse a fixed invocation id, so a continuing process publishes under the invocation that owns the run. This matters because a tree consumer discards a suspend whose invocation id differs from the run's latest resume, treating it as a retired invocation's late suspend. Omit to generate a fresh one. | String |
### Returns
[`AgentRunTransport`](#AgentRunTransport), the same handle `openRun` returns. Its `runId` is the adopted id, available synchronously.
## 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, so the agent can assemble prior conversation context for a model call. Requires [`connect()`](#connect).
Each call returns the next older slice and leaves the cursor paused. 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, while [`locateInput`](#locate-input)'s throwaway scan stays separate. 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. | `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. | Boolean |
## Close the transport
`close(): void`
Unsubscribe the transport's listener from the channel and stop all delivery and cancel routing.
Terminal: after close, `connect`, `history` and `locateInput` reject, `openRun` and `adoptRun` throw, a second `close()` does nothing, and `subscribe` and `on` still register handlers that never fire. Open runs are **not** aborted or ended, and their write handles simply stop receiving signals, so end or suspend a run before closing or it stays open on the channel. The channel is caller-owned and is not detached.
## Example
### Javascript
```
import * as Ably from 'ably';
import { createAgentTransport } from '@ably/ai-transport';
import { createUIMessageCodec } from '@ably/ai-transport/vercel';
// Agent-side. This runs on your own infrastructure, so it holds the key.
const ably = new Ably.Realtime({ key: process.env.ABLY_API_KEY });
export async function handleInvocation({ conversationId, eventId }) {
const transport = createAgentTransport({
channel: ably.channels.get(conversationId),
codec: createUIMessageCodec(),
clientId: 'agent',
});
await transport.connect();
try {
const trigger = await transport.locateInput(eventId);
const run = transport.openRun({ input: trigger });
// Steering messages arrive as ordinary events carrying this run's id.
const pending = [];
transport.subscribe((event) => {
if (event.kind === 'message' && event.meta.runId === run.runId) pending.push(...event.inputs);
});
let context = [...(await loadConversation(conversationId)), ...(trigger?.inputs ?? [])];
do {
// hasInput() drains pending steering messages, so call it once per iteration,
// immediately before assembling this iteration's context.
context = [...context, ...pending.splice(0)];
const result = streamText({ model, messages: context, abortSignal: run.abortSignal });
await run.pipe(result.toUIMessageStream());
} while (run.hasInput());
await run.end({ reason: run.abortSignal.aborted ? 'cancelled' : 'complete' });
} finally {
transport.close();
}
}
```
## Related Topics
- [Client transport](https://ably.com/docs/ai-transport/api/javascript/transport/client-transport.md): API reference for the AI Transport ClientTransport: the factory, connect, subscribe, publishInput, cancel, steer, history, and the classified transport event stream.
- [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.