import { type Api, type AssistantMessageEventStream, type Context, createAssistantMessageEventStream, type Model, type SimpleStreamOptions, } from "@earendil-works/pi-ai"; import { promptAcpSession } from "./acp-client.js"; import { resolveGrokBinary } from "./binary.js"; import { applyUsage, ContentEmitter, createOutput, mapStopReason } from "./emitter.js"; import { buildFullPrompt, buildIncrementalPrompt, currentAssistantPrefixFingerprint, expectedNextPrefixFingerprint, } from "./prompt.js"; import { SessionController } from "./session-controller.js"; import type { ReasoningEffort } from "./types.js"; function mapReasoningEffort(value: string | undefined): ReasoningEffort | undefined { if (!value || value === "off") return "none"; if ( value === "minimal" || value === "low" || value === "medium" || value === "high" || value === "xhigh" || value === "max" ) { return value; } return undefined; } export function createGrokStream( controller: SessionController, getCwd: () => string, ): (model: Model, context: Context, options?: SimpleStreamOptions) => AssistantMessageEventStream { return (model, context, options) => { const stream = createAssistantMessageEventStream(); const output = createOutput(model); const emitter = new ContentEmitter(stream, output); queueMicrotask(async () => { try { if (options?.signal?.aborted) throw new Error("aborted"); emitter.start(); const reasoningEffort = mapReasoningEffort(options?.reasoning); const currentPrefix = currentAssistantPrefixFingerprint(context); await controller.run( { binary: resolveGrokBinary(), modelId: model.id, ...(reasoningEffort ? { reasoningEffort } : {}), cwd: getCwd(), ...(options?.signal ? { signal: options.signal } : {}), }, async (state) => { const prompt = state.bootstrapped ? buildIncrementalPrompt(context) : buildFullPrompt(context); if (!prompt.trim()) throw new Error("Cannot send an empty prompt to Grok Build"); const result = await promptAcpSession(state.session, prompt, { ...(options?.signal ? { signal: options.signal } : {}), onUpdate: (update) => { const text = typeof update.content?.text === "string" ? update.content.text : ""; if (update.sessionUpdate === "agent_thought_chunk") emitter.appendThinking(text); else if (update.sessionUpdate === "agent_message_chunk") emitter.appendText(text); else if (update.sessionUpdate === "tool_call") { emitter.noteActivity(`[Grok tool: ${update.title ?? update.kind ?? "tool"}]`); } }, }); applyUsage(model, output, result._meta); state.bootstrapped = true; state.expectedPrefixFingerprint = expectedNextPrefixFingerprint(context, output); emitter.done(mapStopReason(result.stopReason)); }, (state) => !state.bootstrapped || (state.expectedPrefixFingerprint !== undefined && state.expectedPrefixFingerprint === currentPrefix), ); } catch (error) { const aborted = options?.signal?.aborted === true || (error instanceof Error && error.message === "aborted"); emitter.error(error instanceof Error ? error.message : String(error), aborted); } }); return stream; }; }