fix(llm): keep system prompts on premium chat, reject incomplete streams, stop double-charging quota

This commit is contained in:
Yun Chan 2026-09-28 00:53:44 +09:00
parent d96601a283
commit d311e8123f
10 changed files with 1225 additions and 237 deletions

View file

@ -13,6 +13,13 @@ import type { LLMStatus, LLMModel, LLMAction, LLMConnectionState } from '@d3ro/c
import { resolveSystemPrompt } from './llm-prompts'
import { getBundledOllamaPath } from '../utils/paths'
import { normalizeLoopbackUrl } from '../utils/loopback'
import { readNdjsonLines } from '../utils/ndjson-reader'
import {
toChatRequest,
toRoleMessages,
type ChatStreamOptions as CoreChatStreamOptions,
type RoleMessage,
} from '@d3ro/core/llm-chat'
const logger = getLogger('LocalLLMService')
@ -73,15 +80,30 @@ interface OllamaGenerateResponse {
eval_count?: number
}
interface ChatStreamOptions {
model?: string
temperature?: number
maxTokens?: number
signal?: AbortSignal
timeoutMs?: number
/** @d3ro/core/llm-chat 공통 옵션 + Ollama 전용 keep_alive. */
interface ChatStreamOptions extends CoreChatStreamOptions {
keepAlive?: string
}
interface OllamaChatFrame {
message?: { content: string }
done: boolean
}
/** NDJSON 한 줄을 프레임 객체로 파싱한다. 깨진 줄은 스트림 전체 실패로 본다. */
function parseOllamaFrame<T extends object>(line: string): T {
let parsed: unknown
try {
parsed = JSON.parse(line)
} catch {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Ollama returned malformed NDJSON')
}
if (typeof parsed !== 'object' || parsed === null) {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Ollama returned malformed NDJSON')
}
return parsed as T
}
type AbortCause = 'timeout' | 'cancelled'
interface ActiveRequest {
@ -489,8 +511,6 @@ class LocalLLMService extends EventEmitter {
}
reader = response.body.getReader()
const decoder = new TextDecoder()
let buffer = ''
let fullText = ''
let lastChunk: OllamaGenerateResponse | null = null
const complete = (): GenerateResult => {
@ -505,52 +525,18 @@ class LocalLLMService extends EventEmitter {
return result
}
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()) continue
try {
const chunk = JSON.parse(line) as OllamaGenerateResponse
fullText += chunk.response
if (chunk.done) {
lastChunk = chunk
doneFrame = true
this.emit('token', { token: chunk.response, done: true })
if (chunk.response) yield chunk.response
return complete()
}
this.emit('token', { token: chunk.response, done: chunk.done })
yield chunk.response
} catch {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Ollama returned malformed NDJSON')
}
}
}
buffer += decoder.decode()
const trailing = buffer.trim()
if (trailing) {
try {
const chunk = JSON.parse(trailing) as OllamaGenerateResponse
fullText += chunk.response
if (chunk.done) {
lastChunk = chunk
doneFrame = true
this.emit('token', { token: chunk.response, done: true })
if (chunk.response) yield chunk.response
return complete()
}
this.emit('token', { token: chunk.response, done: chunk.done })
yield chunk.response
} catch {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Ollama returned malformed NDJSON')
for await (const line of readNdjsonLines(reader)) {
const chunk = parseOllamaFrame<OllamaGenerateResponse>(line)
fullText += chunk.response
if (chunk.done) {
lastChunk = chunk
doneFrame = true
this.emit('token', { token: chunk.response, done: true })
if (chunk.response) yield chunk.response
return complete()
}
this.emit('token', { token: chunk.response, done: chunk.done })
yield chunk.response
}
if (!doneFrame) {
@ -845,7 +831,7 @@ class LocalLLMService extends EventEmitter {
* 각 토큰마다 yield, 완료 시 전체 응답 텍스트를 return.
*/
async *chatStream(
messages: Array<{ role: string; content: string }>,
messages: RoleMessage[],
options?: ChatStreamOptions,
): AsyncGenerator<string, string> {
if (!this._available) {
@ -865,7 +851,8 @@ class LocalLLMService extends EventEmitter {
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model,
messages,
// @d3ro/core/llm-chat 계약: system 은 선두 system 메시지 하나로 정규화해 보낸다.
messages: toRoleMessages(toChatRequest(messages)),
stream: true,
keep_alive: options?.keepAlive,
think: false,
@ -882,55 +869,19 @@ class LocalLLMService extends EventEmitter {
}
reader = response.body.getReader()
const decoder = new TextDecoder()
let buffer = ''
let accumulated = ''
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()) continue
try {
const chunk = JSON.parse(line) as { message?: { content: string }; done: boolean }
if (chunk.done) {
doneFrame = true
}
if (chunk.message?.content) {
accumulated += chunk.message.content
yield chunk.message.content
}
if (chunk.done) {
return accumulated
}
} catch {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Ollama returned malformed NDJSON')
}
for await (const line of readNdjsonLines(reader)) {
const chunk = parseOllamaFrame<OllamaChatFrame>(line)
if (chunk.done) {
doneFrame = true
}
}
buffer += decoder.decode()
const trailing = buffer.trim()
if (trailing) {
try {
const chunk = JSON.parse(trailing) as { message?: { content: string }; done: boolean }
if (chunk.done) {
doneFrame = true
}
if (chunk.message?.content) {
accumulated += chunk.message.content
yield chunk.message.content
}
if (chunk.done) {
return accumulated
}
} catch {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Ollama returned malformed NDJSON')
if (chunk.message?.content) {
accumulated += chunk.message.content
yield chunk.message.content
}
if (chunk.done) {
return accumulated
}
}

View file

@ -15,7 +15,16 @@ import { getCloudSyncService } from './CloudSyncService'
import { D3ROError, ErrorCode } from '@d3ro/core/errors'
import type { LLMAction } from '@d3ro/core/types'
import { resolveSystemPrompt } from './llm-prompts'
import { parseAnthropicSSE } from '../utils/sse-parser'
import { parseAnthropicSSE, AnthropicStreamError } from '../utils/sse-parser'
import {
fitChatRequest,
LLM_PROXY_CHAT_LIMITS,
toChatRequest,
type ChatRequest,
type ChatStreamOptions,
type ChatTurn,
type RoleMessage,
} from '@d3ro/core/llm-chat'
const logger = getLogger('PremiumLLMService')
@ -23,15 +32,9 @@ const logger = getLogger('PremiumLLMService')
// 내부 타입
// ============================================================
/** Ollama-style 메시지 → Claude Messages 변환용 */
interface ChatMessage {
role: 'user' | 'assistant'
content: string
}
/** llm-proxy Edge Function 요청 body */
interface LlmProxyRequest {
messages: ChatMessage[]
messages: ChatTurn[]
system?: string
max_tokens?: number
model?: string
@ -55,9 +58,105 @@ export interface QuotaUsageSnapshot {
overageCredits: number
}
type UpgradeReason = 'quota_exceeded' | 'model_not_allowed' | 'auth_required'
/** 호출 단위 취소·기한 상태. cancelGeneration() 은 활성 호출 전부를 취소한다. */
type AbortCause = 'timeout' | 'cancelled'
interface PremiumCall {
controller: AbortController
abortCause: AbortCause | null
abort: (cause: AbortCause) => void
close: () => void
}
/** llm-proxy 출력 토큰 상한 (llm-contract MAX_OUTPUT_TOKENS) */
const PROXY_MAX_OUTPUT_TOKENS = 4096
const DEFAULT_CHAT_MAX_TOKENS = 2048
/** 클라이언트 측 채팅 기한. 프록시가 공급자 호출에 45초 기한을 두므로 여유 있게 잡는다. */
const DEFAULT_CHAT_TIMEOUT_MS = 120_000
/**
* llm-proxy 오류 메시지를 D3ROError 로 분류한다 (순수 함수).
* upgradeReason 이 있으면 호출자가 'upgrade-required' 를 emit 한다.
*/
export function classifyProxyError(message: string): { error: D3ROError; upgradeReason: UpgradeReason | null } {
if (message.includes('401') || message.includes('Unauthorized') || message.includes('auth')) {
return {
error: new D3ROError(ErrorCode.LLMServerUnreachable, `인증 실패: ${message}`),
upgradeReason: 'auth_required',
}
}
if (message.includes('quota_exceeded') || message.includes('429')) {
return {
error: new D3ROError(ErrorCode.LLMProcessingFailed, `쿼터 초과: ${message}`),
upgradeReason: 'quota_exceeded',
}
}
if (message.includes('model_not_allowed') || message.includes('403')) {
return {
error: new D3ROError(ErrorCode.LLMInvalidAction, `모델 권한 없음: ${message}`),
upgradeReason: 'model_not_allowed',
}
}
return {
error: new D3ROError(ErrorCode.LLMProcessingFailed, `llm-proxy: ${message}`),
upgradeReason: null,
}
}
/**
* invokeFunctionStream 오류가 프록시의 HTTP 응답(`<status>: <body>`)인지 판별한다.
* 프록시가 응답했다면 같은 요청을 비스트리밍으로 다시 보내도 같은 실패(와 쿼터 소비)만
* 되풀이되므로, 비스트리밍 폴백은 HTTP 상태가 없는 전송 계층 실패에서만 쓴다.
*/
export function proxyHttpStatus(message: string): number | null {
const match = /^(\d{3}):/.exec(message)
return match ? Number(match[1]) : null
}
function resolveMaxTokens(maxTokens: number | undefined): number {
if (maxTokens === undefined) return DEFAULT_CHAT_MAX_TOKENS
if (!Number.isSafeInteger(maxTokens) || maxTokens <= 0) {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'maxTokens must be a positive safe integer')
}
return Math.min(maxTokens, PROXY_MAX_OUTPUT_TOKENS)
}
function resolveTimeoutMs(timeoutMs: number | undefined): number {
if (timeoutMs === undefined) return DEFAULT_CHAT_TIMEOUT_MS
if (!Number.isSafeInteger(timeoutMs) || timeoutMs <= 0) {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'timeoutMs must be a positive safe integer')
}
return timeoutMs
}
/** ChatRequest → llm-proxy body. system 은 최상위 system 필드로 보낸다. */
function toProxyBody(
request: ChatRequest,
maxTokens: number,
model: string | undefined,
stream: boolean,
): LlmProxyRequest {
const body: LlmProxyRequest = {
messages: request.turns,
max_tokens: maxTokens,
stream,
}
if (request.system !== undefined) body.system = request.system
if (model !== undefined) body.model = model
return body
}
function firstText(response: ClaudeMessageResponse): string {
const firstBlock = response.content?.[0]
return firstBlock?.type === 'text' ? firstBlock.text : ''
}
interface PremiumLLMEvents {
'quota-warning': (payload: { current: number; limit: number; overageCredits: number }) => void
'upgrade-required': (payload: { reason: 'quota_exceeded' | 'model_not_allowed' | 'auth_required' }) => void
'upgrade-required': (payload: { reason: UpgradeReason }) => void
'fallback-triggered': (payload: { reason: string }) => void
}
@ -66,7 +165,7 @@ interface PremiumLLMEvents {
// ============================================================
class PremiumLLMService extends EventEmitter {
private _abortController: AbortController | null = null
private _activeCalls = new Set<PremiumCall>()
private _disposed = false
private _lastQuota: QuotaUsageSnapshot | null = null
@ -145,64 +244,78 @@ class PremiumLLMService extends EventEmitter {
}
/**
* 스트리밍 대화 (Voice Conversation용).
* 스트리밍 대화 (Voice Conversation / 회의 채팅용).
* SSE 스트리밍: llm-proxy에 stream=true로 요청, Anthropic SSE를 토큰 단위 yield.
*
* @d3ro/core/llm-chat 계약을 따른다:
* - role:'system' 메시지는 버리지 않고 body.system 으로 보낸다 (한도 초과 시 앞부분 보존).
* - message_stop 없이 끊긴 스트림, 공급자 오류 이벤트, 빈 응답은 D3ROError 로 throw 한다.
* - signal / timeoutMs / maxTokens 는 호출 단위로 적용된다. temperature 는 프록시가 받지 않는다.
* - 비스트리밍 폴백은 전송 계층 실패에서만 쓴다 (프록시 HTTP 오류는 그대로 전파).
*/
async *chatStream(
messages: Array<{ role: string; content: string }>,
options?: { model?: string; temperature?: number },
messages: RoleMessage[],
options?: ChatStreamOptions,
): AsyncGenerator<string, string> {
this._ensureAuth()
const claudeMessages: ChatMessage[] = messages
.filter((m) => m.role === 'user' || m.role === 'assistant')
.map((m) => ({ role: m.role as 'user' | 'assistant', content: m.content }))
const body: LlmProxyRequest = {
messages: claudeMessages,
model: options?.model,
max_tokens: 2048,
stream: true,
const request = fitChatRequest(toChatRequest(messages), LLM_PROXY_CHAT_LIMITS)
if (request.turns.length === 0) {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Premium chat requires a user message')
}
const body = toProxyBody(request, resolveMaxTokens(options?.maxTokens), options?.model, true)
const call = this._beginCall(options?.signal, resolveTimeoutMs(options?.timeoutMs))
let completed = false
this._abortController = new AbortController()
const cloud = getCloudSyncService()
const { stream, error } = await cloud.invokeFunctionStream(
'llm-proxy',
body as unknown as Record<string, unknown>,
this._abortController.signal,
)
if (error || !stream) {
const msg = error?.message ?? 'Stream unavailable'
logger.error(`SSE stream failed: ${msg}`)
// SSE 실패 시 비스트리밍 fallback
logger.info('Falling back to non-streaming Premium LLM')
const fallbackBody = { ...body, stream: false }
const response = await this._invokeProxy(fallbackBody)
const firstBlock = response.content?.[0]
const text = firstBlock?.type === 'text' ? firstBlock.text : ''
if (text.length > 0) yield text
return text
}
let accumulated = ''
try {
const cloud = getCloudSyncService()
const { stream, error } = await cloud.invokeFunctionStream(
'llm-proxy',
body as unknown as Record<string, unknown>,
call.controller.signal,
)
if (error || !stream) {
const msg = error?.message ?? 'Stream unavailable'
if (call.controller.signal.aborted) {
throw new D3ROError(ErrorCode.LLMProcessingCancelled, msg)
}
logger.error(`SSE stream failed: ${msg}`)
if (proxyHttpStatus(msg) !== null) {
// 프록시가 응답한 오류 — 재요청은 같은 실패와 쿼터 소비만 되풀이한다.
throw this._rejectProxyError(msg)
}
// 전송 계층 실패 시 비스트리밍 fallback (같은 system/maxTokens/signal 적용)
logger.info('Falling back to non-streaming Premium LLM')
const response = await this._invokeProxy({ ...body, stream: false }, call.controller.signal)
const text = firstText(response)
if (!text.trim()) {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Premium LLM returned empty content')
}
completed = true
yield text
return text
}
let accumulated = ''
for await (const token of parseAnthropicSSE(stream)) {
accumulated += token
yield token
}
} catch (err) {
if ((err as Error).name !== 'AbortError') {
logger.warn(`SSE parse error: ${err instanceof Error ? err.message : String(err)}`)
completed = true
if (!accumulated.trim()) {
throw new D3ROError(ErrorCode.LLMProcessingFailed, 'Premium LLM returned empty content')
}
return accumulated
} catch (err) {
throw this._toChatError(err, call)
} finally {
this._abortController = null
// 완료 전 종료(소비자 break, 오류) 시 연결을 끊어 프록시 스트림을 정리한다.
if (!completed) call.abort('cancelled')
call.close()
}
return accumulated
}
/**
@ -227,11 +340,11 @@ class PremiumLLMService extends EventEmitter {
}
cancelGeneration(): void {
if (this._abortController) {
this._abortController.abort()
this._abortController = null
logger.info('Premium LLM generation cancelled')
if (this._activeCalls.size === 0) return
for (const call of [...this._activeCalls]) {
call.abort('cancelled')
}
logger.info('Premium LLM generation cancelled')
}
dispose(): void {
@ -243,32 +356,19 @@ class PremiumLLMService extends EventEmitter {
// ── 내부: Edge Function 호출 ─────────────────────────────
private async _invokeProxy(body: LlmProxyRequest): Promise<ClaudeMessageResponse> {
private async _invokeProxy(body: LlmProxyRequest, signal?: AbortSignal): Promise<ClaudeMessageResponse> {
const cloud = getCloudSyncService()
// Supabase JS 클라이언트의 functions.invoke() 사용 — auth 헤더를 올바르게 처리.
// raw fetch + Authorization: Bearer 방식은 Supabase gateway가 401로 거부.
const { data, error } = await cloud.invokeFunction('llm-proxy', body as unknown as Record<string, unknown>)
const { data, error } = signal
? await cloud.invokeFunction('llm-proxy', body as unknown as Record<string, unknown>, { signal })
: await cloud.invokeFunction('llm-proxy', body as unknown as Record<string, unknown>)
if (error) {
const msg = error.message ?? 'Edge Function error'
logger.error(`llm-proxy error: ${msg}`)
// 에러 메시지 기반 분류
if (msg.includes('401') || msg.includes('Unauthorized') || msg.includes('auth')) {
this.emit('upgrade-required', { reason: 'auth_required' })
throw new D3ROError(ErrorCode.LLMServerUnreachable, `인증 실패: ${msg}`)
}
if (msg.includes('quota_exceeded') || msg.includes('429')) {
this.emit('upgrade-required', { reason: 'quota_exceeded' })
throw new D3ROError(ErrorCode.LLMProcessingFailed, `쿼터 초과: ${msg}`)
}
if (msg.includes('model_not_allowed') || msg.includes('403')) {
this.emit('upgrade-required', { reason: 'model_not_allowed' })
throw new D3ROError(ErrorCode.LLMInvalidAction, `모델 권한 없음: ${msg}`)
}
throw new D3ROError(ErrorCode.LLMProcessingFailed, `llm-proxy: ${msg}`)
throw this._rejectProxyError(msg)
}
// functions.invoke는 response body를 자동 파싱해서 data에 넣음
@ -281,6 +381,60 @@ class PremiumLLMService extends EventEmitter {
return result
}
/** 프록시 오류를 분류하고 필요 시 upgrade-required 를 알린다. */
private _rejectProxyError(message: string): D3ROError {
const { error, upgradeReason } = classifyProxyError(message)
if (upgradeReason) this.emit('upgrade-required', { reason: upgradeReason })
return error
}
/** 호출 단위 AbortController 를 만들고 외부 signal·기한·cancelGeneration 에 연결한다. */
private _beginCall(externalSignal: AbortSignal | undefined, timeoutMs: number): PremiumCall {
const controller = new AbortController()
const onExternalAbort = (): void => call.abort('cancelled')
const timer = setTimeout(() => call.abort('timeout'), timeoutMs)
const call: PremiumCall = {
controller,
abortCause: null,
abort: (cause) => {
if (call.abortCause !== null) return
call.abortCause = cause
controller.abort()
},
close: () => {
clearTimeout(timer)
externalSignal?.removeEventListener('abort', onExternalAbort)
this._activeCalls.delete(call)
},
}
this._activeCalls.add(call)
if (externalSignal?.aborted) {
onExternalAbort()
} else {
externalSignal?.addEventListener('abort', onExternalAbort, { once: true })
}
return call
}
/** 채팅 스트림 실패를 호출 단위 원인에 맞는 D3ROError 로 정규화한다. */
private _toChatError(error: unknown, call: PremiumCall): D3ROError {
if (call.abortCause === 'timeout') {
return new D3ROError(ErrorCode.LLMProcessingTimeout, 'Premium LLM generation timed out')
}
if (call.abortCause === 'cancelled' || call.controller.signal.aborted) {
return new D3ROError(ErrorCode.LLMProcessingCancelled, 'Premium LLM generation cancelled')
}
if (error instanceof D3ROError) return error
if (error instanceof AnthropicStreamError) {
logger.warn(`Premium SSE stream failed (${error.kind}): ${error.message}`)
return new D3ROError(ErrorCode.LLMProcessingFailed, `Premium LLM stream failed: ${error.message}`)
}
return new D3ROError(
ErrorCode.LLMProcessingFailed,
`Premium LLM failed: ${error instanceof Error ? error.message : String(error)}`,
)
}
// ── EventEmitter 타입 오버라이드 ───────────────────────
on<K extends keyof PremiumLLMEvents>(event: K, listener: PremiumLLMEvents[K]): this {

View file

@ -0,0 +1,33 @@
// src/main/utils/ndjson-reader.ts
// NDJSON 스트림을 줄 단위로 읽는다 (버퍼링·청크 경계·마지막 개행 없는 줄 처리만 담당).
// 줄 파싱 정책(깨진 줄 처리, 완료 프레임 요구)은 호출자가 정한다.
/**
* reader 에서 비어 있지 않은 NDJSON 줄을 차례로 yield 한다.
* 개행 없이 끝난 마지막 줄은 trim 해서 yield 한다.
* reader 의 해제·취소는 호출자 책임이다.
*/
export async function* readNdjsonLines(
reader: ReadableStreamDefaultReader<Uint8Array>,
): AsyncGenerator<string, void> {
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()) continue
yield line
}
}
buffer += decoder.decode()
const trailing = buffer.trim()
if (trailing) yield trailing
}

View file

@ -1,6 +1,12 @@
// src/main/utils/sse-parser.ts
// Anthropic Claude Messages API SSE 스트림 파서
// content_block_delta 이벤트에서 텍스트 토큰을 추출하는 AsyncGenerator
// content_block_delta 이벤트에서 텍스트 토큰을 추출하는 AsyncGenerator.
//
// 종료 계약 (@d3ro/core/llm-chat 참조):
// - message_stop 또는 data: [DONE] 을 받아야 정상 종료한다.
// - `error` 이벤트(예: overloaded_error)는 AnthropicStreamError('provider_error')로 throw.
// - 종료 프레임 없이 스트림이 끝나면 AnthropicStreamError('incomplete')로 throw.
// (프록시 타임아웃 등으로 잘린 응답을 완료로 오인하지 않기 위함)
interface ContentBlockDelta {
type: 'content_block_delta'
@ -10,15 +16,80 @@ interface ContentBlockDelta {
}
}
interface AnthropicErrorEvent {
type: 'error'
error?: { type?: string; message?: string }
}
interface SSEEvent {
type: string
[key: string]: unknown
}
export type AnthropicStreamErrorKind = 'provider_error' | 'incomplete'
/** SSE 스트림이 정상 완료되지 못했음을 나타낸다. */
export class AnthropicStreamError extends Error {
readonly kind: AnthropicStreamErrorKind
readonly providerErrorType: string | null
constructor(kind: AnthropicStreamErrorKind, message: string, providerErrorType: string | null = null) {
super(message)
this.name = 'AnthropicStreamError'
this.kind = kind
this.providerErrorType = providerErrorType
}
}
type LineOutcome =
| { kind: 'skip' }
| { kind: 'token'; text: string }
| { kind: 'stop' }
/** SSE 한 줄을 해석한다. 오류 이벤트는 throw 한다. */
function interpretLine(line: string): LineOutcome {
const trimmed = line.trim()
// 빈 줄, 이벤트 타입 라인 (event:), 주석(:) 건너뜀 — 타입은 data 의 type 필드로 판별
if (!trimmed || !trimmed.startsWith('data:')) return { kind: 'skip' }
const payload = trimmed.slice(5).trimStart()
// "data: [DONE]" — 종료 시그널
if (payload === '[DONE]') return { kind: 'stop' }
let evt: SSEEvent
try {
evt = JSON.parse(payload) as SSEEvent
} catch {
// JSON 파싱 실패 — 건너뜀
return { kind: 'skip' }
}
if (evt.type === 'message_stop') return { kind: 'stop' }
if (evt.type === 'error') {
const errorEvent = evt as unknown as AnthropicErrorEvent
const providerType = errorEvent.error?.type ?? 'unknown_error'
const detail = errorEvent.error?.message ?? 'Provider stream error'
throw new AnthropicStreamError('provider_error', `${providerType}: ${detail}`, providerType)
}
if (evt.type === 'content_block_delta') {
const delta = evt as unknown as ContentBlockDelta
if (delta.delta?.type === 'text_delta' && delta.delta.text) {
return { kind: 'token', text: delta.delta.text }
}
}
// 그 외 이벤트 (message_start, content_block_start, ping 등)는 건너뜀
return { kind: 'skip' }
}
/**
* ReadableStream<Uint8Array>을 파싱하여 텍스트 토큰을 yield.
* Anthropic SSE 형식: "data: {json}\n\n" 라인 단위.
* content_block_delta.delta.text 추출, message_stop 또는 [DONE] 시 종료.
* Anthropic SSE 형식: "event: x\ndata: {json}\n\n" 라인 단위.
* content_block_delta.delta.text 추출, message_stop 또는 [DONE] 시 정상 종료.
* 오류 이벤트나 종료 프레임 없는 EOF 는 AnthropicStreamError 로 throw.
*/
export async function* parseAnthropicSSE(
stream: ReadableStream<Uint8Array>,
@ -40,37 +111,19 @@ export async function* parseAnthropicSSE(
buffer = lines.pop() ?? ''
for (const line of lines) {
const trimmed = line.trim()
// 빈 줄 또는 이벤트 타입 라인 (event:) 건너뜀
if (!trimmed || trimmed.startsWith('event:')) continue
// "data: [DONE]" — 종료 시그널
if (trimmed === 'data: [DONE]') return
// "data: {...}" — JSON 파싱
if (trimmed.startsWith('data: ')) {
const json = trimmed.substring(6)
try {
const evt = JSON.parse(json) as SSEEvent
// message_stop → 스트림 종료
if (evt.type === 'message_stop') return
// content_block_delta → 텍스트 토큰 yield
if (evt.type === 'content_block_delta') {
const delta = evt as unknown as ContentBlockDelta
if (delta.delta?.type === 'text_delta' && delta.delta.text) {
yield delta.delta.text
}
}
// 그 외 이벤트 (message_start, content_block_start 등)는 건너뜀
} catch {
// JSON 파싱 실패 — 건너뜀
}
}
const outcome = interpretLine(line)
if (outcome.kind === 'stop') return
if (outcome.kind === 'token') yield outcome.text
}
}
// 개행 없이 끝난 마지막 줄 처리
buffer += decoder.decode()
const outcome = interpretLine(buffer)
if (outcome.kind === 'stop') return
if (outcome.kind === 'token') yield outcome.text
throw new AnthropicStreamError('incomplete', 'Anthropic stream ended before message_stop')
} finally {
reader.releaseLock()
}