← Public packages
@cameronpak/zo-computer
Typed Zo Computer API client with Result-based error handling and SSE streaming.
src/ask-stream.ts
162 lines · 3.8 KB · TypeScriptimport { Result } from 'better-result'
import { zoRequest } from './core'
import {
ZoConfigError,
ZoParseError,
ZoStreamError,
type ZoApiError,
type ZoNetworkError,
} from './errors'
import type {
ZoAskInput,
ZoClientOptions,
ZoStream,
ZoStreamEvent,
} from './types'
type StreamOpenError =
| ZoConfigError
| ZoNetworkError
| ZoApiError
| ZoParseError
/**
* Parse one raw SSE block into an event.
*/
function parseBlock(block: string): ZoStreamEvent | null {
let event = 'message'
const dataLines: Array<string> = []
for (const line of block.split('\n')) {
if (line.startsWith(':')) continue
if (line.startsWith('event:')) {
event = line.slice(6).trim()
continue
}
if (line.startsWith('data:')) {
dataLines.push(line.slice(5).replace(/^ /, ''))
}
}
if (dataLines.length === 0) return null
const raw = dataLines.join('\n')
let data: unknown = raw
try {
data = JSON.parse(raw) as unknown
} catch {
data = raw
}
return { event, data, raw }
}
async function* readEvents(
response: Response,
url: string,
): AsyncGenerator<ZoStreamEvent> {
const body = response.body
if (body === null) {
yield {
event: 'Error',
data: new ZoStreamError({
url,
cause: null,
message: 'Zo returned a streaming response with no body',
}),
raw: '',
}
return
}
const reader = body.pipeThrough(new TextDecoderStream()).getReader()
let buffer = ''
try {
while (true) {
const { done, value } = await reader.read()
if (done) break
buffer += value
let boundary = buffer.indexOf('\n\n')
while (boundary !== -1) {
const block = buffer.slice(0, boundary)
buffer = buffer.slice(boundary + 2)
const parsedEvent = parseBlock(block)
if (parsedEvent !== null) yield parsedEvent
boundary = buffer.indexOf('\n\n')
}
}
const tail = parseBlock(buffer)
if (tail !== null) yield tail
} catch (cause) {
yield {
event: 'Error',
data: new ZoStreamError({
url,
cause:
cause instanceof Error
? { name: cause.name, message: cause.message }
: cause,
message: 'The Zo event stream ended early',
}),
raw: '',
}
} finally {
await reader.cancel().catch(() => {})
}
}
/**
* Send one message to Zo and read the answer as it arrives.
*
* The `Err` case covers failures before the stream opens. Once the stream is
* open, a failure arrives as an event named `Error` whose `data` is a
* `ZoStreamError`.
*
* This export returns an async iterable, so import it. It cannot cross a
* `packages.invoke` boundary.
*
* @example
* import askStream from 'kody:@cameronpak/zo-computer/ask-stream'
* const opened = await askStream({ input: 'Write a haiku' })
* if (opened.status === 'ok') {
* for await (const event of opened.value.events) console.log(event.event)
* }
*/
export default async function askStream(
params: ZoAskInput & ZoClientOptions,
): Promise<Result<ZoStream, StreamOpenError>> {
if (typeof params?.input !== 'string' || params.input.trim() === '') {
return Result.err(
new ZoConfigError({
field: 'input',
message: 'input must be a non-empty string',
}),
)
}
const body: Record<string, unknown> = { input: params.input, stream: true }
if (params.conversationId !== undefined) {
body.conversation_id = params.conversationId
}
if (params.modelName !== undefined) body.model_name = params.modelName
if (params.personaId !== undefined) body.persona_id = params.personaId
if (params.outputFormat !== undefined) {
body.output_format = params.outputFormat
}
const responseResult = await zoRequest(
'/zo/ask',
{ method: 'POST', accept: 'text/event-stream', body },
params,
)
if (Result.isError(responseResult)) return responseResult
const response = responseResult.value
return Result.ok({
conversationId: response.headers.get('x-conversation-id'),
events: readEvents(response, response.url),
})
}