158 lines
4.7 KiB
TypeScript
158 lines
4.7 KiB
TypeScript
import { formatMessage, messages } from './messages'
|
|
import type { RecommendationRequest, StreamEvent } from './types'
|
|
import type { TransportFailureKind } from './models'
|
|
import { parseStreamEvent, StreamParseError } from './streamParser'
|
|
|
|
type StreamPhase = 'metadata' | 'events' | 'terminal'
|
|
|
|
export class StreamTransportError extends Error {
|
|
readonly kind: TransportFailureKind
|
|
|
|
constructor(kind: TransportFailureKind, message: string) {
|
|
super(message)
|
|
this.name = 'StreamTransportError'
|
|
this.kind = kind
|
|
}
|
|
}
|
|
|
|
function parseLine(line: string): StreamEvent {
|
|
let value: unknown
|
|
try {
|
|
value = JSON.parse(line)
|
|
} catch {
|
|
throw new StreamTransportError('parse', messages.recommendationStreamInvalidJson)
|
|
}
|
|
|
|
try {
|
|
return parseStreamEvent(value)
|
|
} catch (error) {
|
|
if (error instanceof StreamParseError) {
|
|
throw new StreamTransportError('parse', error.message)
|
|
}
|
|
throw error
|
|
}
|
|
}
|
|
|
|
function advancePhase(phase: StreamPhase, event: StreamEvent): StreamPhase {
|
|
if (phase === 'metadata') {
|
|
if (event.type === 'metadata') return 'events'
|
|
if (event.type === 'error') return 'terminal'
|
|
throw new StreamTransportError('protocol', messages.recommendationStreamMissingMetadata)
|
|
}
|
|
if (phase === 'events') {
|
|
if (event.type === 'track' || event.type === 'warning') return 'events'
|
|
if (event.type === 'done' || event.type === 'error') return 'terminal'
|
|
throw new StreamTransportError('protocol', messages.recommendationStreamDuplicateMetadata)
|
|
}
|
|
throw new StreamTransportError('protocol', messages.recommendationStreamAfterFinalEvent)
|
|
}
|
|
|
|
async function cancelBody(body: ReadableStream<Uint8Array> | null): Promise<void> {
|
|
try {
|
|
await body?.cancel()
|
|
} catch {
|
|
return
|
|
}
|
|
}
|
|
|
|
async function cancelReader(reader: ReadableStreamDefaultReader<Uint8Array>): Promise<void> {
|
|
try {
|
|
await reader.cancel()
|
|
} catch {
|
|
return
|
|
}
|
|
}
|
|
|
|
function releaseReader(reader: ReadableStreamDefaultReader<Uint8Array>): void {
|
|
try {
|
|
reader.releaseLock()
|
|
} catch {
|
|
return
|
|
}
|
|
}
|
|
|
|
/** Post a recommendation request and consume its NDJSON event stream. */
|
|
export async function streamRecommendations(
|
|
request: RecommendationRequest,
|
|
signal: AbortSignal,
|
|
onEvent: (event: StreamEvent) => void,
|
|
): Promise<void> {
|
|
let response: Response
|
|
try {
|
|
response = await fetch('/api/recommendations', {
|
|
method: 'POST',
|
|
credentials: 'same-origin',
|
|
headers: { 'Content-Type': 'application/json' },
|
|
body: JSON.stringify(request),
|
|
signal,
|
|
})
|
|
} catch (error) {
|
|
if (signal.aborted) throw error
|
|
throw new StreamTransportError('network', messages.recommendationStreamUnavailable)
|
|
}
|
|
|
|
if (!response.ok) {
|
|
await cancelBody(response.body)
|
|
throw new StreamTransportError(
|
|
'http',
|
|
formatMessage('recommendationStreamHttpFailure', { status: response.status }),
|
|
)
|
|
}
|
|
const contentType = response.headers.get('content-type')?.split(';')[0].trim()
|
|
if (contentType !== 'application/x-ndjson') {
|
|
await cancelBody(response.body)
|
|
throw new StreamTransportError('protocol', messages.recommendationStreamInvalidContentType)
|
|
}
|
|
if (!response.body) {
|
|
throw new StreamTransportError('protocol', messages.recommendationStreamEmpty)
|
|
}
|
|
|
|
const reader = response.body.getReader()
|
|
const decoder = new TextDecoder()
|
|
let buffer = ''
|
|
let phase: StreamPhase = 'metadata'
|
|
|
|
try {
|
|
while (true) {
|
|
const { done, value } = await reader.read()
|
|
buffer += decoder.decode(value, { stream: !done })
|
|
const lines = buffer.split('\n')
|
|
buffer = lines.pop() ?? ''
|
|
|
|
for (const rawLine of lines) {
|
|
const line = rawLine.trim()
|
|
if (!line) continue
|
|
const event = parseLine(line)
|
|
phase = advancePhase(phase, event)
|
|
onEvent(event)
|
|
if (phase === 'terminal') {
|
|
await cancelReader(reader)
|
|
return
|
|
}
|
|
}
|
|
|
|
if (done) break
|
|
}
|
|
|
|
const finalLine = buffer.trim()
|
|
if (finalLine) {
|
|
const event = parseLine(finalLine)
|
|
phase = advancePhase(phase, event)
|
|
onEvent(event)
|
|
if (phase === 'terminal') return
|
|
}
|
|
|
|
throw new StreamTransportError('unexpected_eof', messages.recommendationStreamUnexpectedEnd)
|
|
} catch (error) {
|
|
await cancelReader(reader)
|
|
if (signal.aborted || error instanceof StreamTransportError) throw error
|
|
throw new StreamTransportError('network', messages.recommendationStreamInterrupted)
|
|
} finally {
|
|
releaseReader(reader)
|
|
}
|
|
}
|
|
|
|
/** Return whether a rejected stream operation was intentionally aborted. */
|
|
export function isAbortError(error: unknown, signal: AbortSignal): boolean {
|
|
return signal.aborted || (error instanceof DOMException && error.name === 'AbortError')
|
|
}
|