본문으로 건너뛰기

재개 가능한 스트림: 고급

개요에서는 어댑터를 선택하고 응답을 래핑한 다음 GET 핸들러를 추가하는 일반적인 경우를 다룹니다. 이 페이지에서는 나머지 내용을 설명합니다.

durableStream 옵션

durableStream(request, options)은 외부 Durable Streams 백엔드와 통신합니다.

import {
chat,
chatParamsFromRequest,
toServerSentEventsResponse,
} from '@tanstack/ai'
import { durableStream } from '@tanstack/ai-durable-stream'
import { openaiText } from '@tanstack/ai-openai'
// Your token source.
import { getDurableStreamsToken } from './auth'

const durableOptions = {
server: 'https://streams.example.com',
streamPrefix: 'chat-runs',
headers: async () => ({
Authorization: `Bearer ${await getDurableStreamsToken()}`,
}),
}

export async function POST(request: Request) {
const { messages, threadId, runId } = await chatParamsFromRequest(request)
const stream = chat({
adapter: openaiText('gpt-5.5'),
messages,
threadId,
runId,
})
return toServerSentEventsResponse(stream, {
durability: { adapter: durableStream(request, durableOptions), batch: 32 },
})
}
  • headers는 고정 자격 증명에 정적 객체를 사용하고, 순환 토큰에는 비동기 리졸버를 사용합니다. 리졸버는 생성, 추가, 읽기, 닫기마다 실행됩니다.
  • batch는 로그에 추가할 때 버퍼링하는 청크 수를 제어합니다(기본값 32).
  • 백엔드는 생성, 추가, 닫기에서 비어 있지 않은 Stream-Next-Offset 헤더를 반환해야 합니다. 헤더가 없으면 즉시 오류가 발생합니다. 어댑터는 오프셋을 추측하지 않습니다.
  • durableStream은 일반 StreamDurability를 반환합니다. 오프셋에는 백엔드가 할당한 커서가 포함되므로 호출자가 선택할 수 없으며, 어댑터에는 자신이 만들지 않은 오프셋에 범위를 다시 저장하는 upsert가 없습니다. UpsertableStreamDurability가 필요한 코드는 런타임이 아니라 연결 지점에서 컴파일에 실패합니다. 이를 제공할 수 있는 어댑터는 저장된 범위 다시 저장을 참조하세요. 프로듀서 재시작 후 실행 재개에는 필요하지 않습니다. 이미 스트리밍한 내용을 중복하지 않고 실행 재개를 참조하세요.

ID로 실행에 연결

연결이 끊긴 후 재연결은 자동입니다. 처음부터 실행에 의도적으로 연결하려면(두 번째 탭 또는 실행 ID를 이미 아는 전체 새로고침의 경우) joinRun을 호출합니다. 읽기 전용 GEToffset=-1로 수행하므로 서버에 개요의 GET 핸들러가 필요합니다. 이 핸들러는 resumeServerSentEventsResponse({ adapter })입니다(NDJSON에는 resumeHttpResponse 사용). 로그를 재생하고 요청에 재개 오프셋이 없으면 400을 반환합니다.

import { fetchServerSentEvents } from '@tanstack/ai-client'

async function attach(runId: string) {
const connection = fetchServerSentEvents('/api/chat')
for await (const chunk of connection.joinRun(runId)) {
console.log(chunk)
}
}

네 가지 HTTP 어댑터(fetchServerSentEvents, fetchHttpStream, xhrServerSentEvents, xhrHttpStream) 모두 joinRun을 제공합니다.

연결 해제, 중지 및 오류

내구성 있는 실행의 프로듀서는 전달 소켓과 분리되어 있습니다. 클라이언트 연결이 끊기면(페이지 새로고침 또는 연결 끊김) 응답은 취소되지만 실행은 자체 종료 이벤트까지 로그로 계속 흘러갑니다. 따라서 재연결하거나 마운트 시점에 joinRun을 호출하면 완료될 때까지 따라갈 수 있습니다. 이것이 세션 내 재연결뿐 아니라 전체 새로고침 후에도 실행을 재개할 수 있게 하는 이유입니다. 재생이 클라이언트 도구에서 끝나면 클라이언트는 재생 후 도구 결과를 보냅니다. 라이브 스트림과 같은 경로입니다.

실행은 진정한 취소 또는 실패가 있을 때만 조기에 종료됩니다.

  • 취소 — 응답에 AbortControllerabortController로 전달한 다음 중단합니다(사용자의 중지 버튼 또는 request.signalAbortSignal을 자체 컨트롤러에 전달하는 경우). 프로듀서가 중지되고 종료 RUN_ERROR가 추가됩니다. 단순한 클라이언트 연결 해제만으로는 이렇게 되지 않습니다. 연결 해제 시 실행도 중지하려면 컨트롤러를 전달하세요.
  • 프로바이더 실패 — 모델 스트림에서 예외가 발생하면 오류가 종료 RUN_ERROR로 추가됩니다.

두 경우 모두 프로듀서는 종료 시 close()를 기다리고, 닫기 전에 종료 이벤트를 기록하므로 재연결하거나 연결하는 클라이언트가 멈추지 않고 종료 이벤트를 봅니다. 이벤트 추가 또는 닫기가 실패하면 원인은 기본적으로 서버에 기록됩니다(연결한 클라이언트에는 일반적인 미완료 오류만 보이므로 실제 원인은 서버 로그에 있습니다). 자체 로거로 전달하려면 debug를 전달하세요.

import { memoryStream, toServerSentEventsResponse } from '@tanstack/ai'
import type { StreamChunk } from '@tanstack/ai'

function respond(request: Request, stream: AsyncIterable<StreamChunk>) {
return toServerSentEventsResponse(stream, {
durability: { adapter: memoryStream(request) },
debug: true, // or { logger } for a custom Logger
})
}

내구성 있는 소스는 자체 종료 이벤트(RUN_FINISHED/RUN_ERROR)로 끝나야 합니다. 정상 완료 시 로그는 소스가 종료 이벤트를 내보낸 경우에만 종료 상태가 됩니다. 종료 이벤트가 없으면 내구성 소비자는 한 번 재연결하고 진행하지 못한 후 DurableStreamIncompleteError와 함께 실패합니다. chat()은 항상 RUN_FINISHED를 내보내므로 직접 만든 스트림에만 영향을 줍니다.

프로덕션에서의 memoryStream

단일 실행 프로세스에서 memoryStream은 클라이언트 연결 해제 후에도 유지됩니다. 프로듀서가 소켓보다 오래 살아 로그로 계속 흘려보내므로 이후 재연결 또는 joinRun이 아직 실행 중인 실행을 재개합니다. 그래도 다음 두 가지 이유로 개발 및 단일 프로세스 배포에만 적합합니다.

  1. 로그가 한 프로세스의 메모리에 있으므로 다른 워커에 연결된 재연결에서는 아무것도 찾지 못합니다.
  2. 프로세스 자체가 종료되면 로그도 함께 사라져 중단된 실행은 스스로 재개하거나 종료 상태가 될 수 없습니다. 프로덕션 백엔드는 인메모리 프로세스에는 없는 리스/정리기(프로세스 종료 참조)를 추가합니다.

완료된 실행은 유예 기간 후 제거되므로 만료되었거나 알 수 없는 실행을 재개하면 멈추지 않고 즉시 오류가 발생하며, 아무것도 생성하지 않는 실행에 처음부터 연결하면 firstChunkDeadlineMs 후 실패합니다. 프로세스 간 처리를 지원하는 것이 1번이며 durableStream이 이를 추가합니다.

재연결 횟수 제한

연결이 끊기면 마지막 오프셋부터 재개합니다. 오프셋이 유지되는 동안 전송 오류가 재시도되며, 해당 시도에서 재생된 중복 부분만 전달되었어도 마찬가지입니다. 내구성 있는 실행이 종료 이벤트 없이 정상적으로 끝나고 앞으로 진행되지 않으면 DurableStreamIncompleteError와 함께 실패합니다. 내구성이 없는(태그가 없는) 스트림만 정상 종료 시 완료된 실행으로 간주됩니다. 이는 의도적인 구분입니다. 정상적인 닫힘은 서버가 응답을 끝냈다는 뜻이므로 내구성 전송은 프로듀서가 살아 있는 동안 빈 롱 폴링 창을 끝내서는 안 됩니다.

클라이언트는 시도 사이에 속도를 제한하고 maxAttempts로 재연결을 제한하며, 초과하면 StreamReconnectLimitError와 함께 실패합니다. 한도는 새 이벤트를 전달하지 못한 연속 재연결만 계산하고 앞으로 진행되면 0으로 초기화됩니다. 프록시가 이벤트마다 소켓을 교체해도 정상적으로 오래 실행되는 작업은 한도에 가까워지지 않습니다. 실행이 실제로 멈췄을 때만 발생합니다.

import { fetchServerSentEvents } from '@tanstack/ai-client'

function makeConnection() {
return fetchServerSentEvents('/api/chat', {
reconnect: { maxAttempts: 5, delayMs: 250 }, // defaults shown
})
}

durableStream도 같은 방식으로 자체 읽기 루프를 제한합니다. 창 중간에 본문 읽기가 실패하면 마지막 유효 위치부터 재시도하고 연속 실패 횟수를 제한합니다(reconnect: { maxReadFailures: 10, delayMs: 250 }). 정상적인 롱 폴링 진행은 절대 속도를 제한하지 않습니다.

WebSocket 전송

위의 내용은 SSE와 NDJSON, 즉 턴마다 연결을 하나씩 사용하는 방식을 설명합니다. WebSocket은 전체 대화 동안 열린 상태로 유지되므로 일부 동작이 다릅니다. 전체 프로토콜과 코드는 WebSockets을 참조하세요. 이 절에서는 이 페이지의 나머지 내용과 비교해 달라지는 점만 다룹니다.

재연결 정보는 헤더가 아니라 URL에 포함됩니다. SSE와 NDJSON은 Last-Event-ID로 재개합니다. fetch/XHR은 요청을 열기 전에 임의의 헤더를 설정할 수 있기 때문입니다. 브라우저의 WebSocket 생성자는 핸드셰이크에서 사용자 지정 헤더를 설정할 수 없으므로 오프셋은 URL에 포함됩니다: ?runId=<id>&offset=<lastId>. webSocket()은 내구성 있는 실행의 연결이 끊기면 해당 URL로 자동 재연결합니다. fetchServerSentEvents와 동일한 중복 제거 보장을 제공하지만 전송 방식이 다릅니다.

소켓은 실행 범위가 아니라 대화 범위입니다. 하나의 소켓이 여러 턴을 전달합니다. 따라서 toWebSocketStreamdurability 옵션은 한 번 만든 단일 값이 아니라 각 턴의 runId를 키로 하는 팩토리입니다(턴별로 생성한 요청을 통해). 소켓 자체는 클라이언트가 닫거나 유휴 시간 초과 또는 프로세스 종료가 발생할 때 닫히며, 한 턴의 RUN_FINISHED가 도착할 때 닫히지 않습니다.

하트비트와 유휴 시간 초과가 SSE에서 무료로 제공되는 전송 수준 연결 유지를 대신합니다. toWebSocketStreamheartbeatMs마다 핑을 보냅니다(기본값 30초). 인바운드 프레임이 없는 상태로 idleTimeoutMs가 지나면 소켓을 닫습니다(기본값 5분). 단, 턴이 아직 스트리밍 중일 때는 닫지 않으므로 단일 생성이 오래 걸려도 안전합니다.

중단 프레임은 소켓이 아니라 한 턴을 대상으로 합니다. { type: 'abort', runId }는 해당 턴의 onRun 반복만 중단하며 소켓은 다음 턴을 위해 열려 있습니다. 소켓을 닫으면 아직 진행 중인 모든 턴이 중단됩니다.

호스팅은 런타임에 따라 다릅니다. Cloudflare Workers 또는 Durable Objects에서는 toWebSocketResponse가 전역 WebSocketPair를 사용하므로 수동 업그레이드가 필요하지 않습니다. Node(또는 WebSocketPair가 없는 환경)에서는 직접 연결을 업그레이드해야 합니다. 일반적으로 wsWebSocketServer({ noServer: true })를 HTTP 서버의 upgrade 이벤트에 연결하고, 결과 소켓을 toWebSocketStream / resumeWebSocketStream에 직접 전달합니다.

프로덕션에서의 memoryStream에 설명된 프로듀서와 소켓의 주의사항은 동일하게 적용됩니다. memoryStream에서는 연결이 끊긴 WebSocket이 이를 뒷받침하는 chat() 호출을 중단합니다. 연결이 끊긴 SSE와 같습니다. 재연결하면 아직 실행 중인 모델 호출을 재개하는 대신 로그를 재생합니다. durableStream은 SSE와 NDJSON에서처럼 둘을 분리합니다.

오프셋 소유권

StreamDurability<TOffset>가 오프셋 형식을 소유합니다. 코어는 어댑터가 반환한 값을 해당 어댑터에 다시 전달하고 와이어에 기록할 뿐입니다(SSE id: 필드 또는 NDJSON 봉투의 id{ id, chunk }). 추가되는 각 배치에 대해:

  1. 코어는 청크를 전달하기 전에 append(chunks)를 호출합니다.
  2. 어댑터는 순서대로 청크마다 정확히 하나의 오프셋을 반환합니다.
  3. 코어는 누락되거나 추가되거나 비어 있거나 공백/CR/LF를 포함하는 오프셋을 거부합니다.
  4. 재개 시 제공된 오프셋 이후만 엄격하게 읽습니다.

코어는 배열 인덱스에서 오프셋을 도출하지 않으며 StreamChunk에 오프셋을 기록하지도 않습니다.

이미 스트리밍한 내용을 중복하지 않고 실행 재개

실행 중간에 재시작한 프로듀서가 소스의 일부를 재생하는 경우 중복 부분을 두 번 추가해서는 안 됩니다. 클라이언트에는 오프셋 중복 제거 외에 안전망이 없습니다. 다시 추가된 청크는 새 오프셋을 가지므로 seen이 이를 억제하지 못하고, 스트림 프로세서는 텍스트 델타와 도구 호출 인수를 무조건 연결합니다. 이는 성능 저하가 아니라 중복된 문장과 {"a":1}{"a":1} 인수입니다.

이를 처리하는 방법은 호출자가 선택한 오프셋이 아니라 snapshot()과 일반 append를 함께 사용하는 것입니다. 로그에 이미 있는 내용을 읽고 재생 내용과 비교한 다음 일치하는 접두사를 제거하고 나머지만 추가하세요.

import { memoryStream } from '@tanstack/ai'
import type { StreamChunk } from '@tanstack/ai'

async function appendAfterStored(
request: Request,
replayed: Array<StreamChunk>,
) {
const durability = memoryStream(request)
const stored = await durability.snapshot()
// `snapshot` returns without waiting, even while the log is open, so this
// works on a log whose previous producer died without calling `close()`.
const remainder = replayed.slice(stored.length)
if (remainder.length > 0) await durability.append(remainder)
}

재생된 에이전트 실행에 대해 @tanstack/ai-sandbox가 수행하는 작업도 이와 같습니다. 개수가 아니라 타임스탬프에 영향을 받지 않는 지문으로 비교하고, 접두사 내부에서 재생 내용과 로그가 일치하지 않으면 예외를 발생시킵니다. 전체 경로는 실행 저널을 참조하세요.

로그는 어댑터가 만든 오프셋과 함께 추가 전용으로 유지되므로 durableStream에서도 작동합니다. 이는 프로덕션에 권장되는 어댑터이며 upsert가 없습니다. 특히 durableStream.snapshot()에는 두 가지 주의사항이 있습니다. 백엔드가 따라잡기를 계속 보고하지 않으면 포기하기 전 가져올 윈도우 수에 내부 상한이 있으며, -1에서 읽으므로 백엔드가 생성하지 않은 스트림에 대해 []을 반환할 수 없습니다. 이는 빈 실행이 아니라 실패한 호출입니다.

호출자가 선택한 키에 쓸 수 있는 저장소를 가진 어댑터에는 선택적 기능으로 upsert가 계속 제공됩니다. 위의 재개 경로에서는 이를 사용하지 않으며 대부분의 통합에서도 요구하지 않습니다. memoryStream은 이를 제공하지만 durableStream은 제공하지 않습니다(위 참조). 직접 구축하려면 저장된 범위 다시 저장을 참조하세요.

Cloudflare Durable Streams

Durable Streams는 이 프로토콜을 사용하는 Cloudflare Workers 및 Durable Objects 백엔드를 제공하므로 durableStream이 새 어댑터 없이 직접 통신합니다.

TanStack AI 엔드포인트도 Workers에서 실행된다면 공개 URL 대신 서비스 바인딩을 통해 백엔드에 접근하세요. 어댑터에 주입하는 fetch가 모든 요청을 바인딩으로 라우팅하므로 트래픽은 Cloudflare 네트워크에 머물고 호출 권한은 bearer token이 아니라 바인딩이 부여합니다.

import { durableStream } from '@tanstack/ai-durable-stream'

interface Env {
// A service binding to the deployed Durable Streams Worker.
DURABLE_STREAMS: { fetch: typeof fetch }
}

function cloudflareAdapter(request: Request, env: Env) {
// No `server` needed: the binding routes by path, so the adapter uses an
// internal placeholder base and only the `/streams/...` path matters.
return durableStream(request, {
streamPrefix: 'chat-runs',
fetch: env.DURABLE_STREAMS.fetch.bind(env.DURABLE_STREAMS),
})
}

백엔드가 다른 곳에서 실행된다면 대신 server를 Worker의 공개 URL로 지정하세요.

import { durableStream } from '@tanstack/ai-durable-stream'

function urlAdapter(request: Request) {
return durableStream(request, {
server: 'https://durable-streams.example.workers.dev',
streamPrefix: 'chat-runs',
})
}

Durable Object에서 실행하는 것도 프로세스 종료에 설명된 lease/reaper를 충족합니다. DO 알람이 프로듀서가 종료된 실행을 종료 상태로 만들 수 있으므로 재연결하는 클라이언트는 영원히 기다리는 대신 종료 상태를 봅니다.

프로세스 종료

이미 종료된 프로세스는 정리 코드를 실행할 수 없으므로 실제 프로세스 종료를 finally 또는 close()만으로 보장할 수 없습니다. 프로덕션 백엔드는 lease/reaper를 추가해야 합니다.

  1. 프로듀서는 쓰는 동안 lease를 획득하거나 갱신합니다.
  2. 타이머, 알람 또는 백그라운드 워커가 만료된 lease를 감지합니다.
  3. reaper가 중단된 종료 상태를 기록하고 로그를 닫습니다.
  4. 리더는 영원히 기다리는 대신 해당 종료 상태를 확인합니다.

이는 프로세스 내부 응답 헬퍼가 아니라 내구성 서비스 또는 배포에 속합니다.

조언만이 아니라 메커니즘이 있는 경우도 있습니다. 샌드박스된 실행의 작업은 종료된 프로세스가 아니라 샌드박스 안에서 계속되므로 이후 요청이 이를 인수할 수 있습니다. sandboxRunDriver가 실행을 점유하고 종료된 호스트를 차단하며 같은 로그에 계속 추가합니다. 이는 종료된 실행을 종료 상태로 만드는 것이 아니라 라이브 프로듀서를 인수하는 것이며 샌드박스된 실행에만 적용됩니다. 인수 및 분리된 실행을 참조하세요. 그 외의 경우에는 위의 조언이 계속 적용됩니다. memoryStream은 참여할 수 없습니다. 로그가 프로듀서 자체 프로세스에 있으므로 해당 프로세스에서 살아남아 인수될 것이 없습니다.

전달은 상태가 아닙니다

내구성 로그는 청크를 재생합니다. 스레드 메시지나 대화 기록을 조회할 수 있는 소스는 아닙니다. "이 실행이 무엇을 스트리밍했는가?"에는 답하지만 "이 사용자가 무엇을 말했는가?"에는 답하지 않습니다. 신뢰할 수 있는 상태는 자체 저장소에 보관하세요. 클라이언트 측 옵션은 클라이언트 영속성을 참조하세요.