Skip to content
← Public packages

@kody/codex

Create and manage OpenAI Agents API (Codex harness) cloud agent sessions.

src/events/stream.ts

101 lines · 3.4 KB · TypeScript
import { boolean, number, object, optional, parse, string } from 'remix/data-schema'
import {
	DOCS_EVENTS,
	agentsFetch,
	pickAuth,
	requireSessionId,
} from '../client.ts'

const streamInput = object(
	{
		sessionId: string(),
		/** When true (default), return URL/headers guidance without opening SSE. */
		guidanceOnly: optional(boolean()),
		/** Open SSE and collect up to maxEvents (isolate-bounded). */
		collect: optional(boolean()),
		maxEvents: optional(number()),
		apiKeySecret: optional(string()),
		dryRun: optional(boolean()),
	},
	{ unknownKeys: 'error' },
)

/**
 * Event-stream guidance and optional bounded SSE collect for
 * `GET /v1/agents/sessions/{session_id}/events?stream=true`.
 * Prefer `guidanceOnly` (default) inside Kody isolates — long-lived streams
 * belong in your app. Subscribe before sending follow-up input.
 *
 * @param raw.sessionId - Session id
 * @param raw.guidanceOnly - Return connect recipe without calling OpenAI (default true unless collect)
 * @param raw.collect - Open SSE and collect a bounded event list
 * @param raw.maxEvents - Cap collected events (default 50)
 * @param raw.dryRun - Same as guidance preview
 * @param raw.apiKeySecret - Optional alternate secret name
 * @returns Guidance object and/or collected events
 *
 * @example
 * import streamEvents from 'kody:@kody/codex/events/stream'
 * const guide = await streamEvents({ sessionId: 'sess_123' })
 * // => { mode: 'guidance', url, headers, terminalEvents, … }
 */
export default async function streamEvents(raw: unknown = {}) {
	const input = parse(streamInput, raw ?? {})
	const auth = pickAuth(input)
	const sessionId = requireSessionId(input.sessionId)
	const path = `/agents/sessions/${encodeURIComponent(sessionId)}/events`
	const url = `https://api.openai.com/v1${path}?stream=true`
	const collect = input.collect === true
	const guidanceOnly =
		input.dryRun === true || input.guidanceOnly === true || !collect

	const guidance = {
		mode: 'guidance' as const,
		sessionId,
		method: 'GET' as const,
		url,
		headers: {
			Authorization: 'Bearer $OPENAI_API_KEY',
			'OpenAI-Beta': 'agents=v1',
			Accept: 'text/event-stream',
		},
		curl: `curl -N "${url}" -H "OpenAI-Beta: agents=v1" -H "Authorization: Bearer $OPENAI_API_KEY" -H "Accept: text/event-stream"`,
		docs: DOCS_EVENTS,
		terminalEvents: [
			'agent.session.turn.completed',
			'agent.session.turn.failed',
			'agent.session.turn.cancelled',
		],
		notes: [
			'Open the stream before POST .../events follow-up so early events are not missed.',
			'agent.session.idle alone is not success — wait for a terminal turn event and inspect output.',
			'Streams do not replay. After disconnect: open a new stream, GET session + items, then resume.',
			'On agent.session.requires_action, GET the session and handle required_actions.',
			'Long-lived SSE is a poor fit for Kody Worker isolates — use collect for a short sample only.',
		],
	}

	if (guidanceOnly) {
		return guidance
	}

	const response = await agentsFetch(path, {
		...auth,
		method: 'GET',
		query: { stream: true },
		expectSse: true,
		maxSseEvents: input.maxEvents ?? 50,
	})

	return {
		mode: 'collected' as const,
		sessionId,
		ok: true,
		status: response.status,
		eventCount: response.events?.length ?? 0,
		events: response.events,
		guidance,
		warning:
			'Bounded collect only — for production UIs, stream outside Kody and recover via sessions/get + sessions/items.',
	}
}