WireCodec

The WireCodec is the whole contract the transports require: a way to turn your events into Ably messages, and a way to turn Ably messages back into your events. It is the wire tier of the full Codec, so anything satisfying Codec satisfies WireCodec too, and a transport-only application supplies just this part of it.

Build one with defineCodec, which assembles the encoder and decoder from a declarative descriptor table. Pass the result to createClientTransport or createAgentTransport as codec.

JavaScript

1

2

3

4

5

import { createClientTransport } from '@ably/ai-transport';
import { createUIMessageCodec } from '@ably/ai-transport/vercel';

// The shipped codecs are full Codecs, which satisfy WireCodec.
const transport = createClientTransport({ channel, codec: createUIMessageCodec() });

Properties

adapterTagString
Optional. An Ably-Agent identifier registered on the channel, so traffic is attributed to this codec. Omit to opt out of registration.
createEncoderFunction
Create a stateful encoder bound to a channel. See Create an encoder.
createDecoderFunction
Create a stateful decoder for inbound messages. See Create a decoder.

Create an encoder

createEncoder(channel: ChannelWriter, options?: EncoderOptions): Encoder<TInput, TOutput>

Create a stateful encoder bound to one channel. The transports call this once per published message: one encoder for each client input, each steering message, and each pipe or send call, closed straight afterwards. You implement it, and only reach for the result directly when writing to a channel outside a transport.

Stream-tracker state lives inside the encoder and is shared across both directions, so one encoder cannot be reused across channels.

Parameters

channelrequiredChannelWriter
The channel the encoder publishes to.
optionsoptionalEncoderOptions
Encoder configuration, including the default extras and headers to set on every message.

Returns

publishInputFunction
(input, options?) => Promise<void>. Encode and publish one client input on the ai-input wire. Rejects if the codec cannot encode the given input variant.
publishOutputFunction
(output, options?) => Promise<void>. Encode and publish one agent output on the ai-output wire. Rejects if the codec cannot encode the given output variant.
cancelStreamsFunction
() => Promise<void>. Close every in-progress streamed message as status: cancelled and flush pending appends. Pure transport mechanics, emitting no codec output; run termination is signalled separately by ai-run-end. Idempotent, and throws if called after close.
closeFunction
() => Promise<void>. Flush pending appends and release encoder resources.

Create a decoder

createDecoder(): Decoder<TInput, TOutput>

Create a stateful decoder for one channel subscription.

The decoder keeps stream-tracker state across messages, so on a mid-stream join, whether from history compaction, a partial history page, or a rewind miss, the decoder synthesises the missing start events before it emits any delta. Your code therefore always sees a clean start, delta, end sequence.

The decoder's stream trackers are version-guarded: a delivery whose append version serial is at or below the version already incorporated decodes to nothing. One decoder instance is therefore safe to share between the live subscription and history paging, which is exactly what the transports do.

Returns

decodeFunction
(message: Ably.InboundMessage) => DecodedMessage. Decode one inbound Ably message into its input and output halves.

The shape decode returns is a

.

Build one with defineCodec

function defineCodec<TInput, TOutput>(): (config: DefineCodecConfig<TInput, TOutput>) => WireCodec<TInput, TOutput>

Assemble a codec from declarative descriptor tables rather than hand-writing createEncoder and createDecoder. defineCodec is curried on the input and output unions, so the descriptor callbacks narrow to each member without casts.

defineCodec produces a WireCodec and nothing more: there is no reducer slot and no message extraction, so the descriptor tables are the whole codec. That is all either transport needs, and the custom wire codec quickstart builds one end to end.

The codecs a durable session consumes are assembled by a different factory, which the SDK does not currently export. Until it does, a session takes one of the two pre-built session codecs: createUIMessageSessionCodec() or ResponsesSessionCodec.

Parameters

configrequiredDefineCodecConfig
The descriptor tables.
eventFunction
(type, spec?) => OutputDescriptor. Declare one discrete output event, published as a single message. The type literal is set as the wire kind dispatch header.
streamFunction
(kind, spec) => OutputDescriptor. Declare a streamed group of start, delta, and end chunks that become one durable message growing by append. The group id is set as the wire kind header on every phase.
dropFunction
(type) => OutputDescriptor. Declare an output type deliberately kept off the wire. The encoder publishes nothing for it, silently, so comment each entry with why the event is redundant.
eventFunction
(kind, spec?) => InputDescriptor. Declare a single-event input. fields and data read and write that member's payload; a member with no payload may only be declared wireOnly.
batchFunction
(kind, spec) => InputDescriptor. Declare a multi-part input: one domain message fanned out into one wire event per part, sharing the kind and codec-message-id, each carrying a partType. Use it for a user message with several parts.

Returns

WireCodec<TInput, TOutput>: adapterTag, createEncoder and createDecoder. It passes straight to either transport.

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

import { defineCodec, strField } from '@ably/ai-transport';

const fId = strField('id');
const fReason = strField('reason', '');
const fMessageId = strField('messageId', '');

export const myCodec = defineCodec()({
  adapterTag: 'my-provider',

  output: ({ event, stream, drop }) => [
    // One growing message per assistant reply.
    stream('text', {
      streamId: (chunk) => chunk.id,
      fields: [fId],
      start: { type: 'text-start' },
      delta: { type: 'text-delta', field: 'delta', decode: ({ rebuild }) => rebuild([fId]) },
      end: { type: 'text-end' },
    }),
    // One publish per occurrence.
    event('done', { fields: [fReason] }),
    // Deliberately off the wire: a keepalive nothing downstream reads.
    drop('ping'),
  ],

  input: ({ batch, event }) => [
    batch('user-message', {
      explode: (input) => [{ type: 'text', text: input.message.text }],
      partTypeOf: (part) => part.type,
      parts: (p) => [
        p('text', { data: { encode: (part) => part.text, decode: (d) => ({ text: String(d ?? '') }) } }),
      ],
      messageHeaders: (input) => ({
        codecHeaders: { messageId: input.message.id },
        transportHeaders: { role: input.message.role },
      }),
      assemble: (part, ctx) => ({
        message: { id: fMessageId.read(ctx.codecHeaders), role: 'user', text: part.text },
      }),
    }),
    event('regenerate', { wireOnly: true }),
  ],
});