WebSockets
SSE와 NDJSON은 턴마다 연결을 하나씩 엽니다. WebSocket은 다릅니다. 하나의 소켓이 전체 대화 동안 열려 있고 모든 턴을 전달하며 요청을 기다리지 않고 서버가 청크를 푸시할 수 있습니다. 메시지마다 요청을 보내는 모델 대신 영속 채널이 필요할 때 사용합니다. 이미 실행 중인 WebSocket 게이트웨이, 반복적인 핸드셰이크를 피하려는 모바일 클라이언트, 응답 외에도 서버가 푸시해야 하는 UI가 그 예입니다(이 페이지에서는 요청/응답 방식만 다루며, 서버가 시작하는 푸시에는 그 위에 자체 프레이밍이 필요합니다).
이 페이지를 끝까지 읽으면 연결이 끊긴 뒤에도 재개되는 소켓을 통해 모델을 다시 실행하지 않고 채팅 턴을 실행할 수 있습니다.
구조
하나의 소켓, 여러 턴입니다.
- 클라이언트가 소켓을 한 번 엽니다.
- 각 사용자 메시지는 소켓을 통해 전송되는 하나의 JSON 프레임이 됩니다(SSE/NDJSON POST 본문과 같은 형태입니다).
- 서버는 프레임마다 하나의
chat()턴을 실행하고 같은 소켓을 통해 청크를 다시 스트리밍합니다. - 클라이언트가 소켓을 닫거나, 유휴 시간 초과가 발생하거나, 프로세스가 중단될 때까지
RUN_FINISHED이후에도 다음 프레임을 기다리며 소켓이 열린 상태로 유지됩니다.
서버: onRun
이미 수락한 소켓을 toWebSocketStream과 연결합니다. 인바운드 프레임을 디코딩하고,
프레임마다 하나의 chat() 턴을 onRun을 통해 시작한 다음 결과 청크를 다시 전송합니다.
import { chat, memoryStream, toWebSocketStream } from '@tanstack/ai'
import { openaiText } from '@tanstack/ai-openai'
import type { WebSocketLike } from '@tanstack/ai'
// Bridge the per-turn AbortSignal WsRunContext hands you into the
// AbortController chat() expects.
function abortControllerFromSignal(signal: AbortSignal): AbortController {
const controller = new AbortController()
if (signal.aborted) controller.abort(signal.reason)
else
signal.addEventListener('abort', () => controller.abort(signal.reason), {
once: true,
})
return controller
}
// `socket` is a WHATWG-shaped server socket you already accepted (see
// "Hosting" below for where it comes from on Node vs Cloudflare) and
// `request` is the original handshake request.
export function handleChatSocket(socket: WebSocketLike, request: Request) {
toWebSocketStream(socket, request, {
// Per-turn durability, keyed by the frame's runId (see "Durability is
// per turn" below).
durability: (ctx) => memoryStream(ctx.request),
onRun: ({ messages, threadId, runId, signal }) =>
chat({
adapter: openaiText('gpt-5.5'),
messages,
threadId,
runId,
abortController: abortControllerFromSignal(signal),
}),
})
}
onRun은 인바운드 프레임마다 하나의 WsRunContext를 받습니다.
| 필드 | 설명 |
|---|---|
messages | 프레임에서 디코딩된 해당 턴의 UIMessage[] / ModelMessage[]입니다. |
threadId / runId | 이 턴의 AG-UI 식별자입니다. |
forwardedProps | 클라이언트가 프레임과 함께 보낸 추가 데이터입니다. |
request | 영속성 어댑터의 키로 사용할 수 있도록 턴별 합성 Request의 URL에 ?runId=를 담았습니다. |
signal | 소켓이 닫히거나 이 턴이 abort 프레임을 받으면 중단됩니다(아래 참조). 같은 소켓의 다른 턴은 중단하지 않습니다. |
연결이 끊겨도 유지될 필요가 없는 소켓이라면 durability를 완전히 생략합니다.
toWebSocketStream은 계속 청크를 전달하지만 재연결 시 재생할 수는 없습니다.
클라이언트: webSocket()
webSocket()은 SubscribeConnectionAdapter이며 useChat을 위한 어댑터입니다(@tanstack/ai-solid, @tanstack/ai-vue, @tanstack/ai-svelte, @tanstack/ai-angular에서도 내보냅니다). 처음
sendMessage를 호출할 때 소켓을 지연하여 열고, 대화의 이후 모든 턴에 재사용합니다.
import { useChat, webSocket } from '@tanstack/ai-react'
const connection = webSocket('/api/chat-ws')
export function Chat() {
const { messages, sendMessage } = useChat({ connection })
return <button onClick={() => void sendMessage('Hello')}>Send</button>
}
그 밖에는 아무것도 바뀌지 않습니다. messages, sendMessage, stop(), 도구 호출,
영속성은 다른 연결 어댑터와 동일하게 작동합니다. 옵션은 다음과 같습니다.
| 옵션 | 설명 |
|---|---|
protocols | WebSocket 생성자에 그대로 전달되는 WebSocket 서브프로토콜입니다. |
body | 전송되는 모든 RunAgentInput 프레임에 병합되는 정적 추가 필드입니다(HTTP 어댑터의 body와 같습니다). |
reconnect | 재연결 범위(maxAttempts, delayMs)이며, fetchServerSentEvents와 의미를 공유합니다. 고급을 참조하세요. |
WebSocketImpl | WebSocket 구현을 재정의합니다(테스트, 브라우저가 아닌 런타임). |
다른 어댑터 중 webSocket()의 위치는 연결 어댑터를 참조하세요.
와이어 프로토콜
| 방향 | 프레임 | 의미 |
|---|---|---|
| 클라이언트 → 서버 | RunAgentInput 형태의 JSON 객체(SSE/NDJSON POST 본문과 같은 형태) | 하나의 chat() 턴을 시작합니다. |
| 클라이언트 → 서버 | { "type": "abort", "runId": "…" } | 진행 중인 한 턴을 중단합니다(아래 참조). |
| 서버 → 클라이언트 | { "id": "…", "chunk": <StreamChunk> } | 영속성 오프셋이 태그된 하나의 청크입니다(durability가 구성된 경우에만 해당). |
| 서버 → 클라이언트 | 순수 StreamChunk | 태그가 없는 하나의 청크입니다(durability가 구성되지 않은 경우). |
| 서버 → 클라이언트 | { "type": "ping" } | 하트비트입니다. webSocket()은 이를 자동으로 삭제하며, 직접 작성한 클라이언트는 type: "ping"인 모든 항목을 무시해야 합니다. |
서버→클라이언트의 두 형태는 명확히 구분됩니다. 래퍼에는 최상위 type이 절대 없고,
모든 순수 StreamChunk에는 최상위 타입이 있습니다.
재개 정보는 헤더가 아닌 URL에 포함됩니다
SSE와 NDJSON은 fetch/XHR 요청이 열리기 전에 임의의 헤더를 설정할 수 있으므로
Last-Event-ID 헤더로 재개합니다. 브라우저의 WebSocket 생성자는 핸드셰이크에서
사용자 지정 헤더를 설정할 수 없으므로, 오프셋은 대신 URL에 포함됩니다: ?runId=<id>&offset=<lastId>.
webSocket()이 이를 대신 처리합니다. 실행의 종료 청크(RUN_FINISHED / RUN_ERROR)가
생성되기 전에 소켓이 끊기고 해당 실행이 영속성을 사용했다면(오프셋이 태그된 래퍼),
?runId=&offset=에서 다시 열고 재생된 경계를 중복 제거합니다. 이는
fetchServerSentEvents가 제공하는 것과 동일한 재연결 보장이며, 와이어에서 전달되는
방식만 다릅니다. 오프셋을 한 번도 내보내지 않은 실행(durability가 구성되지 않음)은
재개할 대상이 없으므로, 무한히 재시도하는 대신 연결 오류가 발생합니다.
서버에서 ?offset=을 담은 URL은 새로운 턴이 아니라 재개 요청입니다. 이를 모델 호출
없이 영속성 로그를 읽기 전용으로 재생하는 resumeWebSocketStream으로 라우팅합니다.
import { memoryStream, resumeWebSocketStream } from '@tanstack/ai'
import type { WebSocketLike } from '@tanstack/ai'
export function handleResumeSocket(socket: WebSocketLike, request: Request) {
resumeWebSocketStream(socket, { adapter: memoryStream(request) })
}
요청에 재개 오프셋이 없으면 재생할 항목이 없으므로 resumeWebSocketStream은 코드
1008로 소켓을 닫습니다.
영속성은 턴별로 적용됩니다
대화 범위의 소켓은 여러 턴을 전달하므로 전체 연결에 사용할 영속성 어댑터를 한 번만
만들 수 없습니다. 각 턴에는 해당 턴의 runId를 키로 사용하는 자체 로그가 필요합니다.
따라서 toWebSocketStream의 durability는 값이 아니라 팩토리입니다. 턴별 ctx를
받고 URL에 이미 해당 턴의 ?runId=가 포함된 ctx.request를 읽습니다
(memoryStream/durableStream은 이를 자동으로 키로 사용합니다).
턴 중단
{ type: 'abort', runId } 프레임은 해당 턴의 onRun 반복만 중단하며
(ctx.signal이 발생합니다), 다음 턴을 위해 소켓 자체는 열린 상태로 유지됩니다.
내장 webSocket() 클라이언트 어댑터는 실행의 중단 시그널이 발생할 때(useChat의
stop()) 이 프레임을 전송하므로, 서버는 더 이상 수신하지 않는 클라이언트에 대해
모델 호출을 끝까지 실행하지 않고 실제로 생성을 중단합니다. 직접 작성한 클라이언트도
턴을 취소하려면 같은 프레임을 보내야 합니다.
전체 소켓을 닫으면 해당 소켓에서 아직 진행 중인 모든 턴이 중단됩니다.
하트비트와 유휴 시간 초과
toWebSocketStream은 유휴 소켓을 끊는 프록시를 거쳐서도 연결을 유지하도록
heartbeatMs마다(기본값 30초) { type: 'ping' } 프레임을 보냅니다. idleTimeoutMs 동안
인바운드 프레임이 없으면 소켓을 닫습니다(기본값 5분). 하트비트 자체는 활동으로 간주되지
않으며 클라이언트가 보낸 프레임만 활동으로 간주됩니다. 턴이 계속 스트리밍되는 동안에는
유휴 시간 초과가 발생하지 않으므로, 단일 생성이 길어도(에이전트 루프 또는 시간 초과보다
긴 턴) 중간에 끊기지 않습니다.
import { chat, toWebSocketStream } from '@tanstack/ai'
import { openaiText } from '@tanstack/ai-openai'
import type { WebSocketLike } from '@tanstack/ai'
function handleChatSocket(socket: WebSocketLike, request: Request) {
toWebSocketStream(socket, request, {
onRun: ({ messages, threadId, runId }) =>
chat({ adapter: openaiText('gpt-5.5'), messages, threadId, runId }),
heartbeatMs: 15_000,
idleTimeoutMs: 60_000,
})
}
Node에서 호스팅
일반 Node(그리고 전역 WebSocketPair가 없는 다른 환경)에는 WebSocket 업그레이드를 수락하는 내장 방법이 없으므로 직접 처리한 뒤 결과를 toWebSocketStream에 전달해야 합니다. 패턴은 HTTP 서버의 upgrade 이벤트를 연결하고, ws의 WebSocketServer({ noServer: true })로 소켓을 수락한 다음 결과 소켓을 그대로 전달하는 것입니다. ws의 소켓은 이미 WebSocketLike에 필요한 send/close/addEventListener 인터페이스를 구현합니다.
import { WebSocketServer } from 'ws'
import { memoryStream, resumeWebSocketStream } from '@tanstack/ai'
import type { Plugin } from 'vite'
import { handleChatSocket } from './handle-chat-socket'
const WS_PATH = '/api/chat-ws'
export function webSocketChatPlugin(): Plugin {
return {
name: 'websocket-chat-plugin',
configureServer(server) {
if (!server.httpServer) return
const wss = new WebSocketServer({ noServer: true })
server.httpServer.on('upgrade', (req, socket, head) => {
const url = new URL(req.url ?? '/', `http://${req.headers.host}`)
if (url.pathname !== WS_PATH) return
wss.handleUpgrade(req, socket, head, (ws) => {
// `ws`'s socket satisfies WebSocketLike structurally — pass it
// straight through, no adapter needed.
const request = new Request(url)
if (url.searchParams.get('offset') !== null) {
resumeWebSocketStream(ws, { adapter: memoryStream(request) })
} else {
handleChatSocket(ws, request)
}
})
})
},
}
}
이는 작동하는
examples/ts-react-chat
WebSocket 예제와 e2e 스위트의
durable-delivery-ws-plugin.ts.
다른 곳에도 동일한 연결 패턴을 적용할 수 있습니다. 플랫폼의 소켓을 수락한 다음 toWebSocketStream / resumeWebSocketStream을 호출합니다. Deno의 Deno.upgradeWebSocket에서 가져온 소켓은 WHATWG 형태이므로 WebSocketLike를 직접 충족하지만, Bun의 ServerWebSocket(핸들러 객체 API)는 먼저 작은 어댑터가 필요합니다.
Cloudflare에서 호스팅
Cloudflare Workers(및 Durable Objects)는 전역 WebSocketPair를 노출하므로 직접 업그레이드할 필요가 없습니다. toWebSocketResponse가 쌍을 생성하고 서버 측을 수락한 뒤 toWebSocketStream에 연결하고 101 업그레이드 Response를 반환합니다.
import { chat, toWebSocketResponse } from '@tanstack/ai'
import { openaiText } from '@tanstack/ai-openai'
export default {
fetch(request: Request): Response {
return toWebSocketResponse(request, {
onRun: ({ messages, threadId, runId }) =>
chat({ adapter: openaiText('gpt-5.5'), messages, threadId, runId }),
})
},
}
resumeWebSocketResponse({ adapter })는 ?offset= 재개에 대응하는 읽기 전용 래퍼입니다. 두 함수 모두 WebSocketPair가 없는 곳(예: Node)에서 호출되면 대응하는 *Stream을 가리키는 메시지와 함께 오류를 발생시키므로, 소켓 자체를 업그레이드할 수 없는 런타임에 Cloudflare 래퍼를 실수로 배포할 수 없습니다.
프로듀서와 소켓 수명 주기
SSE와 NDJSON에 적용되는 동일한 주의사항이 여기에도 적용됩니다. memoryStream을 사용하면
실행의 프로듀서와 전달 소켓이 같은 프로세스에 존재합니다. 소켓이 끊기면 이를 뒷받침하는 chat() 호출도 중단되므로, 재연결은 아직 진행 중인 실행을 재개하는 대신 이미 기록된 내용만 재생합니다. durableStream은 둘을 분리하므로(프로듀서가 클라이언트의 소켓이 아니라 백엔드에 연결되어 실행됨), 소켓이 끊겨도 아직 생성 중인 실행에 재연결할 수 있습니다. 자세한 설명은 프로덕션에서의 memoryStream을 참조하세요. 해당 설명은 WebSocket에도 그대로 적용됩니다.