Skip to content
← 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 · TypeScript
import { 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),
	})
}