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.

JavaScript

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

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

Returns

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

Subscribe 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

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, codec decode failures, a throw from a run's onCancel hook when that run set no onError, and a throw from a run's onSteer hook whether or not that run set one. Each is an ErrorInfo with a distinguishing code.

Parameters

eventrequired'event' | 'ably-message' | 'error'
Which stream to subscribe to.
handlerrequiredFunction
Called with each payload on that stream.

Returns

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

JavaScript

1

2

const trigger = await transport.locateInput(eventId);
const run = transport.openRun({ input: trigger });

Parameters

eventIdrequiredString
The event-id to match, which is the invocation's input event id.
optsoptionalTransportHistoryOptions
Scan bounds. limit 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.

metaWireMeta
The triggering input's transport-tier metadata. Pass the whole LocatedInput to openRun as input and it reads this metadata to select the opening event.
inputsTInput[]
The decoded input events the triggering wire message carried, in wire order.

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.

JavaScript

1

2

3

4

const run = transport.openRun(
  { input: trigger },
  { onCancel: async (request) => request.message.clientId === userId },
);

Parameters

optsoptionalOpenRunOptions
Run identity and structure.
hooksoptionalOpenRunHooks
Per-run callbacks and an external abort signal.

Returns

AgentRunTransport

A handle for publishing an open run's output and lifecycle.

runIdString
This run's id, generated at openRun or the reused continuation id. Available synchronously.
abortSignalAbortSignal
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.
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<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.

stepIdString
This step's id, stable across retry attempts of the same step.
pipeFunction
(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'
How the run ended.
erroroptionalErrorInfo
The terminal error to surface to clients, allowed only when 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.

JavaScript

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

runIdrequiredString
The id of the run to adopt. Must be non-empty.
optsoptionalAdoptRunOptions
Invocation-id pin.
hooksoptionalOpenRunHooks
Per-run callbacks and an external abort signal, the same shape openRun 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

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.

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

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