/** * Testable Factory Droid stream engine (poll-based session adapter). * * Factory is NOT an OpenAI-compatible endpoint — it is a session-based agent * orchestration API. This engine drives one pi "turn" through the Droid * Sessions API: * * 1. ensure a computer exists (env override → reuse active e2b → auto-create * from a cloud template) [user-selected option B] * 2. ensure a Factory session for the pi conversation (created once, reused; * the session retains history server-side) * 3. POST the NEWEST user message only (not the full history — the agent * loop already has the conversation context) * 4. poll `GET /sessions/{id}` until status → idle, extracting assistant * text deltas along the way (simulated streaming: Factory has no SSE) * 5. push pi events (start / text_start / text_delta / text_end / done) and * per-turn usage from the session's cumulative tokenUsage delta * * NOTE — nested-agent semantics: when you send a message to a Factory session, * the DROID runs its own agentic loop (its own tools: file reads, command * exec, edits) on the Factory computer and returns the final result text. pi's * own tool definitions are IGNORED (Droid manages its own tools). Each pi turn * therefore equals ONE complete Factory agent turn. The tool_use/tool_result * blocks inside the assistant content are Droid's internal actions and are not * surfaced as pi tool calls. */ import { randomUUID } from "node:crypto" import { resolveFactoryDroidApiKey } from "./auth.ts" import { createFactoryClient } from "./client.ts" import { COMPUTER_STATUS_ACTIVE, COMPUTER_STATUS_PROVISIONING, DEFAULT_COMPUTER_PROVISION_TIMEOUT_MS, DEFAULT_POLL_INTERVAL_MS, DEFAULT_POLL_TIMEOUT_MS, FACTORY_COMPUTER_ID_ENV, FACTORY_DROID_COMPUTER_NAME, FACTORY_MACHINE_TEMPLATE_ID_ENV, SESSION_STATUS_IDLE, } from "./config.ts" import { assistantTextFromMessages, assistantThinkingFromMessages, lastUserMessageText, numberValue, usageDelta, } from "./converters.ts" import type { AssistantMessageEventStreamLike, AssistantMessageLike, ContextLike, FactoryDroidDependencies, FactoryMessage, FactoryTokenUsage, ModelLike, StreamOptions, TextContent, ThinkingContent, Usage, } from "./types.ts" export * from "./auth.ts" export * from "./client.ts" export * from "./config.ts" export * from "./converters.ts" export * from "./cost.ts" export * from "./pricing.ts" export * from "./types.ts" function abortError(message = "The operation was aborted"): DOMException { return new DOMException(message, "AbortError") } function defaultUsage(): Usage { return { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, } } export function createStreamFactoryDroid(deps: FactoryDroidDependencies) { const baseUrl = deps.apiBase const fetchImpl = deps.fetchImpl ?? fetch const now = deps.now ?? (() => Date.now()) const uuid = deps.uuid ?? (() => randomUUID()) const cwd = deps.cwd ?? (() => process.cwd()) const state = deps.state ?? { sessionByConversation: new Map(), lastUsageBySession: new Map(), } const delay = deps.delay ?? ((ms: number, signal: AbortSignal) => { if (signal.aborted) return Promise.reject(abortError()) return new Promise((resolve, reject) => { const id = setTimeout(() => { signal.removeEventListener("abort", onAbort) resolve() }, ms) const onAbort = () => { clearTimeout(id) reject(abortError()) } signal.addEventListener("abort", onAbort, { once: true }) }) }) function raceAbort(promise: Promise, signal: AbortSignal): Promise { if (signal.aborted) return Promise.reject(abortError()) return new Promise((resolve, reject) => { const onAbort = () => reject(abortError()) signal.addEventListener("abort", onAbort, { once: true }) promise.then( (value) => { signal.removeEventListener("abort", onAbort) resolve(value) }, (error: unknown) => { signal.removeEventListener("abort", onAbort) reject(error) }, ) }) } // ----------------------------------------------------------------------- // Computer resolution (option B: auto-create) // ----------------------------------------------------------------------- async function ensureComputer(apiKey: string, signal: AbortSignal): Promise { // 1. explicit env override const env = deps.env ?? process.env const explicit = env[FACTORY_COMPUTER_ID_ENV]?.trim() if (explicit) return explicit const client = createFactoryClient({ apiKey, baseUrl, fetchImpl }) // 2. reuse an existing active e2b computer try { const { computers } = await raceAbort(client.listComputers(signal), signal) const active = computers.find( (c) => c.providerType === "e2b" && c.status === COMPUTER_STATUS_ACTIVE, ) if (active) return active.id } catch (error) { // Listing is best-effort; fall through to creation. if (isAbort(error)) throw error } // 3. auto-create from a cloud template const templateId = env[FACTORY_MACHINE_TEMPLATE_ID_ENV]?.trim() let resolvedTemplateId: string | undefined = templateId if (!resolvedTemplateId) { try { const { templates } = await raceAbort( client.listMachineTemplates(signal), signal, ) resolvedTemplateId = templates.find( (t) => t.buildStatus?.status === "success", )?.templateId } catch (error) { if (isAbort(error)) throw error } } if (!resolvedTemplateId) { throw new Error( `No reusable Factory computer and no machine template to auto-create one. ` + `Set ${FACTORY_COMPUTER_ID_ENV} to an existing computer id, or create a cloud template.`, ) } const created = await raceAbort( client.createComputer( { name: FACTORY_DROID_COMPUTER_NAME, provider: "e2b", source: { kind: "template", templateId: resolvedTemplateId }, }, signal, ), signal, ) // 4. poll until provisioning → active const provisionTimeout = DEFAULT_COMPUTER_PROVISION_TIMEOUT_MS const startedAt = now() for (;;) { if (signal.aborted) throw abortError() const status = created.status if (status === COMPUTER_STATUS_ACTIVE) return created.id if (status === "error") { throw new Error( `Factory computer ${created.id} entered error state during provisioning`, ) } if (now() - startedAt > provisionTimeout) { throw new Error( `Factory computer ${created.id} did not become active within ${provisionTimeout}ms`, ) } await delay(DEFAULT_POLL_INTERVAL_MS, signal) const fresh = await raceAbort(client.getComputer(created.id, signal), signal) if (fresh.status === COMPUTER_STATUS_ACTIVE) return fresh.id if (fresh.status === "error") { throw new Error( `Factory computer ${fresh.id} entered error state during provisioning`, ) } } } // ----------------------------------------------------------------------- // Session resolution // ----------------------------------------------------------------------- async function ensureSession( apiKey: string, computerId: string, model: ModelLike, options: StreamOptions | undefined, signal: AbortSignal, ): Promise { const conversationId = options?.sessionId ?? uuid() const existing = state.sessionByConversation.get(conversationId) if (existing) return existing const client = createFactoryClient({ apiKey, baseUrl, fetchImpl }) const sessionSettings: Record = { model: model.id } if (options?.reasoningEffort) sessionSettings.reasoningEffort = options.reasoningEffort const session = await raceAbort( client.createSession({ computerId, cwd: cwd(), sessionSettings }, signal), signal, ) state.sessionByConversation.set(conversationId, session.sessionId) // Baseline the cumulative usage at session creation so the FIRST turn's // delta excludes session-creation overhead (system prompt, etc.). if (session.tokenUsage && Object.keys(session.tokenUsage).length > 0) { state.lastUsageBySession.set(session.sessionId, session.tokenUsage) } return session.sessionId } function isAbort(error: unknown): boolean { return ( (error instanceof DOMException && error.name === "AbortError") || (error instanceof Error && error.name === "AbortError") ) } return function streamFactoryDroid( model: ModelLike, context: ContextLike, options?: StreamOptions, ): AssistantMessageEventStreamLike { const stream = deps.createStream() async function run() { const output: AssistantMessageLike = { role: "assistant", content: [], api: model.api, provider: model.provider, model: model.id, usage: defaultUsage(), stopReason: "stop", timestamp: now(), } const signal = options?.signal const controller = new AbortController() const abortUpstream = () => { if (!controller.signal.aborted) controller.abort() } if (signal?.aborted) abortUpstream() else signal?.addEventListener("abort", abortUpstream, { once: true }) let textBlock: TextContent | undefined let currentTextIdx = -1 let thinkingIdx = -1 let sessionId: string | undefined let resolvedApiKey: string | undefined const endTextBlock = () => { if (!textBlock) return stream.push({ type: "text_end", contentIndex: currentTextIdx, content: textBlock.text, partial: output, }) textBlock = undefined currentTextIdx = -1 } const endThinking = () => { if (thinkingIdx < 0) return const tc = output.content[thinkingIdx] if (tc && tc.type === "thinking") { stream.push({ type: "thinking_end", contentIndex: thinkingIdx, content: (tc as ThinkingContent).thinking, partial: output, }) } thinkingIdx = -1 } // Emit the growth of `assistantText` (and thinking) as deltas. const emitDeltas = (assistantText: string, thinkingText: string) => { if (thinkingText.length > 0 && thinkingIdx < 0) { output.content.push({ type: "thinking", thinking: thinkingText }) thinkingIdx = output.content.length - 1 stream.push({ type: "thinking_start", contentIndex: thinkingIdx, partial: output, }) } if (thinkingIdx >= 0) { const tc = output.content[thinkingIdx] if (tc && tc.type === "thinking") { const prev = (tc as ThinkingContent).thinking if (thinkingText.length > prev.length) { ;(tc as ThinkingContent).thinking = thinkingText stream.push({ type: "thinking_delta", contentIndex: thinkingIdx, delta: thinkingText.slice(prev.length), partial: output, }) } } } if (assistantText.length > 0 && !textBlock) { textBlock = { type: "text", text: "" } output.content.push(textBlock) currentTextIdx = output.content.length - 1 stream.push({ type: "text_start", contentIndex: currentTextIdx, partial: output, }) } if (textBlock) { const prev = textBlock.text if (assistantText.length > prev.length) { textBlock.text = assistantText stream.push({ type: "text_delta", contentIndex: currentTextIdx, delta: assistantText.slice(prev.length), partial: output, }) } } } try { stream.push({ type: "start", partial: output }) // 1. API key const apiKey = options?.apiKey && !options.apiKey.startsWith("$") ? options.apiKey : resolveFactoryDroidApiKey({ env: deps.env, authPaths: deps.authPaths, homeDir: deps.homeDir, }) resolvedApiKey = apiKey if (!apiKey) { throw new Error( `No Factory API key. Set the FACTORY_API_KEY env var ` + `or configure ~/.factory/settings.json / auth.json.`, ) } // 2. computer const computerId = await ensureComputer(apiKey, controller.signal) // 3. session (created once per pi conversation) sessionId = await ensureSession( apiKey, computerId, model, options, controller.signal, ) // 4. send the newest user message const userText = lastUserMessageText(context.messages) if (!userText) throw new Error("No user message text to send to the Factory session") const client = createFactoryClient({ apiKey, baseUrl, fetchImpl }) const { messageId } = await raceAbort( client.addMessage(sessionId, { text: userText, computerId }), controller.signal, ) // 5. poll until the agent loop finishes (simulated streaming) const pollInterval = DEFAULT_POLL_INTERVAL_MS const pollTimeout = options?.timeoutMs ?? DEFAULT_POLL_TIMEOUT_MS const startedAt = now() let lastAssistantText = "" let lastThinking = "" let finalUsage: FactoryTokenUsage | undefined for (;;) { if (controller.signal.aborted) throw abortError() const session = await raceAbort( client.getSession(sessionId, controller.signal), controller.signal, ) finalUsage = session.tokenUsage // Extract the assistant messages belonging to this turn: all // assistant messages created after our user message. const { messages } = await raceAbort( client.getMessages( sessionId, { role: "assistant", limit: 100 }, controller.signal, ), controller.signal, ) const turnMessages = assistantMessagesAfterUserMessage( messages, messageId, ) const assistantText = assistantTextFromMessages(turnMessages) const thinkingText = assistantThinkingFromMessages(turnMessages) emitDeltas(assistantText, thinkingText) lastAssistantText = assistantText lastThinking = thinkingText if (session.status === SESSION_STATUS_IDLE) break if (now() - startedAt > pollTimeout) { throw new Error( `Factory session ${sessionId} did not finish within ${pollTimeout}ms`, ) } await delay(pollInterval, controller.signal) } endThinking() endTextBlock() // 6. usage: per-turn delta from the session's cumulative tokenUsage const lastUsage = state.lastUsageBySession.get(sessionId) const delta = usageDelta(finalUsage, lastUsage) state.lastUsageBySession.set(sessionId, finalUsage ?? {}) output.usage.input = numberValue(delta.inputTokens) ?? 0 output.usage.output = numberValue(delta.outputTokens) ?? 0 output.usage.cacheRead = numberValue(delta.cacheReadTokens) ?? 0 output.usage.cacheWrite = numberValue(delta.cacheCreationTokens) ?? 0 output.usage.totalTokens = output.usage.input + output.usage.output + output.usage.cacheRead + output.usage.cacheWrite deps.calculateCost(model, output.usage) output.stopReason = "stop" stream.push({ type: "done", reason: "stop", message: output }) stream.end() } catch (error: unknown) { const aborted = isAbort(error) || controller.signal.aborted // Best-effort interrupt of the running agent loop. if (sessionId && !aborted && resolvedApiKey) { try { const client = createFactoryClient({ apiKey: resolvedApiKey, baseUrl, fetchImpl, }) void client.interrupt(sessionId).catch(() => undefined) } catch { // Interrupt is best-effort. } } const reason = aborted ? "aborted" : "error" output.stopReason = reason output.errorMessage = reason === "aborted" ? "Request aborted" : error instanceof Error ? error.message : String(error) stream.push({ type: "error", reason, error: output }) stream.end() } finally { signal?.removeEventListener("abort", abortUpstream) } } run().catch((error: unknown) => { const msg: AssistantMessageLike = { role: "assistant", content: [], api: model.api, provider: model.provider, model: model.id, usage: defaultUsage(), stopReason: "error", errorMessage: error instanceof Error ? error.message : String(error), timestamp: now(), } stream.push({ type: "error", reason: "error", error: msg }) stream.end() }) return stream } } /** * Filter a session's assistant messages to only those belonging to the turn * started by `userMessageId`. The session API returns messages newest-first or * oldest-first depending on pagination; we conservatively keep every assistant * message whose parent chain can be traced from the user message OR whose * creation time is after it. When no user message id matches (e.g. list * pagination), we return all assistant messages (worst case: the last turn's * answer, which is the correct one for a fresh conversation). */ export function assistantMessagesAfterUserMessage( messages: readonly FactoryMessage[], userMessageId: string, ): FactoryMessage[] { // Build a map id → message for parent-chain resolution. const byId = new Map() for (const m of messages) byId.set(m.id, m) const userMessage = byId.get(userMessageId) if (!userMessage) { // Fall back to every assistant message (first turn). return messages.filter((m) => m.role === "assistant") } // Find all messages whose parent chain reaches the user message. const result: FactoryMessage[] = [] const visited = new Set() const walk = (id: string) => { if (visited.has(id)) return visited.add(id) const msg = byId.get(id) if (!msg) return if (msg.role === "assistant") result.push(msg) if (msg.parentId) walk(msg.parentId) } // Walk from every assistant message up to the user message; we want // descendants of the user message, so instead walk DOWN is not possible // without children index. Simpler: a message belongs to this turn if its // parent chain, walked up, passes through the user message id. for (const m of messages) { if (m.role !== "assistant") continue let cursor: FactoryMessage | undefined = m let belongs = false let guard = 0 while (cursor && guard++ < 100) { if (cursor.id === userMessageId) { belongs = true break } cursor = cursor.parentId ? byId.get(cursor.parentId) : undefined } if (belongs) result.push(m) } return result }