/** Domain: quota transport. Owns: Pi auth, Codex app-server fallback and errors. Excludes: shared leadership and UI. */ import { type ChildProcessWithoutNullStreams, spawn } from "node:child_process"; import { createInterface } from "node:readline"; import type { ExtensionContext } from "@earendil-works/pi-coding-agent"; import { assertObject, errorMessage, normalizeAppServerResponse, normalizeBackendPayload, parseJsonObject, selectUsageSnapshot, type AppServerRateLimitResponse, type CodexUsageModel, type CodexUsageReport, type RateLimitStatusPayload, type UsageQueryError } from "./usage.ts"; const CODEX_PROVIDER_ID = "openai-codex"; const CODEX_USAGE_LIMIT_ID = "codex"; const CODEX_USAGE_URL = "https://chatgpt.com/backend-api/wham/usage"; const MAX_ERROR_BODY_CHARS = 600; type PiModel = NonNullable; type QueryUsageOptions = { timeoutMs: number }; export type QueryUsageResult = { ok: true; report: CodexUsageReport } | { ok: false; errors: UsageQueryError[] }; type RpcResponse = { id?: unknown; result?: unknown; error?: { message?: unknown; code?: unknown }; }; type PendingRpc = { resolve: (value: unknown) => void; reject: (error: Error) => void; }; export async function queryUsage( ctx: ExtensionContext, options: Pick, model: CodexUsageModel | undefined, mayQuery: () => boolean, ): Promise { const errors: UsageQueryError[] = []; const sources = ["pi-auth", "codex-app-server"] as const; for (const source of sources) { if (!mayQuery()) return undefined; try { const report = source === "pi-auth" ? await queryViaPiAuth(ctx, options.timeoutMs, mayQuery) : await queryViaCodexAppServer(options.timeoutMs, mayQuery); if (!report) return undefined; if ( selectUsageSnapshot(report, CODEX_USAGE_LIMIT_ID) || report.credits ) { return { ok: true, report }; } errors.push({ source, message: `${source} returned no displayable codex rate-limit windows`, }); } catch (cause) { errors.push({ source, message: errorMessage(cause), cause }); } } return { ok: false, errors }; } async function queryViaPiAuth( ctx: ExtensionContext, timeoutMs: number, mayQuery: () => boolean, ): Promise { const auth = await resolvePiCodexAuth(ctx); if (!auth) { throw new Error( "No Pi OpenAI Codex subscription auth was available. Use a Pi OpenAI Codex model or run /login for OpenAI ChatGPT Plus/Pro (Codex).", ); } if (!mayQuery()) return undefined; const response = await fetchWithTimeout( CODEX_USAGE_URL, { headers: auth.headers }, timeoutMs, ); const text = await response.text(); if (!response.ok) { throw new Error( `Codex usage endpoint returned ${response.status} ${response.statusText}: ${redactErrorBody(text)}`, ); } const payload = parseJsonObject(text, "Codex usage endpoint response"); return normalizeBackendPayload( payload as RateLimitStatusPayload, Date.now(), "pi-auth", ); } async function resolvePiCodexAuth( ctx: ExtensionContext, ): Promise<{ headers: Record } | undefined> { const models = codexAuthCandidateModels(ctx); const errors: string[] = []; for (const model of models) { const auth = await ctx.modelRegistry.getApiKeyAndHeaders(model); if (!auth.ok) { errors.push(auth.error); continue; } const headers: Record = {}; for (const [name, value] of Object.entries(auth.headers ?? {})) { if (value !== null) headers[name] = value; } if (!hasHeader(headers, "Authorization") && auth.apiKey) { headers.Authorization = `Bearer ${auth.apiKey}`; } if (!hasHeader(headers, "User-Agent")) { headers["User-Agent"] = "pi-codex-usage"; } if (hasHeader(headers, "Authorization")) { return { headers }; } } if (errors.length > 0) { throw new Error(errors.join("; ")); } return undefined; } function codexAuthCandidateModels(ctx: ExtensionContext): PiModel[] { const candidates: PiModel[] = []; const seen = new Set(); const add = (model: PiModel | undefined) => { if (!model || model.provider !== CODEX_PROVIDER_ID) return; const key = `${model.provider}/${model.id}`; if (seen.has(key)) return; seen.add(key); candidates.push(model); }; add(ctx.model); for (const model of ctx.modelRegistry.getAvailable()) add(model); for (const model of ctx.modelRegistry.getAll()) add(model); return candidates; } async function fetchWithTimeout( url: string, init: RequestInit, timeoutMs: number, ): Promise { const controller = new AbortController(); const timeout = setTimeout(() => controller.abort(), timeoutMs); try { return await fetch(url, { ...init, signal: controller.signal }); } catch (error) { if (controller.signal.aborted) { throw new Error( `Timed out after ${Math.round(timeoutMs / 1000)}s while fetching Codex usage.`, ); } throw error; } finally { clearTimeout(timeout); } } async function queryViaCodexAppServer( timeoutMs: number, mayQuery: () => boolean, ): Promise { const client = new CodexAppServerClient(timeoutMs); try { await client.start(); if (!mayQuery()) return undefined; await client.request("initialize", { clientInfo: { name: "pi_codex_usage", title: "Pi Codex Usage", version: "0.1.0", }, capabilities: { experimentalApi: false, requestAttestation: false, optOutNotificationMethods: [], }, }); client.notify("initialized"); if (!mayQuery()) return undefined; const result = await client.request("account/rateLimits/read", undefined); return normalizeAppServerResponse( assertObject( result, "account/rateLimits/read result", ) as AppServerRateLimitResponse, Date.now(), ); } finally { client.dispose(); } } class CodexAppServerClient { private child?: ChildProcessWithoutNullStreams; private nextId = 1; private stderr = ""; private readonly pending = new Map(); private startPromise?: Promise; private exitError?: Error; private readonly timeoutMs: number; constructor(timeoutMs: number) { this.timeoutMs = timeoutMs; } start(): Promise { if (this.startPromise) return this.startPromise; this.startPromise = new Promise((resolve, reject) => { const child = spawn("codex", ["app-server", "--listen", "stdio://"], { stdio: ["pipe", "pipe", "pipe"], }); this.child = child; const startupTimeout = setTimeout(() => { reject( new Error( `Timed out after ${Math.round(this.timeoutMs / 1000)}s starting codex app-server.`, ), ); }, this.timeoutMs); child.once("spawn", () => { clearTimeout(startupTimeout); resolve(); }); child.once("error", (error) => { clearTimeout(startupTimeout); reject(new Error(`Failed to start codex app-server: ${error.message}`)); this.rejectAll(error); }); child.once("exit", (code, signal) => { const suffix = this.stderr ? ` stderr: ${redactErrorBody(this.stderr)}` : ""; this.exitError = new Error( `codex app-server exited before completing the request (code ${code ?? "unknown"}, signal ${signal ?? "none"}).${suffix}`, ); this.rejectAll(this.exitError); }); child.stderr.setEncoding("utf8"); child.stderr.on("data", (chunk: string) => { this.stderr = truncateEnd(this.stderr + chunk, MAX_ERROR_BODY_CHARS); }); const lines = createInterface({ input: child.stdout }); lines.on("line", (line) => this.handleLine(line)); }); return this.startPromise; } request(method: string, params: unknown): Promise { const child = this.child; if (!child?.stdin.writable) { throw new Error("codex app-server is not running."); } if (this.exitError) throw this.exitError; const id = this.nextId++; const payload = params === undefined ? { method, id } : { method, id, params }; const response = new Promise((resolve, reject) => { const timeout = setTimeout(() => { this.pending.delete(id); reject( new Error( `Timed out after ${Math.round(this.timeoutMs / 1000)}s waiting for ${method}.`, ), ); }, this.timeoutMs); this.pending.set(id, { resolve: (value) => { clearTimeout(timeout); resolve(value); }, reject: (error) => { clearTimeout(timeout); reject(error); }, }); }); child.stdin.write(`${JSON.stringify(payload)}\n`); return response; } notify(method: string): void { const child = this.child; if (!child?.stdin.writable) return; child.stdin.write(`${JSON.stringify({ method })}\n`); } dispose(): void { for (const [id, pending] of this.pending) { pending.reject(new Error(`codex app-server request ${id} cancelled.`)); } this.pending.clear(); const child = this.child; if (!child) return; child.stdin.end(); if (!child.killed) child.kill(); this.child = undefined; } private handleLine(line: string): void { let parsed: RpcResponse; try { parsed = JSON.parse(line) as RpcResponse; } catch { return; } if (typeof parsed.id !== "number") return; const pending = this.pending.get(parsed.id); if (!pending) return; this.pending.delete(parsed.id); if (parsed.error) { const message = typeof parsed.error.message === "string" ? parsed.error.message : "unknown error"; pending.reject(new Error(`codex app-server request failed: ${message}`)); return; } pending.resolve(parsed.result); } private rejectAll(error: Error): void { for (const pending of this.pending.values()) pending.reject(error); this.pending.clear(); } } function hasHeader(headers: Record, name: string): boolean { return Object.keys(headers).some( (key) => key.toLowerCase() === name.toLowerCase(), ); } function redactErrorBody(body: string): string { return truncateEnd( body .replace(/Bearer\s+[A-Za-z0-9._~+/=-]+/gi, "Bearer ") .replace(/"access_token"\s*:\s*"[^"]+"/gi, '"access_token":""') .trim(), MAX_ERROR_BODY_CHARS, ); } function truncateEnd(value: string, maxChars: number): string { if (value.length <= maxChars) return value; return `${value.slice(0, maxChars - 1)}…`; }