import { randomUUID } from "node:crypto"; import { existsSync, readFileSync } from "node:fs"; import { arch, homedir, platform, release } from "node:os"; import { join } from "node:path"; import type { AgentMessage, ThinkingLevel } from "@earendil-works/pi-agent-core"; import type { Model } from "@earendil-works/pi-ai"; import { convertToLlm, type ToolInfo } from "@earendil-works/pi-coding-agent"; export type AssistantPhase = "commentary" | "final_answer"; export type ResponseContentItem = | { type: "input_text"; text: string } | { type: "input_image"; image_url: string } | { type: "output_text"; text: string }; export type ResponseItem = | { type: "message"; role: string; content: ResponseContentItem[] | string; end_turn?: boolean; phase?: AssistantPhase; } | { type: "reasoning"; summary: Array<{ type: "summary_text"; text: string }>; content?: Array<{ type: "reasoning_text" | "text"; text: string }>; encrypted_content: string | null; } | { type: "function_call"; name: string; arguments: string; call_id: string } | { type: "function_call_output"; call_id: string; output: unknown } | { type: "compaction"; encrypted_content: string } | { type: "compaction_trigger" } | { type: string; [key: string]: unknown }; export interface ResponsesReasoningConfig { effort?: "none" | "minimal" | "low" | "medium" | "high" | "xhigh" | "max"; summary?: "auto" | "concise" | "detailed" | null; } export type ResponsesTextConfig = Record; export interface RemoteModelIdentity { provider: string; api: string; model: string; baseUrl?: string; } export interface RemoteCompactionData extends RemoteModelIdentity { implementation: "responses_compaction_v2"; replacementHistory: ResponseItem[]; } export interface RemoteCompactionResult { output: ResponseItem[]; } interface ContentLike { type?: string; text?: string; data?: string; mimeType?: string; source?: unknown; } const IMAGE_OMITTED = "image content omitted because you do not support image input"; const REMOTE_COMPACTION_FEATURE = "remote_compaction_v2"; const UUID_RE = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i; let fallbackInstallationId: string | undefined; export function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function isCompactionItem(value: unknown): value is Extract { return isRecord(value) && value.type === "compaction" && typeof value.encrypted_content === "string" && value.encrypted_content.trim().length > 0; } export function normalizeBaseUrl(value: unknown): string | undefined { if (typeof value !== "string" || !value.trim()) return undefined; try { const url = new URL(value); const path = url.pathname.replace(/\/+$/, ""); return `${url.protocol}//${url.host.toLowerCase()}${path}`; } catch { return value.trim().replace(/\/+$/, ""); } } function isResponsesModel( model: unknown, provider: string, api: string, ): model is Model { return isRecord(model) && model.provider === provider && model.api === api; } export function isDirectOpenAIResponsesModel(model: unknown): model is Model { return isResponsesModel(model, "openai", "openai-responses"); } export function isOpenAICodexResponsesModel(model: unknown): model is Model { return isResponsesModel(model, "openai-codex", "openai-codex-responses"); } export function supportsRemoteCompaction(model: unknown): model is Model { return isDirectOpenAIResponsesModel(model) || isOpenAICodexResponsesModel(model); } export function remoteModelIdentity(model: Model): RemoteModelIdentity { const baseUrl = normalizeBaseUrl(model.baseUrl); return { provider: String(model.provider), api: String(model.api), model: model.id, ...(baseUrl ? { baseUrl } : {}), }; } export function remoteIdentityMatches(identity: RemoteModelIdentity, model: Model): boolean { const current = remoteModelIdentity(model); return ( identity.provider === current.provider && identity.api === current.api && identity.model === current.model ); } function normalizeEndpointBase(value: unknown, fallback: string): string { return (normalizeBaseUrl(value) ?? fallback).replace(/\/+$/, ""); } export function remoteCompactionEndpointUrl(model: Model): string { if (isDirectOpenAIResponsesModel(model)) { const base = normalizeEndpointBase(model.baseUrl, "https://api.openai.com/v1"); if (base.endsWith("/responses")) return base; return base.endsWith("/v1") ? `${base}/responses` : `${base}/v1/responses`; } if (isOpenAICodexResponsesModel(model)) { const base = normalizeEndpointBase((model as Model).baseUrl, "https://chatgpt.com/backend-api"); if (base.endsWith("/codex/responses")) return base; if (base.endsWith("/codex")) return `${base}/responses`; return `${base}/codex/responses`; } throw new Error("Responses compaction v2 is not supported for this model."); } function resolveCodexInstallationId(): string { const codexHome = process.env.CODEX_HOME?.trim() || join(homedir(), ".codex"); const path = join(codexHome, "installation_id"); try { if (existsSync(path)) { const value = readFileSync(path, "utf8").trim(); if (UUID_RE.test(value)) return value.toLowerCase(); } } catch { // A stable id is useful but not required for local correctness. } fallbackInstallationId ??= randomUUID(); return fallbackInstallationId; } export function buildCodexIdentityHeaders(sessionId?: string): Record { const headers: Record = { "x-codex-installation-id": resolveCodexInstallationId(), }; if (sessionId) { headers["x-codex-window-id"] = `${sessionId}:0`; headers.session_id = sessionId; } return headers; } function extractCodexAccountId(token: string): string { const parts = token.split("."); if (parts.length !== 3) throw new Error("Failed to extract accountId from Codex token."); const payload = JSON.parse(Buffer.from(parts[1], "base64url").toString("utf8")) as unknown; if (!isRecord(payload)) throw new Error("Failed to extract accountId from Codex token."); const auth = payload["https://api.openai.com/auth"]; if (!isRecord(auth) || typeof auth.chatgpt_account_id !== "string" || !auth.chatgpt_account_id) { throw new Error("Failed to extract accountId from Codex token."); } return auth.chatgpt_account_id; } function addRemoteFeature(headers: Record): Record { const existing = Object.entries(headers).find( ([name]) => name.toLowerCase() === "x-codex-beta-features", )?.[1]; const features = new Set( (existing ?? "").split(",").map((value) => value.trim()).filter(Boolean), ); features.add(REMOTE_COMPACTION_FEATURE); const filtered = Object.fromEntries( Object.entries(headers).filter(([name]) => name.toLowerCase() !== "x-codex-beta-features"), ); return { ...filtered, "x-codex-beta-features": [...features].join(",") }; } export function buildRemoteCompactionHeaders(params: { model: Model; apiKey?: string; headers?: Record; sessionId?: string; }): Record { if (!supportsRemoteCompaction(params.model)) { throw new Error("Responses compaction v2 headers are not supported for this model."); } const common = addRemoteFeature({ ...(params.apiKey ? { authorization: `Bearer ${params.apiKey}` } : {}), ...buildCodexIdentityHeaders(params.sessionId), ...(params.headers ?? {}), accept: "text/event-stream", "content-type": "application/json", }); if (isDirectOpenAIResponsesModel(params.model)) return common; if (!params.apiKey) throw new Error("Codex remote compaction requires an access token."); return { ...common, "chatgpt-account-id": extractCodexAccountId(params.apiKey), originator: "pi", "user-agent": `pi-smart-compaction (${platform()} ${release()}; ${arch()})`, "OpenAI-Beta": "responses=experimental", }; } function parsePhase(value: unknown): AssistantPhase | undefined { if (typeof value !== "string" || !value) return undefined; try { const parsed = JSON.parse(value) as unknown; if (!isRecord(parsed)) return undefined; return parsed.phase === "commentary" || parsed.phase === "final_answer" ? parsed.phase : undefined; } catch { return undefined; } } function parseReasoning(value: unknown): ResponseItem | undefined { if (typeof value !== "string" || !value) return undefined; try { const parsed = JSON.parse(value) as unknown; if (!isRecord(parsed) || parsed.type !== "reasoning") return undefined; const summary = Array.isArray(parsed.summary) ? parsed.summary.flatMap((item) => isRecord(item) && typeof item.text === "string" ? [{ type: "summary_text" as const, text: item.text }] : [], ) : []; const content = Array.isArray(parsed.content) ? parsed.content.flatMap((item) => isRecord(item) && typeof item.text === "string" ? [{ type: item.type === "reasoning_text" ? "reasoning_text" as const : "text" as const, text: item.text, }] : [], ) : []; return { type: "reasoning", summary, ...(content.length > 0 ? { content } : {}), encrypted_content: typeof parsed.encrypted_content === "string" ? parsed.encrypted_content : null, }; } catch { return undefined; } } function responseContent(content: unknown, output: boolean): ResponseContentItem[] { if (typeof content === "string") { return content ? [{ type: output ? "output_text" : "input_text", text: content }] : []; } if (!Array.isArray(content)) return []; const result: ResponseContentItem[] = []; for (const value of content as ContentLike[]) { if (value.type === "text" && typeof value.text === "string") { result.push({ type: output ? "output_text" : "input_text", text: value.text }); } else if (value.type === "image" && typeof value.data === "string" && typeof value.mimeType === "string") { result.push({ type: "input_image", image_url: `data:${value.mimeType};base64,${value.data}` }); } } return result; } function toolOutput(content: unknown): unknown { if (typeof content === "string") return content; return responseContent(content, false); } function llmMessageToResponseItems(message: ReturnType[number]): ResponseItem[] { if (message.role === "user") { const content = responseContent(message.content, false); return content.length > 0 ? [{ type: "message", role: "user", content }] : []; } if (message.role === "toolResult") { return [{ type: "function_call_output", call_id: message.toolCallId.split("|", 1)[0], output: toolOutput(message.content), }]; } if (message.role !== "assistant") return []; const items: ResponseItem[] = []; let text = ""; let phase: AssistantPhase | undefined; const flush = () => { if (!text) return; items.push({ type: "message", role: "assistant", content: [{ type: "output_text", text }], ...(phase ? { phase } : {}), }); text = ""; }; for (const part of message.content) { if (part.type === "text") { const signedPhase = parsePhase(part.textSignature); if (signedPhase && phase && signedPhase !== phase) flush(); if (signedPhase) phase = signedPhase; text += part.text; } else if (part.type === "thinking") { flush(); const reasoning = parseReasoning(part.thinkingSignature); if (reasoning) items.push(reasoning); } else if (part.type === "toolCall") { flush(); items.push({ type: "function_call", name: part.name, arguments: JSON.stringify(part.arguments ?? {}), call_id: part.id.split("|", 1)[0], }); } } flush(); return items; } export function messagesToResponseItems(messages: AgentMessage[]): ResponseItem[] { return convertToLlm(messages).flatMap(llmMessageToResponseItems); } function cloneItem(item: T): T { return JSON.parse(JSON.stringify(item)) as T; } function callId(item: ResponseItem): string | undefined { const value = (item as Record).call_id; return typeof value === "string" && value ? value : undefined; } function normalizeCallPairs(items: ResponseItem[]): ResponseItem[] { const callIds = new Set( items.flatMap((item) => item.type === "function_call" && callId(item) ? [callId(item)!] : []), ); const paired = items.filter( (item) => item.type !== "function_call_output" || callIds.has(callId(item) ?? ""), ); const outputIds = new Set( paired.flatMap((item) => item.type === "function_call_output" && callId(item) ? [callId(item)!] : []), ); return paired.flatMap((item) => { const id = item.type === "function_call" ? callId(item) : undefined; return id && !outputIds.has(id) ? [item, { type: "function_call_output", call_id: id, output: "aborted" }] : [item]; }); } function stripImages(item: ResponseItem): ResponseItem { const next = cloneItem(item); if (next.type === "message" && Array.isArray(next.content)) { next.content = next.content.map((part) => part.type === "input_image" ? { type: "input_text", text: IMAGE_OMITTED } : part, ); } if (next.type === "function_call_output" && Array.isArray(next.output)) { next.output = next.output.map((part) => isRecord(part) && part.type === "input_image" ? { type: "input_text", text: IMAGE_OMITTED } : part, ); } return next; } export function normalizeResponseItemsForPrompt( items: ResponseItem[], model: Pick, "input">, ): ResponseItem[] { const paired = normalizeCallPairs(items.map(cloneItem)); return model.input.includes("image") ? paired : paired.map(stripImages); } function isRealUserMessage(item: ResponseItem): boolean { if (item.type !== "message" || item.role !== "user") return false; if (typeof item.content === "string") return item.content.trim().length > 0; return Array.isArray(item.content) && item.content.length > 0; } function messageText(item: ResponseItem): string { if (item.type !== "message") return ""; if (typeof item.content === "string") return item.content; if (!Array.isArray(item.content)) return ""; return item.content.flatMap((part: ResponseContentItem) => "text" in part ? [part.text] : []).join(""); } function retainRecentUsers(items: ResponseItem[], tokenBudget: number): ResponseItem[] { let remaining = Math.max(0, tokenBudget); const result: ResponseItem[] = []; for (const item of [...items].reverse()) { if (!isRealUserMessage(item) || remaining <= 0) continue; const tokens = Math.max(1, Math.ceil(messageText(item).length / 4)); if (tokens <= remaining) { result.push(cloneItem(item)); remaining -= tokens; continue; } if (item.type !== "message" || !Array.isArray(item.content)) continue; let chars = remaining * 4; const content = item.content.flatMap((part) => { if (part.type === "input_image") return [part]; if (chars <= 0) return []; const text = part.text.slice(0, chars); chars -= text.length; return text ? [{ ...part, text }] : []; }); if (content.length > 0) result.push({ ...cloneItem(item), content }); remaining = 0; } return result.reverse(); } export function buildRemoteCompactionHistory( input: ResponseItem[], compactionItem: ResponseItem, keepRecentTokens: number, ): ResponseItem[] { if (!isCompactionItem(compactionItem)) { throw new Error("Responses compaction v2 did not return a valid compaction item."); } return [ ...retainRecentUsers(input, keepRecentTokens), cloneItem(compactionItem), ]; } export function buildToolsPayload(allTools: ToolInfo[], activeToolNames: string[]): Record[] { const active = new Set(activeToolNames); return allTools.filter((tool) => active.has(tool.name)).map((tool) => ({ type: "function", name: tool.name, description: tool.description, parameters: tool.parameters, })); } export function thinkingLevelToReasoning(level: ThinkingLevel | undefined): ResponsesReasoningConfig | undefined { if (!level || level === "off") return { effort: "none", summary: "auto" }; return { effort: level, summary: "auto" }; } export function buildRemoteCompactionRequestBody(params: { model: Model; input: ResponseItem[]; instructions?: string; tools: Record[]; parallelToolCalls: boolean; reasoning?: ResponsesReasoningConfig; text?: ResponsesTextConfig; sessionId?: string; }): Record { return { model: params.model.id, input: [...params.input, { type: "compaction_trigger" }], instructions: params.instructions, tools: params.tools, parallel_tool_calls: params.parallelToolCalls, tool_choice: "auto", stream: true, store: false, include: ["reasoning.encrypted_content"], ...(params.sessionId ? { prompt_cache_key: params.sessionId } : {}), ...(params.reasoning ? { reasoning: params.reasoning } : {}), ...(params.text ? { text: params.text } : {}), }; } export function parseSseData(text: string): unknown[] { return text.replace(/\r\n/g, "\n").split("\n\n").flatMap((block) => { const data = block.split("\n") .filter((line) => line.startsWith("data:")) .map((line) => line.slice(5).trimStart()) .join("\n") .trim(); if (!data || data === "[DONE]") return []; try { return [JSON.parse(data) as unknown]; } catch { throw new Error("Responses compaction v2 returned malformed SSE data."); } }); } export function parseRemoteCompactionEvents(events: unknown[]): { compactionItem: ResponseItem; } { let completed = false; const compactionItems: ResponseItem[] = []; for (const event of events) { if (!isRecord(event)) continue; if (event.type === "error") { throw new Error(`Responses compaction v2 failed: ${String(event.message ?? "unknown error")}`); } if (event.type === "response.failed") { const response = isRecord(event.response) ? event.response : undefined; const error = response && isRecord(response.error) ? response.error : undefined; throw new Error(`Responses compaction v2 failed: ${String(error?.message ?? "response failed")}`); } if (event.type === "response.output_item.done" && isRecord(event.item) && event.item.type === "compaction") { compactionItems.push(event.item as ResponseItem); } if (event.type === "response.completed") completed = true; } if (!completed) throw new Error("Responses compaction v2 stream ended before response.completed."); if (compactionItems.length !== 1) { throw new Error(`Responses compaction v2 expected exactly one compaction item, got ${compactionItems.length}.`); } const compactionItem = compactionItems[0]; if (!isCompactionItem(compactionItem)) { throw new Error("Responses compaction v2 returned an invalid compaction item."); } return { compactionItem }; } export async function callRemoteCompaction(params: { model: Model; apiKey?: string; headers?: Record; sessionId?: string; input: ResponseItem[]; instructions?: string; tools: Record[]; parallelToolCalls: boolean; reasoning?: ResponsesReasoningConfig; text?: ResponsesTextConfig; keepRecentTokens: number; signal?: AbortSignal; fetchFn?: typeof fetch; }): Promise { if (params.signal?.aborted) { throw params.signal.reason instanceof Error ? params.signal.reason : new DOMException("The operation was aborted.", "AbortError"); } const response = await (params.fetchFn ?? fetch)(remoteCompactionEndpointUrl(params.model), { method: "POST", headers: buildRemoteCompactionHeaders(params), body: JSON.stringify(buildRemoteCompactionRequestBody(params)), signal: params.signal, }); if (!response.ok) { const text = await response.text().catch(() => ""); throw new Error(`Responses compaction v2 failed (${response.status}): ${text || response.statusText}`); } const parsed = parseRemoteCompactionEvents(parseSseData(await response.text())); return { output: buildRemoteCompactionHistory(params.input, parsed.compactionItem, params.keepRecentTokens), }; } export function toRemoteCompactionData( model: Model, result: RemoteCompactionResult, ): RemoteCompactionData { return { implementation: "responses_compaction_v2", ...remoteModelIdentity(model), replacementHistory: result.output.map(cloneItem), }; } export function isRemoteCompactionData(value: unknown): value is RemoteCompactionData { if (!isRecord(value)) return false; if (value.implementation !== "responses_compaction_v2") return false; if ( typeof value.provider !== "string" || typeof value.api !== "string" || typeof value.model !== "string" || !Array.isArray(value.replacementHistory) ) return false; const items = value.replacementHistory.filter( (item): item is ResponseItem => isRecord(item) && typeof item.type === "string", ); const compactions = items.filter((item) => item.type === "compaction"); return items.length === value.replacementHistory.length && compactions.length === 1 && isCompactionItem(compactions[0]); }