본문으로 건너뛰기

사용자 지정 내구성 어댑터

스트림이 유지되기를 원하는 저장소가 있습니다: Redis, Postgres, queue, Electric, object store입니다. 이 페이지를 끝까지 읽으면 StreamDurability 어댑터가 완성되어 toServerSentEventsResponse / toHttpResponse에 연결됩니다. 따라서 클라이언트는 모델을 다시 실행하지 않고 진행 중인 실행에 재연결할 수 있습니다.

Core는 사용자의 저장소를 전혀 이해하지 못합니다. Core는 사용자가 전달하는 불투명한 오프셋 문자열을 그대로 왕복 전달할 뿐입니다. 다음 다섯 가지 메서드를 구현합니다:

메서드작업
resumeFrom()이 요청에서 재개 오프셋을 반환하거나, 새 실행이면 null를 반환합니다.
append(chunks)전달 전에 배치를 영속화하고, 순서대로 각 청크에 대해 하나의 오프셋을 반환합니다.
read(offset, signal)offset 이후의 청크만 엄격하게 재생합니다.
close()실행을 완료로 표시하고 대기 중인 리더를 깨웁니다.
snapshot()기다리지 않고 지금 이 실행에 저장된 모든 항목을 반환합니다.

중요한 규칙

이를 잘못 구현하면 재개가 미묘한 방식으로 중단됩니다:

  • 오프셋은 불투명하고, 고유하며, 왕복 전달이 안전해야 합니다. 각 청크마다 서로 다른 오프셋을 반환합니다. 오프셋은 SSE id: 줄이나 NDJSON { id, chunk } 내부의 envelope를 통해 전달되므로 이를 견뎌야 합니다. Core는 빈 오프셋, 다음을 포함하는 오프셋을 거부합니다: NUL/CR/LF, 앞이나 뒤의 공백 또는 중복된 오프셋입니다.

  • read은 오프셋보다 엄격하게 이후의 항목을 재생하며, 가장 오래된 항목부터 처리하고, 다음을 볼 때가 아니라 로그가 닫힐 때 종료됩니다. 이는 대부분 "단순화" 과정에서 되돌아가기 쉬운 규칙이므로 그 이유를 설명하겠습니다. 다음은 불변 조건 memoryStream 자체의 read 루프가 전달하는 packages/ai/src/stream-durability.ts에서 인용한 내용입니다(테스트가 이를 고정합니다):

    터미널 청크(RUN_FINISHED / RUN_ERROR)는 읽기를 종료하지 않습니다. 즉, 에이전트 루프 실행은 반복마다 하나를 내보냅니다(finishReason "tool_calls" 다음에 "stop"). 따라서 첫 번째 청크에서 중지하면 도구 호출 실행이 첫 번째 도구 호출에서 잘립니다. 프로듀서는 close()를 호출하여 실제 완료를 알립니다. (모든 종료 경로에서 호출합니다. StreamDurability.close 참조), 그러면 log.complete가 설정됩니다. 그때까지 또는 호출자가 중단할 때까지 테일을 읽습니다.

    따라서 첫 번째 terminal chunk에서 반환하는 어댑터는 모든 재개된 tool-calling 실행을 잘라 냅니다. 재개된 클라이언트는 첫 번째 도구 호출을 확인한 후 정상적인 종료를 확인하며, 이를 "실행이 끝났습니다"로 해석합니다. close()은 로그 끝을 알리는 유일한 신호이며 core는 모든 producer 종료 시(완료, 취소 및 실패)에 이를 기다리므로 그때까지 tailing하면 항상 종료됩니다.

  • 실행이 아직 생성 중일 때 read은 빈 응답으로 끝나서는 안 됩니다. 대신 Park(다음 append를 기다림)합니다. 새 데이터가 없는 깔끔한 종료는 클라이언트에 실행이 끝났음을 알립니다. 그렇지 않으면 클라이언트는 DurableStreamIncompleteError와 함께 실패합니다. 중단 signal을 준수하여 연결이 끊긴 클라이언트가 대기를 중지하도록 합니다.

  • 순서 지정이나 append-before-deliver를 직접 처리하지 않습니다. Core가 버퍼링하고 append을 호출한 다음, 해당 offset을 반환하면 청크를 전달합니다.

  • snapshot은 절대 기다리지 않습니다. 현재 저장된 내용을 append 순서대로 즉시 확인하며, 실행이 아직 생성 중인 동안에도 확인하고, 저장된 내용이 없는 실행에서는 []로 확인합니다. read과 달리 대기하지 않습니다. 호출자는 로그의 끝을 계속 받으려는 것이 아니라 재개하기 전에 이전 host의 prefix를 확인하려고 합니다. 이 동작을 잘못 구현하면 정확히 처리가 필요한 로그에서 재개가 멈춥니다. 그 이유는 충돌한 producer가 close()을 호출하지 않아 해당 로그가 영원히 열린 상태로 남기 때문입니다. 그리고 그 위에 read를 사용하면 절대 완료되지 않습니다. 특히 다음을 재사용하지 마세요. 알 수 없는 실행 실패 경로에서 처음부터 시작하는 read 조인이 취하는 방식은 다음과 같습니다: read('-1')를 사용하여 빈 로그에서는 실패할 수 있지만, snapshot()[]로 확인되어야 합니다. 다음에 대해 거부하는 것은 전송, 프로토콜 또는 인증 실패는 여전히 올바른 처리입니다.

구현하기

스토어의 작업을 기준으로 어댑터를 작성합니다. 여기서는 사용자가 제공하는 실행별 추가 전용 로그를 사용합니다. 백엔드에 맞게 RunLog를 교체하세요:

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

// Your backend, one append-only log per run. Back it with Redis Streams, a
// Postgres table, a queue. Anything that returns a stable cursor per entry.
interface RunLog {
append: (chunks: Array<StreamChunk>) => Promise<Array<string>>
readAfter: (
cursor: string | null,
) => Promise<Array<{ cursor: string; chunk: StreamChunk }>>
isComplete: () => Promise<boolean>
waitForChange: (signal?: AbortSignal) => Promise<void>
markComplete: () => Promise<void>
// Everything stored so far, in append order. Must not wait for more.
readAll: () => Promise<Array<{ cursor: string; chunk: StreamChunk }>>
}

export function customDurability(
request: Request,
openLog: (runId: string) => RunLog,
): StreamDurability {
const url = new URL(request.url)
// The resume offset: native SSE reconnect header first, then a join's ?offset.
const resume =
request.headers.get('Last-Event-ID') ?? url.searchParams.get('offset')
// Your adapter owns run identity. Resolve it the way core's own
// `resolveResumeRunId` does, and the way `durableStream` does too: the
// X-Run-Id header first (a POST producer sends this), then the ?runId
// query (a GET join sends this). Never mint a fresh id when neither is
// present — a generated id addresses a log no attach request could ever
// name, so the run would appear to work while writing where nobody reads.
const runId =
request.headers.get('X-Run-Id') ?? url.searchParams.get('runId')
if (runId === null) {
throw new Error(
'a runId is required: send it as an X-Run-Id header or a ?runId query param',
)
}
const log = openLog(runId)

return {
resumeFrom: () => resume,
append: (chunks) => log.append(chunks),
close: () => log.markComplete(),
read: async function* (offset, signal) {
// '-1' / 'now' are the from-start / from-tail join sentinels.
let cursor: string | null = offset === '-1' ? null : offset
for (;;) {
if (signal?.aborted) return
const entries = await log.readAfter(cursor)
for (const entry of entries) {
cursor = entry.cursor
// Yield terminal chunks like any other. An agent-loop run emits a
// RUN_FINISHED per iteration, so returning on one would truncate a
// resumed tool-calling run at its first tool call.
yield { offset: entry.cursor, chunk: entry.chunk }
}
// The ONLY end-of-log condition: the producer called `close()`.
if (await log.isComplete()) return
// Park. Do NOT end the response here while the producer is alive.
await log.waitForChange(signal)
}
},
snapshot: async () => {
const entries = await log.readAll()
return entries.map((entry) => ({
offset: entry.cursor,
chunk: entry.chunk,
}))
},
}
}

기본 제공 어댑터와 정확히 동일한 방식으로 연결합니다:

import { chat, chatParamsFromRequest, toServerSentEventsResponse } from '@tanstack/ai'
import { openaiText } from '@tanstack/ai-openai'
// Your modules: the adapter above, and your backend's per-run log factory.
import { customDurability } from './durability'
import { openRunLog } from './run-log'

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: customDurability(request, openRunLog) },
})
}

NDJSON의 경우 toServerSentEventsResponsetoHttpResponse로 교체합니다. 어댑터는 동일하며, 와이어 인코딩만 변경됩니다.

저장된 범위 다시 영속화하기

일부 스토어는 호출자가 선택한 키에 쓸 수 있습니다. 해당 기능을 지원한다면, 이미 스트리밍한 범위를 재생하는 호출자가 각 청크를 이미 차지하고 있는 항목에 다시 기록할 수 있도록 추가 기능으로 노출할 수 있습니다.

실행을 재개하는 데 이 기능은 필요하지 않습니다. 실행 중간에서 다시 시작하는 권장 방법은 snapshot()을 사용해 저장된 접두사를 읽고 이를 억제한 다음, 나머지만 append하여 로그를 추가 전용으로 유지하는 것입니다. 이 방식은 모든 어댑터에서 작동합니다. 자세한 내용은 이미 스트리밍한 내용을 중복하지 않고 실행 재개하기를 참조하세요.

이 기능은 별도의 메서드인 upsert입니다. 이를 구현하면 어댑터는 UpsertableStreamDurability가 되며, 제외하면 일반적인 StreamDurability입니다. 기능이 없다는 것은 정직한 신호입니다. 해당 기능이 필요한 소비자는 UpsertableStreamDurability를 요청하므로, 불일치는 실행 로그 깊숙한 곳의 실패가 아니라 연결 지점에서 컴파일 오류로 발생합니다. 다음과 같은 기능은 제공하세요. 스토어가 호출자가 선택한 키에 쓸 수 있는 경우에만 제공해야 합니다. 예를 들어 Postgres INSERT ... ON CONFLICT (cursor) DO UPDATE 또는 명시적이고 중복 제거된 ID를 사용하는 Redis XADD가 있습니다. 모든 쓰기 작업에 자체 커서를 기록하는 스토어에서는 이 기능을 사용할 수 없으므로 upsert를 노출하지 않습니다. 이것이 durableStream가 선택한 방식입니다. memoryStream는 이 기능을 노출합니다.

각 항목은 청크를 해당 오프셋과 연결하므로 위치에 따라 정렬할 필요가 없습니다. 저장된 상태를 변경하기 전에 전체 배치를 검증하여, 거부된 호출이 일부만 적용된 상태를 남기지 않도록 하세요:

  • 직접 생성하지 않은 오프셋은 거부합니다. 직접 생성한 모든 오프셋은 정의상 재개할 수 있기 때문입니다;
  • 하나의 배치 안에서 반복되는 오프셋은 거부합니다;
  • 아직 저장되지 않은 오프셋은 현재 꼬리 뒤에 오도록 요구합니다. 따라서 겹치는 범위를 재생하는 호출자는 연속된 접미사를 재생해야 하며, 앞에서 건너뛴 위치에 쓸 수 없습니다.
import type { StreamChunk, UpsertableStreamDurability } from '@tanstack/ai'

// Your backend, plus the one operation `upsert` needs beyond `append`: a write
// at a cursor you supply that replaces whatever is already stored there.
interface UpsertableRunLog {
tailSeq: () => Promise<number>
hasOffset: (offset: string) => Promise<boolean>
write: (
entries: Array<{ chunk: StreamChunk; offset: string }>,
) => Promise<Array<string>>
}

// Offsets here are `${runId}:${seq}`. Yours can be any format you can decode
// back into a run id and a position.
function decodeSeq(runId: string, offset: string, index: number): number {
const prefix = `${runId}:`
const seq = offset.startsWith(prefix)
? Number(offset.slice(prefix.length))
: Number.NaN
if (!Number.isSafeInteger(seq) || seq < 1) {
throw new Error(
`entries[${index}].offset ${JSON.stringify(offset)} was not minted by this run`,
)
}
return seq
}

// The returned function is `async`, so every rejection reaches the caller as a
// rejected promise rather than a synchronous throw.
export function makeUpsert(
log: UpsertableRunLog,
runId: string,
): UpsertableStreamDurability['upsert'] {
return async (entries) => {
let tail = await log.tailSeq()
const seen = new Set<string>()
for (const [index, entry] of entries.entries()) {
const seq = decodeSeq(runId, entry.offset, index)
if (seen.has(entry.offset)) {
throw new Error(
`entries[${index}].offset ${JSON.stringify(entry.offset)} is repeated in this batch`,
)
}
seen.add(entry.offset)
if (await log.hasOffset(entry.offset)) continue
if (seq <= tail) {
throw new Error(
`entries[${index}].offset ${JSON.stringify(entry.offset)} is not stored yet but claims position ${seq}, at or before the tail ${tail}`,
)
}
tail = seq
}
// Every entry passed, so the write below cannot reject partway and leave a
// prefix of the batch applied.
return log.write(entries)
}
}

이를 위의 다섯 가지 메서드와 함께 upsert: makeUpsert(log, runId)로 어댑터에 전달하고, 결과에 UpsertableStreamDurability를 주석으로 추가하여 추가 기능이 타입에 나타나도록 하세요.

오프셋을 입력합니다(선택 사항)

StreamDurability<TOffset>은 오프셋 문자열에 대해 제네릭입니다. 이를 브랜딩하여 예상되는 오프셋 대신 원시 문자열을 전달할 수 없도록 합니다:

import type { StreamDurability } from '@tanstack/ai'

type MyOffset = string & { readonly __brand: 'MyOffset' }

// Your adapter is then StreamDurability<MyOffset>; append/read/resumeFrom all
// speak MyOffset, and a plain string won't type-check where one is expected.
type MyAdapter = StreamDurability<MyOffset>

Core는 여전히 이 값을 불투명한 값으로 처리하며, 브랜드는 자체 코드만 더 엄격하게 제한합니다.

터미널화는 사용자의 책임입니다

Core는 모든 프로듀서 종료 시(정상 완료, 취소, close()를 기다리고, 및 실패) 시 취소/실패 전에 종료 RUN_ERROR를 추가합니다. 종료합니다. close()readisComplete()true를 반환하고 깨워야 합니다 대기 중인 리더를 깨워야 하므로, 따라잡은 리더는 멈추며 멈추지 않고 대기하지 않습니다 — 이것이 유일한 read를 종료하는 방법입니다. 추가한 마지막 청크 자체는 종료하지 않습니다(참조 규칙). 백엔드 프로듀서가 close()를 실행하지 않고 종료될 수 있다면(프로세스 크래시), 버려진 로그를 종료 상태로 만드는 리스/리퍼를 추가합니다. 프로세스 종료를 참조합니다.