`WorkflowChatTransport`

WorkflowChatTransport

워크플로 기반 채팅 앱에서 스트림 자동 재연결을 지원하는 useChat용 ChatTransport 구현이에요. 메시지를 채팅 엔드포인트로 POST하고 x-workflow-run-id 응답 헤더를 추출한 뒤, 중단(네트워크 오류, 페이지 새로고침, 함수 타임아웃) 시 /{runId}/stream 엔드포인트로 재연결해요.

전체 응답이 단일 HTTP 요청으로 도착한다고 가정하는 DefaultChatTransport와 달리, WorkflowChatTransport는 초기 응답 스트림이 함수 타임아웃으로 중단될 수 있는 Workflow SDK를 위해 설계됐어요. 이 트랜스포트는 finish 이벤트가 누락되면 자동으로 감지해, 중단된 지점부터 이어받도록 재연결해요.

'use client';

import { useChat } from '@ai-sdk/react';
import { WorkflowChatTransport } from '@ai-sdk/workflow/client';

export default function Chat() {
  const { messages, sendMessage } = useChat({
    transport: new WorkflowChatTransport({
      api: '/api/chat',
      maxConsecutiveErrors: 5,
    }),
  });

  // ... render chat UI
}

출처: 문서

본문

Import

import { WorkflowChatTransport } from "@ai-sdk/workflow/client"

생성자 (Constructor)

파라미터 (Parameters)

이름 타입 선택 설명
api string 선택 채팅 요청용 API 엔드포인트. 재연결 엔드포인트는 {api}/{runId}/stream으로 유도돼요. 기본값: '/api/chat'.
fetch typeof fetch 선택 HTTP 요청에 사용할 커스텀 fetch 구현. 기본값: 전역 fetch.
maxConsecutiveErrors number 선택 포기하기 전 재연결 시도 중 허용되는 연속 오류 최대 횟수. 기본값: 3.
initialStartIndex number 선택 재연결 시 시작할 기본 청크 인덱스. 음수 값은 내구성(durable) UIMessageChunk 스트림의 끝에서부터 읽어요(예: -50은 마지막 50개 청크를 가져옴). 페이지 새로고침 후 전체 대화를 재생하지 않고 이어가기에 유용해요. 원시 ModelCallStreamPart 스트림은 음수 UI 청크 인덱스를 지원하지 않아요. reconnectToStream 옵션으로 호출마다 재정의할 수 있어요. 기본값: 0.
onChatSendMessage (response: Response, options: SendMessagesOptions) => void | Promise<void> 선택 초기 POST 요청이 성공한 후 호출되는 콜백. 응답 헤더 확인(예: 워크플로 실행 ID 추출)이나 클라이언트 측 채팅 기록 추적에 유용해요.
onChatEnd ({ chatId, chunkIndex }) => void | Promise<void> 선택 스트림이 끝날 때(finish 청크 수신) 호출되는 콜백. 채팅 ID와 총 청크 수를 받아요. 정리 작업이나 상태 업데이트에 유용해요.
prepareSendMessagesRequest PrepareSendMessagesRequest 선택 전송 전 POST 요청을 커스터마이즈하는 함수. API 엔드포인트, 헤더, 자격 증명, 본문을 재정의할 수 있어요.
prepareReconnectToStreamRequest PrepareReconnectToStreamRequest 선택 재연결 GET 요청을 커스터마이즈하는 함수. API 엔드포인트, 헤더, 자격 증명을 재정의할 수 있어요.

메서드 (Methods)

sendMessages()

메시지를 채팅 엔드포인트로 POST로 보내고 스트리밍 응답을 반환해요. 스트림이 중단되면(finish 이벤트 미수신), GET으로 {api}/{runId}/stream?startIndex={chunkIndex}에 자동 재연결해 중단된 지점부터 이어받아요.

POST 요청은 메시지를 JSON으로 포함하며, 응답에 워크플로 실행을 식별하는 x-workflow-run-id 헤더가 포함될 것으로 기대해요.

const stream = await transport.sendMessages({
  chatId: 'chat-123',
  trigger: 'submit-message',
  messages: [...],
  abortSignal: controller.signal,
});
이름 타입 설명
chatId string 채팅 세션의 고유 식별자.
trigger "'submit-message' | 'regenerate-message'" 메시지 제출 유형.
messageId string | undefined 재생성할 메시지의 ID. 새 메시지면 undefined.
messages UIMessage[] 대화 기록을 나타내는 UI 메시지 배열.
abortSignal AbortSignal | undefined 요청을 중단하는 신호. 초기 POST와 재연결 GET 요청 모두에 전파돼요.
반환 (Returns)

초기 POST 응답과 자동 재연결의 청크를 모두 포함하는 Promise<ReadableStream<UIMessageChunk>>를 반환해요.

reconnectToStream()

이전에 중단된 기존 채팅 스트림에 재연결해요. 페이지 새로고침 후 이어가거나 클라이언트가 연결을 다시 맺어야 할 때 유용해요.

const stream = await transport.reconnectToStream({
  chatId: 'chat-123',
  startIndex: -50, // Optional: fetch last 50 chunks
});
이름 타입 선택 설명
chatId string 재연결할 채팅 ID. 재연결 URL을 구성하는 데 사용돼요.
abortSignal AbortSignal | undefined 재연결 요청을 중단하는 신호.
startIndex number 선택 이 재연결의 시작 인덱스를 재정의해요. 서버의 내구성 스트림과 tail-index 헤더가 동일한 UIMessageChunk 인덱스 공간을 사용할 때 음수 값은 끝에서부터 읽어요. 생략하면 생성자의 initialStartIndex를 사용해요.
반환 (Returns)

Promise<ReadableStream<UIMessageChunk> | null>을 반환해요.

재연결 동작 방식 (How Reconnection Works)

트랜스포트는 다음 흐름을 따라요:

  1. POST {api}로 메시지 전송. 응답에 x-workflow-run-id 헤더가 포함되어야 해요.
  2. 스트리밍: SSE 응답을 스트리밍하며 청크가 도착할 때마다 카운트해요.
  3. 중단 감지: 스트림이 finish 이벤트 없이 닫히면(예: 함수 타임아웃, 네트워크 오류), 응답이 불완전함을 감지해요.
  4. 재연결: GET으로 {api}/{runId}/stream?startIndex={chunkIndex}에 재연결해 마지막으로 받은 청크부터 이어받아요.
  5. 재시도: 재연결 스트림도 중단되면 maxConsecutiveErrors 횟수까지 재시도해요.
  6. 완료: finish 이벤트를 받으면 onChatEnd를 호출하고 스트림을 닫아요.

음수 시작 인덱스 (Negative Start Index)

initialStartIndex가 음수(예: -50)면, 트랜스포트는 첫 재연결 요청에 그대로 보내요. 서버는 이를 절대 위치로 해석하고 x-workflow-stream-tail-index 응답 헤더를 반환해서, 트랜스포트가 이후 재시도의 올바른 위치를 계산할 수 있게 해야 해요.

헤더가 없거나 유효하지 않으면, 트랜스포트는 처음(startIndex=0)부터 재생하는 걸로 폴백해요.

음수 인덱스는 저장된 객체가 이미 UIMessageChunk인 내구성 서버 스트림이 필요해요. 아래 보이는 원시 WorkflowAgent 변환은 음수가 아닌 인덱스만 지원해요.

서버 요구 사항 (Server Requirements)

WorkflowChatTransport가 동작하려면 서버가 두 엔드포인트를 제공해야 해요:

POST {api} (예: /api/chat)

  • 메시지를 JSON 본문으로 받기
  • UIMessageChunk 이벤트의 SSE 스트림 반환
  • x-workflow-run-id 응답 헤더 포함

GET {api}/{runId}/stream (예: /api/chat/{runId}/stream)

  • startIndex 쿼리 파라미터 받기
  • 주어진 청크 인덱스부터 시작하는 SSE 스트림 반환
  • 음수 startIndex의 경우 tail로 해석하고 x-workflow-stream-tail-index 응답 헤더 포함

전체 엔드포인트 예시는 WorkflowAgent 가이드를 참고하세요.

예시 (Examples)

useChat 기본 사용법 (Basic Usage with useChat)

'use client';

import { useChat } from '@ai-sdk/react';
import { WorkflowChatTransport } from '@ai-sdk/workflow/client';
import { useMemo } from 'react';

export default function Chat() {
  const transport = useMemo(
    () => new WorkflowChatTransport({ api: '/api/chat' }),
    [],
  );

  const { messages, sendMessage, status } = useChat({ transport });

  return (
    <div>
      {messages.map(message => (
        <div key={message.id}>
          {message.role === 'user' ? 'User: ' : 'AI: '}
          {message.parts.map((part, index) =>
            part.type === 'text' ? <span key={index}>{part.text}</span> : null,
          )}
        </div>
      ))}
      <button onClick={() => sendMessage({ text: 'Hello!' })}>Send</button>
    </div>
  );
}

콜백 사용 (With Callbacks)

'use client';

import { useChat } from '@ai-sdk/react';
import { WorkflowChatTransport } from '@ai-sdk/workflow/client';
import { useMemo } from 'react';

export default function Chat() {
  const transport = useMemo(
    () =>
      new WorkflowChatTransport({
        api: '/api/chat',
        maxConsecutiveErrors: 5,
        onChatSendMessage: response => {
          const runId = response.headers.get('x-workflow-run-id');
          console.log('Workflow run started:', runId);
        },
        onChatEnd: ({ chatId, chunkIndex }) => {
          console.log(`Chat ${chatId} complete, ${chunkIndex} chunks`);
        },
      }),
    [],
  );

  const { messages, sendMessage } = useChat({ transport });

  // ... render chat UI
}

서버 사이드 엔드포인트 (Next.js)

import { createModelCallToUIChunkTransform } from '@ai-sdk/workflow';
import { createUIMessageStreamResponse, type UIMessage } from 'ai';
import { start } from 'workflow/api';
import { chat } from '@/workflow/agent-chat';

export async function POST(request: Request) {
  const { messages }: { messages: UIMessage[] } = await request.json();
  const run = await start(chat, [messages]);

  return createUIMessageStreamResponse({
    stream: run.readable.pipeThrough(createModelCallToUIChunkTransform()),
    headers: {
      'x-workflow-run-id': run.runId,
    },
  });
}
import { createModelCallToUIChunkTransform } from '@ai-sdk/workflow';
import { createUIMessageStreamResponse } from 'ai';
import type { NextRequest } from 'next/server';
import { getRun } from 'workflow/api';

export async function GET(
  request: NextRequest,
  { params }: { params: Promise<{ runId: string }> },
) {
  const { runId } = await params;
  const startIndex = Number(
    new URL(request.url).searchParams.get('startIndex') ?? '0',
  );
  if (!Number.isSafeInteger(startIndex) || startIndex < 0) {
    return Response.json(
      { error: 'startIndex must be a non-negative safe integer' },
      { status: 400 },
    );
  }

  const run = await getRun(runId);
  const readable = run
    .getReadable({ startIndex: 0 })
    .pipeThrough(
      createModelCallToUIChunkTransform({ uiStartIndex: startIndex }),
    );

  return createUIMessageStreamResponse({
    stream: readable,
    headers: {
      'x-workflow-run-id': runId,
    },
  });
}

이 WorkflowAgent 엔드포인트는 인덱스 0부터 원시 ModelCallStreamPart 객체를 재생한 뒤, UIMessageChunk 객체로 변환한 후 트랜스포트의 음수가 아닌 커서를 적용해요. 음수 시작 인덱스는 저장된 객체가 이미 UIMessageChunk인 내구성 스트림이 필요해요.

더 알아보기 (Learn more)