// @generated by scripts/build-runtime.mjs; do not edit. // @ts-nocheck -- the generated entry uses a .ts extension for Pi's Jiti loader. import { resolveCompactionRoute, terminalText } from "./chunks/chunk-BD6B7C5K.ts"; // src/codex-compact.ts import { buildContextEntries, buildSessionContext, convertToLlm, sessionEntryToContextMessages } from "@earendil-works/pi-coding-agent"; // src/checkpoint.ts import { createHash, randomUUID } from "node:crypto"; // src/protocol.ts var MAX_SSE_BYTES = 8 * 1024 * 1024; var MAX_COMPACT_JSON_BYTES = 8 * 1024 * 1024; var MAX_COMPACTION_ITEM_BYTES = 2 * 1024 * 1024; var CodexCompactionProtocolError = class extends Error { constructor(message) { super(message); this.name = "CodexCompactionProtocolError"; } }; function isObject(value) { return typeof value === "object" && value !== null && !Array.isArray(value); } function byteLength(value) { return Buffer.byteLength(JSON.stringify(value), "utf8"); } function isCompactionItem(value) { return isObject(value) && value.type === "compaction" && typeof value.encrypted_content === "string" && value.encrypted_content.length > 0; } function validateCompactionItem(value, maxBytes = MAX_COMPACTION_ITEM_BYTES) { if (!isCompactionItem(value)) { throw new CodexCompactionProtocolError("Remote response did not contain a valid compaction item"); } if (byteLength(value) > maxBytes) { throw new CodexCompactionProtocolError("Remote compaction item exceeded the size limit"); } return structuredClone(value); } function compactionItemsFromEvent(event) { const items = []; if (event.type === "response.output_item.done" && isObject(event.item)) { items.push(event.item); } if (event.type === "response.completed" && isObject(event.response)) { const output = event.response.output; if (Array.isArray(output)) items.push(...output); } return items.filter((item) => isObject(item) && item.type === "compaction"); } async function collectCompactionSse(stream, options = {}) { const maxBytes = options.maxBytes ?? MAX_SSE_BYTES; const reader = stream.getReader(); const onAbort = () => { void reader.cancel(new DOMException("Compaction aborted", "AbortError")).catch(() => void 0); }; options.signal?.addEventListener("abort", onAbort, { once: true }); const decoder = new TextDecoder(); let bytes = 0; let pending = ""; let dataLines = []; let completedResponse; const items = /* @__PURE__ */ new Map(); const checkAbort = () => { if (options.signal?.aborted) throw new DOMException("Compaction aborted", "AbortError"); }; const dispatch = () => { if (dataLines.length === 0) return; const data = dataLines.join("\n"); dataLines = []; if (data === "[DONE]") return; let parsed; try { parsed = JSON.parse(data); } catch { throw new CodexCompactionProtocolError("Remote compaction returned malformed SSE JSON"); } if (!isObject(parsed)) return; if (parsed.type === "response.completed") { completedResponse = isObject(parsed.response) ? parsed.response : {}; } for (const candidate of compactionItemsFromEvent(parsed)) { const item = validateCompactionItem(candidate, options.maxItemBytes); items.set(JSON.stringify(item), item); } }; const processLine = (line) => { if (line === "") { dispatch(); return; } if (line.startsWith(":")) return; if (line === "data") dataLines.push(""); else if (line.startsWith("data:")) dataLines.push(line.slice(5).replace(/^ /, "")); }; try { while (true) { checkAbort(); const { done, value } = await reader.read(); if (done) break; bytes += value.byteLength; if (bytes > maxBytes) { throw new CodexCompactionProtocolError("Remote compaction stream exceeded the size limit"); } pending += decoder.decode(value, { stream: true }); let newline = pending.indexOf("\n"); while (newline !== -1) { const rawLine = pending.slice(0, newline); pending = pending.slice(newline + 1); processLine(rawLine.endsWith("\r") ? rawLine.slice(0, -1) : rawLine); newline = pending.indexOf("\n"); } } pending += decoder.decode(); if (pending.length > 0) processLine(pending.endsWith("\r") ? pending.slice(0, -1) : pending); dispatch(); checkAbort(); } catch (error) { await reader.cancel(error).catch(() => void 0); throw error; } finally { options.signal?.removeEventListener("abort", onAbort); reader.releaseLock(); } if (!completedResponse) { throw new CodexCompactionProtocolError("Remote compaction stream ended without response.completed"); } if (items.size !== 1) { throw new CodexCompactionProtocolError( `Remote compaction returned ${items.size} distinct compaction items; expected exactly one` ); } return { item: [...items.values()][0], completedResponse }; } function isRetainedCompactContent(value) { if (!isObject(value)) return false; if (value.type === "input_text") return typeof value.text === "string"; if (value.type !== "input_image") return false; if (value.detail !== void 0 && value.detail !== null && value.detail !== "auto" && value.detail !== "low" && value.detail !== "high" && value.detail !== "original") { return false; } if (value.file_id !== void 0 && value.file_id !== null && typeof value.file_id !== "string" || value.image_url !== void 0 && value.image_url !== null && typeof value.image_url !== "string") { return false; } return typeof value.file_id === "string" && value.file_id.length > 0 || typeof value.image_url === "string" && value.image_url.length > 0; } function isRetainedCompactMessage(value) { return isObject(value) && value.role === "user" && (value.type === void 0 || value.type === "message") && Array.isArray(value.content) && value.content.length > 0 && value.content.every(isRetainedCompactContent); } function validateCompactedResponse(value, options = {}) { if (!isObject(value) || !Array.isArray(value.output)) { throw new CodexCompactionProtocolError("Responses Compact returned an invalid response object"); } const maxBytes = options.maxBytes ?? MAX_COMPACT_JSON_BYTES; const maxItemBytes = options.maxItemBytes ?? MAX_COMPACTION_ITEM_BYTES; if (byteLength(value) > maxBytes) { throw new CodexCompactionProtocolError("Responses Compact response exceeded the size limit"); } if (value.output.length === 0) { throw new CodexCompactionProtocolError("Responses Compact returned no output items"); } const output = value.output.map((item2) => { if (!isObject(item2)) { throw new CodexCompactionProtocolError("Responses Compact returned a non-object output item"); } if (byteLength(item2) > maxItemBytes) { throw new CodexCompactionProtocolError("Responses Compact output item exceeded the size limit"); } return structuredClone(item2); }); const compactionItems = output.filter((item2) => item2.type === "compaction"); if (compactionItems.length !== 1 || output.at(-1)?.type !== "compaction") { throw new CodexCompactionProtocolError( "Responses Compact must return retained messages followed by one compaction item" ); } for (const item2 of output.slice(0, -1)) { if (!isRetainedCompactMessage(item2)) { throw new CodexCompactionProtocolError("Responses Compact returned an unsupported retained output item"); } } const item = validateCompactionItem(output.at(-1), maxItemBytes); return { item, output: [...output.slice(0, -1), item], response: structuredClone(value) }; } async function collectCompactResponse(response, options = {}) { const maxBytes = options.maxBytes ?? MAX_COMPACT_JSON_BYTES; const declaredLength = Number(response.headers.get("content-length")); if (Number.isFinite(declaredLength) && declaredLength > maxBytes) { const error = new CodexCompactionProtocolError("Responses Compact response exceeded the size limit"); await response.body?.cancel(error).catch(() => void 0); throw error; } if (!response.body) { throw new CodexCompactionProtocolError("Responses Compact response did not contain a body"); } const reader = response.body.getReader(); const chunks = []; let bytes = 0; const onAbort = () => { void reader.cancel(new DOMException("Compaction aborted", "AbortError")).catch(() => void 0); }; options.signal?.addEventListener("abort", onAbort, { once: true }); try { while (true) { if (options.signal?.aborted) throw new DOMException("Compaction aborted", "AbortError"); const { done, value } = await reader.read(); if (done) break; bytes += value.byteLength; if (bytes > maxBytes) { throw new CodexCompactionProtocolError("Responses Compact response exceeded the size limit"); } chunks.push(value); } if (options.signal?.aborted) throw new DOMException("Compaction aborted", "AbortError"); } catch (error) { await reader.cancel(error).catch(() => void 0); throw error; } finally { options.signal?.removeEventListener("abort", onAbort); reader.releaseLock(); } const body = new Uint8Array(bytes); let offset = 0; for (const chunk of chunks) { body.set(chunk, offset); offset += chunk.byteLength; } let parsed; try { parsed = JSON.parse(new TextDecoder().decode(body)); } catch { throw new CodexCompactionProtocolError("Responses Compact returned malformed JSON"); } return validateCompactedResponse(parsed, options); } function markerTextFromItem(item) { if (!isObject(item) || item.role !== "user" || !Array.isArray(item.content)) return void 0; if (item.content.length !== 1) return void 0; const content = item.content[0]; if (!isObject(content) || content.type !== "input_text" || typeof content.text !== "string") { return void 0; } return content.text; } function rewriteCheckpointMarker(payload, marker, replacementHistory) { if (!isObject(payload) || !Array.isArray(payload.input)) { throw new CodexCompactionProtocolError("Codex Responses payload is missing an input array"); } const matches = payload.input.map((item, index2) => markerTextFromItem(item) === marker ? index2 : -1).filter((index2) => index2 >= 0); if (matches.length !== 1) { throw new CodexCompactionProtocolError( `Provider payload contained ${matches.length} checkpoint markers; expected exactly one` ); } const index = matches[0]; return { ...payload, input: [ ...payload.input.slice(0, index), ...structuredClone(replacementHistory), ...payload.input.slice(index + 1) ] }; } function appendCompactionTrigger(payload) { if (!isObject(payload) || !Array.isArray(payload.input)) { throw new CodexCompactionProtocolError("Codex Responses payload is missing an input array"); } if (payload.input.some((item) => isObject(item) && item.type === "compaction_trigger")) { throw new CodexCompactionProtocolError("Provider payload already contains a compaction trigger"); } return { ...payload, input: [...payload.input, { type: "compaction_trigger" }] }; } function expandRemoteCompactionPayload(payload, checkpoint) { if (checkpoint) { return rewriteCheckpointMarker(payload, checkpoint.marker, checkpoint.replacementHistory); } if (!isObject(payload) || !Array.isArray(payload.input)) { throw new CodexCompactionProtocolError("Responses payload is missing an input array"); } return structuredClone(payload); } function prepareRemoteCompactionPayload(payload, checkpoint) { return appendCompactionTrigger(expandRemoteCompactionPayload(payload, checkpoint)); } function hasCheckpointMarker(payload, marker) { return isObject(payload) && Array.isArray(payload.input) && payload.input.some((item) => markerTextFromItem(item) === marker); } // src/checkpoint.ts var CHECKPOINT_KIND = "pi-codex-remote-compaction"; var CHECKPOINT_VERSION = 3; var REPLACEMENT_TOKEN_BUDGET = 64e3; var REPLACEMENT_BYTE_BUDGET = 8 * 1024 * 1024; var MAX_MEDIA_ITEM_BYTES = 2 * 1024 * 1024; var MAX_CHECKPOINT_DETAILS_BYTES = 10 * 1024 * 1024; var MAX_CHECKPOINT_ID_LENGTH = 128; var MAX_PROVIDER_ID_LENGTH = 256; var MAX_API_ID_LENGTH = 256; var MAX_MODEL_ID_LENGTH = 512; var MAX_KEPT_FINGERPRINTS = 1e5; function isObject2(value) { return typeof value === "object" && value !== null && !Array.isArray(value); } function stableValue(value) { if (Array.isArray(value)) return value.map(stableValue); if (!isObject2(value)) return value; return Object.fromEntries( Object.entries(value).sort(([left], [right]) => left.localeCompare(right)).map(([key, child]) => [key, stableValue(child)]) ); } function serializedBytes(value) { return Buffer.byteLength(JSON.stringify(value), "utf8"); } function fingerprintMessage(message) { return createHash("sha256").update(JSON.stringify(stableValue(message))).digest("hex"); } function checkpointMarker(checkpointId) { return [ `[PI_CODEX_REMOTE_CHECKPOINT:${checkpointId}]`, "Opaque checkpoint injection failed. Do not infer missing history; tell the user to re-enable", "@narumitw/pi-codex-compact with the same model and Responses API." ].join(" "); } function fallbackSummary(checkpointId) { return [ `Responses compaction checkpoint ${checkpointId} stores the older history opaquely.`, "Full replay requires @narumitw/pi-codex-compact and the same model through a compatible Responses provider.", "Without them, only Pi's retained recent messages remain available." ].join(" "); } function markerMessage(checkpointId, timestamp) { return { role: "user", content: [{ type: "text", text: checkpointMarker(checkpointId) }], timestamp }; } function parseCheckpointDetails(value) { if (!isObject2(value)) return void 0; try { if (serializedBytes(value) > MAX_CHECKPOINT_DETAILS_BYTES) return void 0; } catch { return void 0; } const isVersionOne = value.version === 1 && value.api === "openai-codex-responses" && value.protocol === "remote-compaction-v2"; const isVersionTwo = value.version === 2 && (value.api === "openai-codex-responses" || value.api === "openai-responses" || value.api === "azure-openai-responses") && (value.protocol === "remote-v2" || value.protocol === "responses-compact"); const isVersionThree = value.version === CHECKPOINT_VERSION && typeof value.api === "string" && value.api.length > 0 && value.api.length <= MAX_API_ID_LENGTH && (value.profile === "codex-responses-v1" || value.profile === "openai-responses-v1") && (value.protocol === "remote-v2" || value.protocol === "responses-compact"); if (value.kind !== CHECKPOINT_KIND || !isVersionOne && !isVersionTwo && !isVersionThree || typeof value.checkpointId !== "string" || value.checkpointId.length < 8 || value.checkpointId.length > MAX_CHECKPOINT_ID_LENGTH || typeof value.provider !== "string" || value.provider.length === 0 || value.provider.length > MAX_PROVIDER_ID_LENGTH || typeof value.modelId !== "string" || value.modelId.length === 0 || value.modelId.length > MAX_MODEL_ID_LENGTH || !Array.isArray(value.replacementHistory) || !Array.isArray(value.keptMessageFingerprints) || value.keptMessageFingerprints.length > MAX_KEPT_FINGERPRINTS || typeof value.createdAt !== "string" || value.createdAt.length > 64) { return void 0; } const api = value.api; const profile = api === "openai-codex-responses" ? "codex-responses-v1" : api === "openai-responses" || api === "azure-openai-responses" ? "openai-responses-v1" : value.profile; if (profile !== "codex-responses-v1" && profile !== "openai-responses-v1" || api !== "openai-codex-responses" && api !== "openai-responses" && api !== "azure-openai-responses" && profile !== "codex-responses-v1" || isVersionThree && value.profile !== profile) { return void 0; } if (value.replacementHistory.length === 0 || !value.replacementHistory.every(isObject2) || !value.keptMessageFingerprints.every( (fingerprint) => typeof fingerprint === "string" && /^[a-f0-9]{64}$/.test(fingerprint) ) || serializedBytes(value.replacementHistory) > REPLACEMENT_BYTE_BUDGET) { return void 0; } const last = value.replacementHistory.at(-1); try { validateCompactionItem(last); } catch { return void 0; } return { kind: CHECKPOINT_KIND, version: CHECKPOINT_VERSION, checkpointId: value.checkpointId, provider: value.provider, api, profile, modelId: value.modelId, protocol: isVersionOne ? "remote-v2" : value.protocol, replacementHistory: structuredClone(value.replacementHistory), keptMessageFingerprints: [...value.keptMessageFingerprints], createdAt: value.createdAt }; } function latestCheckpoint(entries) { for (let index = entries.length - 1; index >= 0; index--) { const entry = entries[index]; if (entry.type !== "compaction") continue; const details = parseCheckpointDetails(entry.details); return details ? { entry, details } : void 0; } return void 0; } function isOlderCompactionSummary(message, timestamp) { return message.role === "compactionSummary" && Number.isFinite(message.timestamp) && Number.isFinite(timestamp) && message.timestamp < timestamp; } function projectCheckpointContext(messages, details, checkpointSummary) { const summaryIndex = messages.findIndex( (message) => message.role === "compactionSummary" && message.summary === checkpointSummary ); if (summaryIndex < 0) return void 0; const timestamp = messages[summaryIndex].timestamp; let messageIndex = summaryIndex + 1; let fingerprintIndex = 0; while (fingerprintIndex < details.keptMessageFingerprints.length) { if (messageIndex >= messages.length) return void 0; const message = messages[messageIndex]; if (fingerprintMessage(message) === details.keptMessageFingerprints[fingerprintIndex]) { messageIndex += 1; fingerprintIndex += 1; continue; } if (isOlderCompactionSummary(message, timestamp)) { messageIndex += 1; continue; } return void 0; } while (messageIndex < messages.length && isOlderCompactionSummary(messages[messageIndex], timestamp)) { messageIndex += 1; } return [ ...messages.slice(0, summaryIndex), markerMessage(details.checkpointId, timestamp), ...messages.slice(messageIndex) ]; } function rawText(item) { if (!Array.isArray(item.content)) return ""; return item.content.flatMap( (part) => isObject2(part) && typeof part.text === "string" && part.type === "input_text" ? [part.text] : [] ).join("\n"); } function hasMedia(item) { return Array.isArray(item.content) && item.content.some((part) => isObject2(part) && part.type === "input_image"); } function truncateTextItem(item, maxChars) { if (!Array.isArray(item.content) || maxChars <= 32) return void 0; let remaining = maxChars - 16; const content = [...item.content].reverse().flatMap((part) => { if (!isObject2(part) || part.type !== "input_text" || typeof part.text !== "string" || remaining <= 0) { return []; } const text = part.text.slice(-remaining); remaining -= text.length; return [{ ...part, text: `[truncated] ${text}` }]; }); if (content.length === 0) return void 0; return { ...item, content: content.reverse() }; } function buildReplacementHistory(input, compactionItem, options = {}) { const tokenBudget = options.tokenBudget ?? REPLACEMENT_TOKEN_BUDGET; const byteBudget = options.byteBudget ?? REPLACEMENT_BYTE_BUDGET; const opaque = validateCompactionItem(compactionItem); let remainingBytes = byteBudget - serializedBytes(opaque); let remainingChars = tokenBudget * 4; if (remainingBytes <= 0) throw new Error("Opaque compaction item exceeds replacement history budget"); const retainedNewestFirst = []; const candidates = input.filter( (item) => isObject2(item) && item.role === "user" && item.type !== "compaction_trigger" ); for (let index = candidates.length - 1; index >= 0; index--) { const candidate = candidates[index]; const bytes = serializedBytes(candidate); if (hasMedia(candidate) && bytes > MAX_MEDIA_ITEM_BYTES) continue; const text = rawText(candidate); let retained = candidate; if (text.length > remainingChars) { if (hasMedia(candidate)) continue; const truncated = truncateTextItem(candidate, remainingChars); if (!truncated) continue; retained = truncated; } if (serializedBytes(retained) > remainingBytes) { if (hasMedia(retained)) continue; const maxCharsByBytes = Math.max(0, remainingBytes - 128); const truncated = truncateTextItem(retained, Math.min(remainingChars, maxCharsByBytes)); if (!truncated || serializedBytes(truncated) > remainingBytes) continue; retained = truncated; } retainedNewestFirst.push(structuredClone(retained)); remainingBytes -= serializedBytes(retained); remainingChars -= Math.min(remainingChars, rawText(retained).length); if (remainingBytes <= 128 || remainingChars <= 32) break; } return [...retainedNewestFirst.reverse(), opaque]; } function createCheckpointDetails(input) { const details = { kind: CHECKPOINT_KIND, version: CHECKPOINT_VERSION, checkpointId: input.checkpointId ?? randomUUID(), provider: input.provider, api: input.api, profile: input.profile, modelId: input.modelId, protocol: input.protocol, replacementHistory: structuredClone(input.replacementHistory), keptMessageFingerprints: input.keptMessages.map(fingerprintMessage), createdAt: input.createdAt ?? (/* @__PURE__ */ new Date()).toISOString() }; const parsed = parseCheckpointDetails(details); if (!parsed) throw new Error("Created an invalid Codex checkpoint"); return parsed; } // src/remote-compact.ts import { normalizeContext } from "@earendil-works/pi-ai"; // src/remote-types.ts function isJsonObject(value) { return typeof value === "object" && value !== null && !Array.isArray(value); } function abortError() { return new DOMException("Compaction aborted", "AbortError"); } function assertPreparedInput(payload) { if (!Array.isArray(payload.input) || !payload.input.every(isJsonObject)) { throw new Error("Prepared compaction payload has invalid input items"); } return structuredClone(payload.input); } // src/remote-shared.ts async function collectProviderUsage(stream, signal) { let usage; for await (const event of stream) { if (signal.aborted) throw abortError(); if (event.type === "error") { throw new Error(event.error.errorMessage ?? "Responses compaction request failed"); } if (event.type === "done") usage = event.message.usage; } if (signal.aborted) throw abortError(); if (!usage) throw new Error("Responses provider stream ended without completion usage"); return usage; } // src/remote-compact.ts var OFFICIAL_COMPACT_FIELDS = [ "model", "input", "instructions", "previous_response_id", "prompt_cache_key", "prompt_cache_retention", "service_tier" ]; var CODEX_COMPACT_FIELDS = [ "model", "input", "instructions", "tools", "parallel_tool_calls", "reasoning", "service_tier", "prompt_cache_key", "text", "access_programs" ]; function compactPayload(payload, profile) { const fields = profile === "codex-responses-v1" ? CODEX_COMPACT_FIELDS : OFFICIAL_COMPACT_FIELDS; const result = {}; for (const field of fields) { if (Object.hasOwn(payload, field) && payload[field] !== void 0) { result[field] = structuredClone(payload[field]); } } if (typeof result.model !== "string" || result.model.length === 0) { throw new CodexCompactionProtocolError("Responses payload is missing a model"); } assertPreparedInput(result); return result; } function requestUrl(input) { return new URL(input instanceof Request ? input.url : String(input)); } function responsesCompactUrl(input) { const original = requestUrl(input); if (!original.pathname.endsWith("/responses")) { throw new CodexCompactionProtocolError("Provider request URL does not end with the Responses endpoint"); } const compact = new URL(original); compact.pathname = `${compact.pathname}/compact`; if (compact.origin !== original.origin) { throw new CodexCompactionProtocolError("Responses Compact URL changed origin"); } return compact; } function mergedHeaders(input, init) { const headers = new Headers(input instanceof Request ? input.headers : void 0); new Headers(init?.headers).forEach((value, name) => { headers.set(name, value); }); headers.delete("content-encoding"); headers.delete("content-length"); headers.set("accept", "application/json"); headers.set("content-type", "application/json"); return headers; } function mergedSignal(input, init, ownerSignal) { const signals = [ownerSignal]; if (input instanceof Request) signals.push(input.signal); if (init?.signal) signals.push(init.signal); return signals.length === 1 ? ownerSignal : AbortSignal.any(signals); } function nonRetryableBridgeFailure(error) { const message = error instanceof Error ? error.message : String(error); return Response.json( { error: { message, type: "invalid_request_error", code: "invalid_compact_response" } }, { status: 400 } ); } function nonNegativeInteger(value) { return typeof value === "number" && Number.isSafeInteger(value) && value >= 0; } function optionalUsageDetail(details, field) { if (details === void 0 || details === null) return 0; if (!isJsonObject(details)) { throw new CodexCompactionProtocolError("Responses Compact response has invalid usage details"); } const value = details[field]; if (value === void 0 || value === null) return 0; if (!nonNegativeInteger(value)) { throw new CodexCompactionProtocolError(`Responses Compact response has invalid usage detail ${field}`); } return value; } function validatedUsage(response) { const usage = response.usage; if (!isJsonObject(usage)) { throw new CodexCompactionProtocolError("Responses Compact response is missing usage"); } for (const field of ["input_tokens", "output_tokens", "total_tokens"]) { if (!nonNegativeInteger(usage[field])) { throw new CodexCompactionProtocolError(`Responses Compact response has invalid ${field}`); } } const inputTokens = usage.input_tokens; const outputTokens = usage.output_tokens; const totalTokens = usage.total_tokens; const cachedTokens = optionalUsageDetail(usage.input_tokens_details, "cached_tokens"); const cacheWriteTokens = optionalUsageDetail(usage.input_tokens_details, "cache_write_tokens"); const reasoningTokens = optionalUsageDetail(usage.output_tokens_details, "reasoning_tokens"); if (cachedTokens + cacheWriteTokens > inputTokens || reasoningTokens > outputTokens) { throw new CodexCompactionProtocolError("Responses Compact response has inconsistent usage details"); } if (totalTokens !== inputTokens + outputTokens) { throw new CodexCompactionProtocolError("Responses Compact response has inconsistent total usage"); } return structuredClone(usage); } function syntheticCompletion(result, payload) { const completed = { id: typeof result.response.id === "string" ? result.response.id : "resp_pi_compact_bridge", object: "response", created_at: typeof result.response.created_at === "number" ? result.response.created_at : Math.floor(Date.now() / 1e3), status: "completed", model: payload.model, output: [], parallel_tool_calls: false, tool_choice: "auto", tools: [], usage: validatedUsage(result.response) }; const events = [ { type: "response.created", response: { ...completed, status: "in_progress" } }, { type: "response.completed", response: completed } ]; return new Response(events.map((event) => `data: ${JSON.stringify(event)} `).join(""), { status: 200, headers: { "content-type": "text/event-stream" } }); } async function requestResponsesCompact(request) { if (request.signal.aborted) throw abortError(); let preparedPayload; let sentInput; let compactResult; let bridgeError; let dispatchInFlight = false; let successfulResponses = 0; const baseFetch = request.fetch ?? globalThis.fetch; const bridgeFetch = async (input, init) => { if (request.signal.aborted) throw abortError(); if (successfulResponses > 0) { bridgeError = new CodexCompactionProtocolError("Provider dispatched again after Responses Compact succeeded"); return nonRetryableBridgeFailure(bridgeError); } if (dispatchInFlight) { bridgeError = new CodexCompactionProtocolError("Provider dispatched overlapping Responses Compact requests"); return nonRetryableBridgeFailure(bridgeError); } if (!preparedPayload) { bridgeError = new CodexCompactionProtocolError("Provider dispatched before exposing its request payload"); return nonRetryableBridgeFailure(bridgeError); } let compactUrl; try { compactUrl = responsesCompactUrl(input); } catch (error) { bridgeError = error; return nonRetryableBridgeFailure(error); } const signal = mergedSignal(input, init, request.signal); dispatchInFlight = true; try { const response = await baseFetch(compactUrl, { ...init, method: "POST", headers: mergedHeaders(input, init), body: JSON.stringify(preparedPayload), signal }); if (!response.ok) return response; try { const result = await collectCompactResponse(response, { signal }); successfulResponses += 1; if (successfulResponses !== 1) { throw new CodexCompactionProtocolError( "Provider returned more than one successful Responses Compact response" ); } compactResult = result; return syntheticCompletion(result, preparedPayload); } catch (error) { bridgeError = error; return nonRetryableBridgeFailure(error); } } finally { dispatchInFlight = false; } }; const stream = request.provider.stream(request.model, normalizeContext(request.context), { apiKey: request.apiKey, headers: request.headers, env: request.env, signal: request.signal, transport: "sse", cacheRetention: "none", timeoutMs: request.requestTimeoutMs ?? 5 * 60 * 1e3, maxRetries: request.maxRetries ?? 2, fetch: bridgeFetch, onPayload: (payload) => { if (preparedPayload) { throw new CodexCompactionProtocolError("Provider exposed more than one compaction request payload"); } const expanded = expandRemoteCompactionPayload(payload, request.priorCheckpoint); preparedPayload = compactPayload(expanded, request.profile); sentInput = assertPreparedInput(preparedPayload); return expanded; } }); let usage; try { usage = await collectProviderUsage(stream, request.signal); } catch (error) { if (bridgeError) throw bridgeError; throw error; } if (request.signal.aborted) throw abortError(); if (bridgeError) throw bridgeError; if (!preparedPayload || !sentInput || !compactResult || successfulResponses !== 1) { throw new CodexCompactionProtocolError("Provider did not complete exactly one Responses Compact request"); } return { item: compactResult.item, promptInput: sentInput, compactedOutput: compactResult.output, usage }; } // src/remote-v2.ts import { normalizeContext as normalizeContext2 } from "@earendil-works/pi-ai"; async function requestRemoteCompactionV2(request) { if (request.signal.aborted) throw abortError(); let sentInput; const inspections = []; const baseFetch = request.fetch ?? globalThis.fetch; const inspectedFetch = async (input, init) => { const response = await baseFetch(input, init); if (!response.ok || !response.body) return response; const [providerBody, inspectionBody] = response.body.tee(); const inspection2 = collectCompactionSse(inspectionBody, { signal: request.signal }).then( (value) => ({ ok: true, value }), (error) => ({ ok: false, error }) ); inspections.push(inspection2); return new Response(providerBody, { status: response.status, statusText: response.statusText, headers: response.headers }); }; const stream = request.provider.stream(request.model, normalizeContext2(request.context), { apiKey: request.apiKey, headers: request.headers, env: request.env, signal: request.signal, transport: "sse", cacheRetention: "none", timeoutMs: request.requestTimeoutMs ?? 5 * 60 * 1e3, maxRetries: request.maxRetries ?? 2, fetch: inspectedFetch, onPayload: (payload) => { const prepared = prepareRemoteCompactionPayload(payload, request.priorCheckpoint); sentInput = assertPreparedInput(prepared).slice(0, -1); return prepared; } }); const usage = await collectProviderUsage(stream, request.signal); if (!sentInput) { throw new CodexCompactionProtocolError("Provider did not expose a request payload"); } if (inspections.length !== 1) { throw new CodexCompactionProtocolError( `Provider exposed ${inspections.length} successful SSE responses; expected exactly one` ); } const inspection = await inspections[0]; if (request.signal.aborted) throw abortError(); if (!inspection.ok) throw inspection.error; return { item: inspection.value.item, promptInput: sentInput, usage }; } // src/remote.ts function requestRemoteCompaction(request) { return request.protocol === "responses-compact" ? requestResponsesCompact(request) : requestRemoteCompactionV2(request); } // src/settings.ts import { randomUUID as randomUUID2 } from "node:crypto"; import { constants } from "node:fs"; import { mkdir, open, rename, rm, writeFile } from "node:fs/promises"; import { basename, dirname, join } from "node:path"; import { getAgentDir } from "@earendil-works/pi-coding-agent"; var CODEX_COMPACT_SETTINGS_FILE = "pi-codex-compact.json"; var MAX_SETTINGS_BYTES = 64 * 1024; var DEFAULT_CODEX_COMPACT_SETTINGS = Object.freeze({ enabled: true, protocol: "auto", apiProfiles: {}, requestTimeoutMs: 3e5, maxRetries: 2, replacementTokenBudget: 64e3, notifyOnFallback: true }); var LIMITS = Object.freeze({ requestTimeoutMs: { minimum: 3e4, maximum: 6e5 }, maxRetries: { minimum: 0, maximum: 2 }, replacementTokenBudget: { minimum: 8e3, maximum: 128e3 } }); var BUILT_IN_APIS = /* @__PURE__ */ new Set([ "openai-completions", "mistral-conversations", "openai-codex-responses", "openai-responses", "azure-openai-responses", "anthropic-messages", "bedrock-converse-stream", "google-generative-ai", "google-vertex", "pi-messages" ]); var MAX_API_PROFILE_ID_LENGTH = 256; function isRecord(value) { return typeof value === "object" && value !== null && !Array.isArray(value); } function validInteger(value, minimum, maximum) { return typeof value === "number" && Number.isSafeInteger(value) && value >= minimum && value <= maximum; } function hasWhitespaceOrControl(value) { return [...value].some((character) => { const code = character.codePointAt(0) ?? 0; return code <= 32 || code === 127; }); } function normalizeApiProfiles(value) { if (!isRecord(value) || Object.getPrototypeOf(value) !== Object.prototype && Object.getPrototypeOf(value) !== null) { return void 0; } const profiles = {}; for (const [api, profile] of Object.entries(value)) { if (api.length === 0 || api.length > MAX_API_PROFILE_ID_LENGTH || api.trim() !== api || hasWhitespaceOrControl(api) || BUILT_IN_APIS.has(api) || api === "__proto__" || api === "constructor" || api === "prototype" || profile !== "codex-responses-v1") { return void 0; } profiles[api] = profile; } return profiles; } function normalizeCodexCompactSettings(value) { if (!isRecord(value)) return void 0; if (Object.hasOwn(value, "enabled") && typeof value.enabled !== "boolean") return void 0; if (Object.hasOwn(value, "protocol") && value.protocol !== "auto" && value.protocol !== "remote-v2" && value.protocol !== "responses-compact") { return void 0; } if (Object.hasOwn(value, "notifyOnFallback") && typeof value.notifyOnFallback !== "boolean") { return void 0; } if (Object.hasOwn(value, "apiProfiles") && normalizeApiProfiles(value.apiProfiles) === void 0) return void 0; for (const [field, limits] of Object.entries(LIMITS)) { if (Object.hasOwn(value, field) && !validInteger(value[field], limits.minimum, limits.maximum)) { return void 0; } } return { enabled: typeof value.enabled === "boolean" ? value.enabled : DEFAULT_CODEX_COMPACT_SETTINGS.enabled, protocol: value.protocol === "remote-v2" || value.protocol === "responses-compact" ? value.protocol : DEFAULT_CODEX_COMPACT_SETTINGS.protocol, apiProfiles: normalizeApiProfiles(value.apiProfiles) ?? structuredClone(DEFAULT_CODEX_COMPACT_SETTINGS.apiProfiles), requestTimeoutMs: typeof value.requestTimeoutMs === "number" ? value.requestTimeoutMs : DEFAULT_CODEX_COMPACT_SETTINGS.requestTimeoutMs, maxRetries: typeof value.maxRetries === "number" ? value.maxRetries : DEFAULT_CODEX_COMPACT_SETTINGS.maxRetries, replacementTokenBudget: typeof value.replacementTokenBudget === "number" ? value.replacementTokenBudget : DEFAULT_CODEX_COMPACT_SETTINGS.replacementTokenBudget, notifyOnFallback: typeof value.notifyOnFallback === "boolean" ? value.notifyOnFallback : DEFAULT_CODEX_COMPACT_SETTINGS.notifyOnFallback }; } function codexCompactSettingsPath() { return join(getAgentDir(), CODEX_COMPACT_SETTINGS_FILE); } function aborted(signal) { if (signal?.aborted) throw new DOMException("Settings operation aborted", "AbortError"); } async function loadCodexCompactSettings(path = codexCompactSettingsPath(), signal) { aborted(signal); try { const handle = await open(path, constants.O_RDONLY | constants.O_NOFOLLOW); let text; try { const stats = await handle.stat(); aborted(signal); if (!stats.isFile()) throw new Error("settings path is not a regular file"); if (stats.size > MAX_SETTINGS_BYTES) throw new Error("settings file exceeds 64 KiB"); text = await handle.readFile("utf8"); } finally { await handle.close(); } aborted(signal); const document = JSON.parse(text); const settings = normalizeCodexCompactSettings(document); if (!settings || !isRecord(document)) throw new Error("invalid settings shape or bounds"); return { kind: "loaded", path, settings, document }; } catch (error) { if (signal?.aborted) throw error; if (isNodeError(error) && error.code === "ENOENT") { return { kind: "missing", path, settings: { ...DEFAULT_CODEX_COMPACT_SETTINGS }, document: {} }; } return { kind: "invalid", path, settings: { ...DEFAULT_CODEX_COMPACT_SETTINGS }, issue: isNodeError(error) && error.code === "ELOOP" ? "symbolic links are not accepted" : error instanceof Error ? error.message : String(error) }; } } async function savePatch(path, patch, signal) { const latest = await loadCodexCompactSettings(path, signal); if (latest.kind === "invalid") { throw new Error("Cannot overwrite an invalid pi-codex-compact.json; repair it and reload first"); } const document = { ...latest.document, ...patch }; const settings = normalizeCodexCompactSettings(document); if (!settings) throw new Error("Refusing to save invalid Codex compaction settings"); const temporaryPath = join(dirname(path), `.${basename(path)}.${randomUUID2()}.tmp`); await mkdir(dirname(path), { recursive: true }); aborted(signal); try { await writeFile(temporaryPath, `${JSON.stringify(document, null, 2)} `, { encoding: "utf8", flag: "wx", mode: 384 }); aborted(signal); const current = await loadCodexCompactSettings(path, signal); if (current.kind === "invalid" || current.kind !== latest.kind || JSON.stringify(current.document) !== JSON.stringify(latest.document)) { throw new Error("pi-codex-compact.json changed while saving; reopen settings and retry"); } await rename(temporaryPath, path); } finally { await rm(temporaryPath, { force: true }).catch(() => void 0); } return { kind: "loaded", path, settings, document }; } function createCodexCompactSettingsRuntime(path = codexCompactSettingsPath()) { let state = { kind: "missing", path, settings: { ...DEFAULT_CODEX_COMPACT_SETTINGS }, document: {} }; let queue = Promise.resolve(); const enqueue = (operation) => { const result = queue.then(operation, operation); queue = result.then( () => void 0, () => void 0 ); return result; }; return { get: () => structuredClone(state), reload: (signal) => enqueue(async () => { state = await loadCodexCompactSettings(path, signal); return structuredClone(state); }), update: (patch, signal) => enqueue(async () => { state = await savePatch(path, patch, signal); return structuredClone(state); }), flush: () => queue }; } function isNodeError(error) { return error instanceof Error && "code" in error; } // src/codex-compact.ts var STATUS_KEY = "codex-compact"; function activeCheckpoint(ctx) { return latestCheckpoint(ctx.sessionManager.getBranch()); } function isCheckpointCompatible(details, model, settings) { const route = resolveCompactionRoute(model, settings); return route.kind === "remote" && model !== void 0 && route.api === details.api && route.profile === details.profile && model.id === details.modelId; } function keptMessages(event) { const leafId = event.branchEntries.at(-1)?.id ?? null; const contextEntries = buildContextEntries(event.branchEntries, leafId); const keptIndex = contextEntries.findIndex((entry) => entry.id === event.preparation.firstKeptEntryId); if (keptIndex < 0) { throw new Error("Pi compaction cut point is not present in the active context"); } return contextEntries.slice(keptIndex).flatMap(sessionEntryToContextMessages); } function activeTools(pi) { const available = new Map(pi.getAllTools().map((tool) => [tool.name, tool])); return pi.getActiveTools().flatMap((name) => { const tool = available.get(name); return tool ? [ { name: tool.name, description: tool.description, parameters: tool.parameters } ] : []; }); } function projectedCurrentMessages(event, model, route) { const leafId = event.branchEntries.at(-1)?.id ?? null; const session = buildSessionContext(event.branchEntries, leafId); const prior = latestCheckpoint(event.branchEntries); if (!prior) return { messages: session.messages }; if (prior.details.api !== route.api || prior.details.profile !== route.profile || prior.details.modelId !== model.id) { throw new Error("The active opaque checkpoint belongs to a different Responses model"); } const projected = projectCheckpointContext(session.messages, prior.details, prior.entry.summary); if (!projected) { throw new Error("The previous opaque checkpoint could not be projected safely"); } return { messages: projected, prior: prior.details }; } function notifyFailure(ctx, error, settings) { if (!ctx.hasUI || !settings.notifyOnFallback) return; const message = terminalText(error instanceof Error ? error.message : String(error)); ctx.ui.notify(`Responses compaction failed; using Pi compaction. ${message}`, "warning"); } function sessionStillOwned(ctx, sessionId, signal) { return !signal.aborted && ctx.sessionManager.getSessionId() === sessionId; } async function compactRemotely(pi, event, ctx, settings, ownerSignal, fetch) { const model = ctx.model; const route = resolveCompactionRoute(model, settings); if (route.kind === "native" || !model) return void 0; const signal = AbortSignal.any([event.signal, ownerSignal]); if (signal.aborted) return { cancel: true }; const sessionId = ctx.sessionManager.getSessionId(); ctx.ui.setStatus(STATUS_KEY, route.protocol === "remote-v2" ? "Responses Remote V2\u2026" : "Responses Compact API\u2026"); try { const auth = await ctx.modelRegistry.getApiKeyAndHeaders(model); if (!sessionStillOwned(ctx, sessionId, signal)) return { cancel: true }; if (!auth.ok) throw new Error(auth.error); const provider = ctx.modelRegistry.getProvider(model.provider); if (!provider) throw new Error("The active Responses provider is unavailable"); const current = projectedCurrentMessages(event, model, route); const context = { systemPrompt: ctx.getSystemPrompt(), messages: convertToLlm(current.messages), tools: activeTools(pi) }; const response = await requestRemoteCompaction({ provider, model, context, protocol: route.protocol, profile: route.profile, apiKey: auth.apiKey, headers: auth.headers, env: auth.env, signal, priorCheckpoint: current.prior ? { marker: checkpointMarker(current.prior.checkpointId), replacementHistory: current.prior.replacementHistory } : void 0, requestTimeoutMs: settings.requestTimeoutMs, maxRetries: settings.maxRetries, fetch }); if (!sessionStillOwned(ctx, sessionId, signal)) return { cancel: true }; const replacementHistory = buildReplacementHistory( response.compactedOutput?.slice(0, -1) ?? response.promptInput, response.item, { tokenBudget: settings.replacementTokenBudget } ); const details = createCheckpointDetails({ provider: model.provider, api: route.api, profile: route.profile, modelId: model.id, protocol: route.protocol, replacementHistory, keptMessages: keptMessages(event) }); return { compaction: { summary: fallbackSummary(details.checkpointId), firstKeptEntryId: event.preparation.firstKeptEntryId, tokensBefore: event.preparation.tokensBefore, usage: response.usage, details } }; } catch (error) { if (signal.aborted || ctx.sessionManager.getSessionId() !== sessionId) { return { cancel: true }; } notifyFailure(ctx, error, settings); return void 0; } finally { if (ctx.sessionManager.getSessionId() === sessionId) ctx.ui.setStatus(STATUS_KEY, void 0); } } function createCodexCompactExtension(options = {}) { return (pi) => { const providerWarnings = /* @__PURE__ */ new Set(); const settingsRuntime = options.settingsRuntime ?? createCodexCompactSettingsRuntime(); let sessionController = new AbortController(); let generation = 0; pi.registerCommand("codex-compact", { description: "Compact now or configure Responses compaction", handler: async (args, ctx) => { if (args.trim()) throw new Error("Usage: /codex-compact"); const ownerGeneration = generation; const controller = sessionController; const { showCodexCompactMenu } = await import("./chunks/settings-menu-SZZ437BH.ts"); if (ownerGeneration !== generation || controller.signal.aborted) return; await showCodexCompactMenu(settingsRuntime, ctx, { signal: controller.signal, isCurrent: () => ownerGeneration === generation && !controller.signal.aborted }); } }); pi.on("session_start", async (_event, ctx) => { sessionController.abort(); sessionController = new AbortController(); generation += 1; const ownerGeneration = generation; const sessionId = ctx.sessionManager.getSessionId(); providerWarnings.clear(); let state; try { state = await settingsRuntime.reload(sessionController.signal); } catch (error) { if (sessionController.signal.aborted || ownerGeneration !== generation) return; if (ctx.hasUI) { ctx.ui.notify( `Could not load pi-codex-compact.json; using defaults. ${terminalText(error instanceof Error ? error.message : String(error))}`, "warning" ); } return; } if (sessionController.signal.aborted || ownerGeneration !== generation || ctx.sessionManager.getSessionId() !== sessionId) { return; } if (ctx.hasUI && state.kind === "invalid") { ctx.ui.notify( `Invalid pi-codex-compact.json; using defaults without overwriting it. ${terminalText(state.issue ?? "unknown validation error")}`, "warning" ); } }); pi.on( "session_before_compact", (event, ctx) => compactRemotely(pi, event, ctx, settingsRuntime.get().settings, sessionController.signal, options.fetch) ); pi.on("context", (event, ctx) => { if (!settingsRuntime.get().settings.enabled) return void 0; const settings = settingsRuntime.get().settings; const checkpoint = activeCheckpoint(ctx); if (!checkpoint || !isCheckpointCompatible(checkpoint.details, ctx.model, settings)) return void 0; const messages = projectCheckpointContext(event.messages, checkpoint.details, checkpoint.entry.summary); return messages ? { messages } : void 0; }); pi.on("before_provider_request", (event, ctx) => { if (!settingsRuntime.get().settings.enabled) return void 0; const settings = settingsRuntime.get().settings; const checkpoint = activeCheckpoint(ctx); if (!checkpoint || !isCheckpointCompatible(checkpoint.details, ctx.model, settings)) return void 0; const marker = checkpointMarker(checkpoint.details.checkpointId); if (!hasCheckpointMarker(event.payload, marker)) return void 0; return rewriteCheckpointMarker(event.payload, marker, checkpoint.details.replacementHistory); }); pi.on("model_select", (event, ctx) => { if (!settingsRuntime.get().settings.enabled) return; const settings = settingsRuntime.get().settings; const checkpoint = activeCheckpoint(ctx); if (!checkpoint || isCheckpointCompatible(checkpoint.details, event.model, settings)) return; const key = `${ctx.sessionManager.getSessionId()}:${event.model.provider}:${event.model.id}`; if (providerWarnings.has(key)) return; providerWarnings.add(key); if (ctx.hasUI) { ctx.ui.notify( "The active Responses checkpoint cannot replay on this model; Pi will expose only its fallback marker and retained recent messages.", "warning" ); } }); pi.on("session_shutdown", async (_event, ctx) => { generation += 1; sessionController.abort(); providerWarnings.clear(); ctx.ui.setStatus(STATUS_KEY, void 0); await settingsRuntime.flush(); }); }; } var codex_compact_default = createCodexCompactExtension(); export { codex_compact_default as default }; //# sourceMappingURL=index.ts.map