TanStack
Transports

Custom Transports

No built-in transport fits your protocol. Write your own connection adapter. There are two shapes:

Request-scoped adapters

If your transport makes one request per user message, implement ConnectConnectionAdapter. This is the lowest-level option before a persistent transport:

ts
import { useChat, type ConnectConnectionAdapter } from "@tanstack/ai-react";
import type { StreamChunk } from "@tanstack/ai";

const myAdapter: ConnectConnectionAdapter = {
  async *connect(messages, data, abortSignal, runContext) {
    const response = await fetch("/api/chat", {
      method: "POST",
      headers: {
        "Content-Type": "application/json",
        ...runContext?.headers,
      },
      body: JSON.stringify({
        threadId: runContext?.threadId,
        runId: runContext?.runId,
        messages,
        ...data,
      }),
      ...(abortSignal ? { signal: abortSignal } : {}),
    });

    if (!response.ok) throw new Error(`HTTP ${response.status}`);
    if (!response.body) throw new Error("Response has no body");

    // Example: newline-delimited JSON. Replace this loop with whatever
    // framing your wire format uses, yielding one `StreamChunk` per event.
    const reader = response.body.getReader();
    const decoder = new TextDecoder();
    let buffer = "";
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;
      buffer += decoder.decode(value, { stream: true });
      const lines = buffer.split("\n");
      buffer = lines.pop() ?? "";
      for (const line of lines) {
        if (line.trim()) {
          const chunk: StreamChunk = JSON.parse(line);
          yield chunk;
        }
      }
    }
  },
};

const { messages } = useChat({ connection: myAdapter });

runContext has these fields:

  • threadId, runId, clientTools, forwardedProps: put these in your JSON payload, so that the server can build an AG-UI-compliant response.
  • headers: copy these onto the POST request headers, not into the body.

The built-in fetch and XHR adapters copy runContext.headers. stream() and rpcStream() do not, because they never get runContext. If your custom connect does not copy headers, it drops BYOK keys (x-byok-*). See Bring Your Own Key.

The runtime adds the terminal event for you:

  • If your connect stream completes with no RUN_FINISHED, the runtime adds one.
  • If your connect stream throws, the runtime adds a RUN_ERROR.

Persistent transports

A WebSocket, a BroadcastChannel, postMessage between iframes, or a shared worker works in a different way from request and response. You open the channel once. Then you send and receive on it for the full life of the client. connect cannot model this, because it expects one async iterable per request.

For a resumable WebSocket, use the built-in webSocket() adapter. For other persistent channels, implement SubscribeConnectionAdapter (full definition in The adapter interface):

ts
import type { SubscribeConnectionAdapter } from "@tanstack/ai-react";

// subscribe(abortSignal?): AsyncIterable<StreamChunk>   (long-lived)
// send(messages, data?, abortSignal?, runContext?): Promise<void>   (one per user message)
  • subscribe(): the ChatClient calls it once. It returns a long-lived async iterable of every chunk that the channel produces.
  • send(): the ChatClient calls it once per user message, to put a request frame on the channel. It returns when the frame is written. The chunks arrive separately through subscribe().

The runtime connects the two. The chunks that arrive between send() and the next terminal event (RUN_FINISHED or RUN_ERROR) belong to that run.

Custom WebSocket example

If webSocket() does not fit (a different wire format, no need for resume, or a server that you do not control), write the adapter yourself:

ts
import { useChat, type SubscribeConnectionAdapter } from "@tanstack/ai-react";
import type { StreamChunk } from "@tanstack/ai";

function websocketConnection(url: string): SubscribeConnectionAdapter {
  const ws = new WebSocket(url);
  const queue: Array<StreamChunk> = [];
  let pending: ((chunk: StreamChunk | null) => void) | null = null;
  let closed = false;

  const ready = new Promise<void>((resolve) => {
    ws.addEventListener("open", () => resolve(), { once: true });
  });

  function deliver(chunk: StreamChunk | null) {
    const resolve = pending;
    if (resolve) {
      pending = null;
      resolve(chunk);
    } else if (chunk !== null) {
      queue.push(chunk);
    }
  }

  ws.addEventListener("message", (event) => {
    const chunk: StreamChunk = JSON.parse(event.data);
    deliver(chunk);
  });
  ws.addEventListener("close", () => {
    closed = true;
    deliver(null);
  });

  return {
    async *subscribe(abortSignal) {
      // Register the abort listener once (not per-iteration) so it can't
      // accumulate on a long-lived socket.
      const onAbort = () => deliver(null);
      abortSignal?.addEventListener("abort", onAbort, { once: true });
      try {
        while (!abortSignal?.aborted) {
          // Drain buffered chunks BEFORE honoring `closed`: a burst of messages
          // followed by a close event (common within one macrotask) must still
          // deliver the queued chunks (including a trailing RUN_FINISHED),
          // otherwise the client would hang waiting for a terminal it dropped.
          const buffered = queue.shift();
          if (buffered !== undefined) {
            yield buffered;
            continue;
          }
          if (closed) return;
          const chunk = await new Promise<StreamChunk | null>((resolve) => {
            pending = resolve;
          });
          if (chunk === null) return;
          yield chunk;
        }
      } finally {
        abortSignal?.removeEventListener("abort", onAbort);
      }
    },

    async send(messages, data, _abortSignal, runContext) {
      await ready;
      ws.send(
        JSON.stringify({
          threadId: runContext?.threadId,
          runId: runContext?.runId,
          messages,
          data,
        }),
      );
    },
  };
}

const { messages } = useChat({
  connection: websocketConnection("wss://example.com/chat"),
});

Tip: Your server must emit RUN_FINISHED (or RUN_ERROR) at the end of each run. Without it, the client does not know that the assistant turn ended, and it waits forever. See Stream Events for the full event lifecycle.

When to use a persistent transport

Use subscribe / send if one or more of these is true:

  • One connection carries many runs (the chat thread keeps the socket open across messages).
  • The server pushes chunks outside of a request (presence updates, server-initiated tool calls, broadcast notifications).
  • You want to share one connection across tabs (BroadcastChannel) or workers.

In other cases, use fetchServerSentEvents or stream(). They are simpler, and you do not manage a connection lifecycle.

The adapter interface

A ConnectionAdapter is a union. Give it connect, or give it both subscribe and send. Never give it both modes.

ts
import type { UIMessage } from "@tanstack/ai-client";
import type { ModelMessage, StreamChunk } from "@tanstack/ai";

export interface RunAgentInputContext {
  threadId: string;
  runId: string;
  parentRunId?: string;
  clientTools?: Array<{ name: string; description: string; parameters: unknown }>;
  forwardedProps?: Record<string, unknown>;
  headers?: Record<string, string>;
}

export interface ConnectConnectionAdapter {
  connect(
    messages: UIMessage[] | ModelMessage[],
    data?: Record<string, any>,
    abortSignal?: AbortSignal,
    runContext?: RunAgentInputContext,
  ): AsyncIterable<StreamChunk>;
}

export interface SubscribeConnectionAdapter {
  subscribe(abortSignal?: AbortSignal): AsyncIterable<StreamChunk>;
  send(
    messages: UIMessage[] | ModelMessage[],
    data?: Record<string, any>,
    abortSignal?: AbortSignal,
    runContext?: RunAgentInputContext,
  ): Promise<void>;
}

export type ConnectionAdapter =
  | ConnectConnectionAdapter
  | SubscribeConnectionAdapter;

Internally, ChatClient changes both shapes to one subscribe / send pair with normalizeConnectionAdapter():

  • connect: the client wraps it in an async queue. The wrapped send() waits until the active subscriber processes all events or exits.
  • subscribe + send: the client uses them as they are.

Pass your adapter as connection and send a message. The chunks from your transport stream into messages like any built-in adapter.