/** * OpenAI-completions streaming provider. * * One wire protocol — `POST {baseUrl}/chat/completions` with `stream: true` * and SSE — covers a long list of vendors: * OpenAI, DeepSeek, Groq, xAI (Grok), Cerebras, OpenRouter, Mistral, * Ollama, LM Studio, plus any custom `/v1`-compatible endpoint. * * Only base URL + auth key change between them, so this single * implementation is the workhorse of the pi harness. */ import { log } from '../../../../shared/logger.js'; import type { PiStreamRequest, PiStreamEvent, PiMessage, PiContentBlock, PiStopReason, PiUsage, } from './types.js'; import { fetchWithRetry, readWithIdleTimeout } from './retry.js'; import { classifyPiError, classifyPiNetworkError } from './humanize-error.js'; /* ── SSE parser (LF or CRLF tolerant, flushes the trailing event) ── */ async function* parseSse(res: Response): AsyncIterable { if (!res.body) return; const reader = res.body.getReader(); const decoder = new TextDecoder(); let buffer = ''; try { while (true) { const { value, done } = await readWithIdleTimeout(reader, 'OpenAI-compat'); if (done) break; buffer += decoder.decode(value, { stream: true }); let idx; while ( (idx = (() => { const a = buffer.indexOf('\n\n'); const b = buffer.indexOf('\r\n\r\n'); if (a < 0) return b; if (b < 0) return a; return Math.min(a, b); })()) !== -1 ) { const isCrlf = buffer.slice(idx, idx + 4) === '\r\n\r\n'; const raw = buffer.slice(0, idx); buffer = buffer.slice(idx + (isCrlf ? 4 : 2)); const parsed = parseEvent(raw); if (parsed !== undefined) yield parsed; } } buffer += decoder.decode(); if (buffer.trim()) { const parsed = parseEvent(buffer); if (parsed !== undefined) yield parsed; } } finally { try { reader.releaseLock(); } catch {} } } function parseEvent(raw: string): any | undefined { const lines = raw.split(/\r?\n/); const dataLines = lines .filter((l) => l.startsWith('data:')) .map((l) => l.slice(5).trimStart()); if (!dataLines.length) return undefined; const data = dataLines.join('\n'); if (!data || data === '[DONE]') return undefined; try { return JSON.parse(data); } catch { return undefined; } } /* ── Message conversion (pi → OpenAI) ── */ function isToolResultMessage(m: PiMessage): boolean { return m.role === 'user' && m.content.length > 0 && m.content.every((b) => b.type === 'tool_result'); } /** * The OpenAI Chat Completions schema wants tool results as their own * `role: "tool"` messages, one per result, immediately following the * assistant message that emitted the tool_calls. Our unified PiMessage * keeps tool_result blocks inside a user-role message instead — split them * here so the wire payload is valid. */ function toOpenAIMessages(pi: PiMessage[]): any[] { const out: any[] = []; for (const m of pi) { if (isToolResultMessage(m)) { for (const block of m.content) { if (block.type !== 'tool_result') continue; out.push({ role: 'tool', tool_call_id: block.toolUseId, content: block.content, }); } continue; } if (m.role === 'assistant') { const textParts: string[] = []; const toolCalls: any[] = []; for (const b of m.content) { if (b.type === 'text') textParts.push(b.text); else if (b.type === 'tool_use') { toolCalls.push({ id: b.id, type: 'function', function: { name: b.name, arguments: JSON.stringify(b.input || {}), }, }); } } const msg: any = { role: 'assistant', content: textParts.join('') || null }; if (toolCalls.length > 0) msg.tool_calls = toolCalls; out.push(msg); continue; } // role === 'user' with non-tool-result content (text + optional images). // Media parts go first; text is appended last (parity with the other // providers and pi/index's media-first block ordering). const contentBlocks: any[] = []; let plainText = ''; let hasMedia = false; for (const b of m.content) { if (b.type === 'text') { plainText += (plainText ? '\n' : '') + b.text; } else if (b.type === 'image') { hasMedia = true; contentBlocks.push({ type: 'image_url', image_url: { url: `data:${b.mediaType};base64,${b.data}` }, }); } else if (b.type === 'document') { // The Chat Completions schema has no document part — degrade to a text // note rather than crashing. The file is also on disk (saved-files // note), so the agent can open it with its tools. This shouldn't // normally happen: buildUserMessage gates documents on canNativeDocument // (false for this flavor), so a PDF here rides as the disk pointer. plainText += (plainText ? '\n' : '') + `[Attached document${b.name ? ` "${b.name}"` : ''} (${b.mediaType}) could not be inlined for this model — it is saved to disk; open it with your file tools.]`; } } if (hasMedia) { if (plainText) contentBlocks.push({ type: 'text', text: plainText }); out.push({ role: 'user', content: contentBlocks }); } else { out.push({ role: 'user', content: plainText }); } } return out; } function toOpenAITools(tools: { name: string; description: string; inputSchema: Record }[]) { return tools.map((t) => ({ type: 'function', function: { name: t.name, description: t.description, parameters: t.inputSchema, }, })); } function mapFinishReason(reason?: string): PiStopReason { switch (reason) { case 'stop': return 'end_turn'; case 'length': return 'max_tokens'; case 'tool_calls': case 'function_call': return 'tool_use'; case 'content_filter': return 'error'; default: return 'end_turn'; } } /* ── Streaming entry point ── */ interface PartialToolCall { id: string; name: string; argsBuf: string; } export async function* streamOpenAICompletions(req: PiStreamRequest): AsyncIterable { const url = `${req.baseUrl.replace(/\/+$/, '')}/chat/completions`; const openaiMessages = toOpenAIMessages(req.messages); // Inline the system prompt as the first message — most providers honour it // there; the few that prefer `instructions` (Responses API) aren't us. if (req.systemPrompt?.trim()) { openaiMessages.unshift({ role: 'system', content: req.systemPrompt }); } const body: any = { model: req.modelId, messages: openaiMessages, stream: true, // gpt-5.x / o-series reject the legacy `max_tokens`; the openai-api // sub-provider routes the cap through `max_completion_tokens` instead. [req.maxTokensField ?? 'max_tokens']: req.maxOutputTokens ?? 8192, }; // Without this opt-in, OpenAI/OpenRouter streams carry NO usage at all — and // usage is what feeds the supervisor's proactive session recycling. Gated // per sub-provider: Mistral's strict schema 422s on unknown fields // (noStreamUsage in sub-providers.ts); everyone else tolerates or needs it. if (req.includeStreamUsage !== false) { body.stream_options = { include_usage: true }; } if (req.tools && req.tools.length > 0) { body.tools = toOpenAITools(req.tools); // 'none' = the round-cap wrap-up round: the model must summarize, not // start more work. Tools stay declared so histories containing tool calls // remain valid. body.tool_choice = req.toolChoice === 'none' ? 'none' : 'auto'; } let res: Response; try { const headers: Record = { 'content-type': 'application/json', 'accept': 'text/event-stream', }; if (req.apiKey) headers['authorization'] = `Bearer ${req.apiKey}`; res = await fetchWithRetry(url, { method: 'POST', headers, body: JSON.stringify(body), signal: req.signal, }); } catch (err: any) { if (err?.name === 'AbortError') { yield { type: 'done', stopReason: 'aborted' }; return; } const cls = classifyPiNetworkError('OpenAI-compat', err); yield { type: 'error', error: cls.message, kind: cls.kind, retryable: cls.retryable }; return; } if (!res.ok) { let detail = ''; try { detail = await res.text(); } catch {} const cls = classifyPiError('OpenAI-compat', res.status, res.statusText, detail); yield { type: 'error', error: cls.message, status: cls.status, kind: cls.kind, retryable: cls.retryable }; return; } let accumulated = ''; let lastFinish: string | undefined; let usage: PiUsage | undefined; const toolCallsByIndex = new Map(); let chunkCount = 0; let firstChunkSummary = ''; let thinkingEmitted = false; // Vendors disagree on where streamed usage lives: spec says a final // choice-less chunk's `usage`, Groq defaults to nesting under `x_groq.usage`, // Moonshot tucks it onto the choice itself. Read all three. const readUsage = (u: any) => { if (!u || (u.prompt_tokens === undefined && u.completion_tokens === undefined)) return; usage = { inputTokens: u.prompt_tokens, outputTokens: u.completion_tokens }; }; try { for await (const chunk of parseSse(res)) { chunkCount++; if (chunkCount === 1) { try { firstChunkSummary = JSON.stringify(chunk).slice(0, 600); } catch {} } readUsage(chunk?.x_groq?.usage); const choice = chunk?.choices?.[0]; if (!choice) { readUsage(chunk?.usage); continue; } readUsage(choice?.usage); const delta = choice.delta || {}; // Reasoning models stream hidden thinking under vendor-specific fields // (DeepSeek/OpenRouter: reasoning_content; others: reasoning / // reasoning_text — upstream pi's field priority). Emit ONE liveness // pulse so the chat doesn't look hung; never forward the text itself. const reasoningDelta = delta.reasoning_content ?? delta.reasoning ?? delta.reasoning_text; if (!thinkingEmitted && typeof reasoningDelta === 'string' && reasoningDelta.length > 0) { thinkingEmitted = true; yield { type: 'thinking' }; } if (typeof delta.content === 'string' && delta.content.length > 0) { accumulated += delta.content; yield { type: 'text_delta', delta: delta.content }; } // Tool-call deltas: function name + arguments stream in pieces keyed by // `index`. Accumulate per index and only emit the full tool_use once we // see finish_reason: 'tool_calls' (or at stream end). const toolDeltas: any[] = Array.isArray(delta.tool_calls) ? delta.tool_calls : []; for (const td of toolDeltas) { const idx = typeof td.index === 'number' ? td.index : 0; let partial = toolCallsByIndex.get(idx); if (!partial) { partial = { id: td.id || '', name: '', argsBuf: '' }; toolCallsByIndex.set(idx, partial); } if (td.id) partial.id = td.id; const fn = td.function || {}; if (fn.name) partial.name = fn.name; if (typeof fn.arguments === 'string') partial.argsBuf += fn.arguments; } if (choice.finish_reason) lastFinish = choice.finish_reason; readUsage(chunk?.usage); } } catch (err: any) { if (err?.name === 'AbortError') { yield { type: 'done', stopReason: 'aborted' }; return; } const cls = classifyPiNetworkError('OpenAI-compat', err); yield { type: 'error', error: cls.message, kind: cls.kind, retryable: cls.retryable }; return; } log.info( `[pi/openai-compat] stream done — chunks=${chunkCount} text=${accumulated.length} ` + `toolCalls=${toolCallsByIndex.size} finishReason=${lastFinish || 'none'} ` + `promptTok=${usage?.inputTokens ?? '?'} outTok=${usage?.outputTokens ?? '?'}`, ); if (chunkCount === 0) { log.warn(`[pi/openai-compat] zero chunks parsed — content-type=${res.headers.get('content-type') || '?'}`); } else if (!accumulated && toolCallsByIndex.size === 0) { log.info(`[pi/openai-compat] first chunk (truncated): ${firstChunkSummary}`); } if (accumulated) yield { type: 'text_end', text: accumulated }; for (const partial of toolCallsByIndex.values()) { let input: any = {}; if (partial.argsBuf) { try { input = JSON.parse(partial.argsBuf); } catch { // Truncated/malformed tool-call JSON — almost always the output-token // cap cutting the arguments mid-stream. Executing a fabricated {_raw} // input produces misleading tool errors the model retries forever; // fail the round loudly instead (the session strips dangling tool_use // blocks from history on errored rounds). const capped = lastFinish === 'length'; yield { type: 'error', error: capped ? `The model's ${partial.name} call was cut off by the output-token limit (${req.maxOutputTokens ?? 8192} tokens) — the arguments did not fit. Try a smaller change, or raise the model's output budget.` : `The model emitted a malformed ${partial.name} tool call (arguments were not valid JSON).`, kind: 'other', retryable: !capped, }; yield { type: 'done', stopReason: 'error', usage }; return; } } yield { type: 'tool_use', id: partial.id || `call_${partial.name}_${Math.random().toString(36).slice(2, 10)}`, name: partial.name, input, }; } if (!accumulated && toolCallsByIndex.size === 0) { yield { type: 'error', error: `Provider returned no output (finishReason=${lastFinish || 'unknown'}). ` + `This usually means the model rejected the request or hit a content filter.`, }; yield { type: 'done', stopReason: 'error', usage }; return; } yield { type: 'done', stopReason: toolCallsByIndex.size > 0 ? 'tool_use' : mapFinishReason(lastFinish), usage, }; }