AgentTransport
The AgentTransport is the agent side of one Ably channel 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.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
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<TInput, TOutput>(options: AgentTransportOptions<TInput, TOutput>): AgentTransport<TInput, TOutput>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() resolves.
Parameters
optionsrequiredAgentTransportOptionsReturns
AgentTransport<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 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<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 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().
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.Parameters
eventrequired'event' | 'ably-message' | 'error'handlerrequiredFunctionReturns
() => void. An unsubscribe function.
Locate the input that woke this invocation
locateInput(eventId: string, opts?: TransportHistoryOptions): Promise<LocatedInput<TInput> | 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().
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.
1
2
const trigger = await transport.locateInput(eventId);
const run = transport.openRun({ input: trigger });Parameters
eventIdrequiredStringevent-id to match, which is the invocation's input event id.optsoptionalTransportHistoryOptionslimit caps the wire messages scanned, page granular, and a bounded scan that finds nothing resolves undefined.Returns
Promise<LocatedInput<TInput> | undefined>. Resolves undefined when no matching input is found in history.
metaWireMetaLocatedInput to openRun as input and it reads this metadata to select the opening event.inputsTInput[]Open a run
openRun(opts?: OpenRunOptions, hooks?: OpenRunHooks<TOutput>): AgentRunTransport<TOutput>Open a run and return a handle that publishes its output and lifecycle. Requires connect().
Pass the located 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.
1
2
3
4
const run = transport.openRun(
{ input: trigger },
{ onCancel: async (request) => request.message.clientId === userId },
);Parameters
optsoptionalOpenRunOptionshooksoptionalOpenRunHooksReturns
AgentRunTransport
A handle for publishing an open run's output and lifecycle.
runIdStringopenRun or the reused continuation id. Available synchronously.abortSignalAbortSignalsignal 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.Check for pending input
hasInput(): booleanWhether 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<TOutput>): Promise<StreamResult>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<TOutput>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.
stepIdStringpipeFunction(source) => Promise<StreamResult>. Pipe an output stream; every output carries this step's identity.sendFunction(event) => Promise<void>. Publish a single discrete output as one assistant message carrying this step's identity. The step must be active.endFunction(params) => Promise<void>. Publish ai-step-end, closing the step. Idempotent. The reason is derived from piped output when omitted.Suspend the run
suspend(): Promise<void>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<void>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<void>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.
reasonrequired'complete' | 'cancelled' | 'error'reason is 'error'. Omit to end in error without detail.Adopt an open run
adoptRun(runId: string, opts?: AdoptRunOptions, hooks?: OpenRunHooks<TOutput>): AgentRunTransport<TOutput>Adopt a run that is already open, without publishing anything. The handle registers for cancel and steer routing exactly as openRun does, and nothing reaches the channel until you publish output or a terminal event. Requires 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, 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.
1
2
3
4
5
6
7
8
9
10
11
12
13
// 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
runIdrequiredStringoptsoptionalAdoptRunOptionshooksoptionalOpenRunHooksopenRun takes.Returns
AgentRunTransport<TOutput>, the same handle openRun returns. Its runId is the adopted id, available synchronously.
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, so the agent can assemble prior conversation context for a model call. Requires 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's throwaway scan stays separate. 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 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
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
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();
}
}