Files
cc-web/lib/codex-app-server-client.js
2026-09-28 08:13:09 +08:00

358 lines
11 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

'use strict';
const readline = require('readline');
const { spawn } = require('child_process');
const {
acquireCodexHomeLock,
} = require('./codex-home-lock');
const DEFAULT_START_RETRY_ATTEMPTS = 4;
const DEFAULT_START_RETRY_DELAY_MS = 500;
const DEFAULT_START_RETRY_MAX_DELAY_MS = 5000;
function positiveInt(value, fallback, { min = 1, max = Number.MAX_SAFE_INTEGER } = {}) {
const parsed = Number.parseInt(String(value ?? ''), 10);
if (!Number.isFinite(parsed) || parsed < min) return fallback;
return Math.min(parsed, max);
}
function isCodexHomeBusyError(value) {
const parts = [];
const visit = (item, depth = 0) => {
if (!item || depth > 3) return;
if (typeof item === 'string') parts.push(item);
else if (typeof item === 'object') {
for (const key of ['message', 'stderr', 'cause', 'code', 'details']) visit(item[key], depth + 1);
}
};
visit(value);
const text = parts.join(' ');
return /database\s+(?:is\s+)?locked|database\s+table\s+is\s+locked|SQLITE_BUSY(?:_TIMEOUT)?|failed\s+to\s+(?:open|initialize)\s+(?:state|log|database)|failed\s+to\s+initialize\s+state\s+runtime/i.test(text);
}
function retryDelayMs(attempt, base, max) {
const exponential = Math.min(max, base * (2 ** Math.max(0, attempt - 1)));
const jitter = Math.floor(Math.random() * Math.max(1, Math.min(100, exponential * 0.2)));
return Math.min(max, exponential + jitter);
}
function createCodexAppServerClient(options = {}) {
const command = options.command || 'codex';
const args = Array.isArray(options.args) && options.args.length > 0
? options.args.slice()
: ['app-server', '--stdio'];
const env = options.env || process.env;
const cwd = options.cwd || process.cwd();
const clientInfo = options.clientInfo || {
name: 'ccweb_codexapp',
title: 'CC-Web Codex App',
version: '0.1.0',
};
const onNotification = typeof options.onNotification === 'function' ? options.onNotification : () => {};
const onServerRequest = typeof options.onServerRequest === 'function' ? options.onServerRequest : null;
const onExit = typeof options.onExit === 'function' ? options.onExit : () => {};
const onLog = typeof options.onLog === 'function' ? options.onLog : () => {};
const postInitialize = typeof options.postInitialize === 'function' ? options.postInitialize : null;
let proc = null;
let rl = null;
let nextId = 1;
let initPromise = null;
let exited = false;
let processExitPromise = null;
let resolveProcessExit = null;
let lockHandle = null;
let lastStartError = null;
let stopRequested = false;
const pending = new Map();
function rejectAllPending(err) {
for (const [, pendingRequest] of pending) {
clearTimeout(pendingRequest.timer);
pendingRequest.reject(err);
}
pending.clear();
}
function sendRaw(message) {
if (!proc || !proc.stdin || proc.stdin.destroyed) {
throw new Error('Codex app-server 未启动。');
}
proc.stdin.write(`${JSON.stringify(message)}\n`);
}
function respondToServerRequest(id, result, error) {
try {
if (error) {
sendRaw({ id, error });
} else {
sendRaw({ id, result: result || {} });
}
} catch (err) {
onLog('WARN', 'codex_app_server_response_failed', { error: err.message });
}
}
function handleServerRequest(message) {
const id = message.id;
const method = message.method;
const params = message.params || {};
if (onServerRequest) {
Promise.resolve()
.then(() => onServerRequest({ method, params, id }))
.then((result) => respondToServerRequest(id, result || {}))
.catch((err) => respondToServerRequest(id, null, {
code: -32603,
message: err?.message || 'cc-web 无法处理 Codex app-server 请求。',
}));
return;
}
respondToServerRequest(id, null, {
code: -32601,
message: `cc-web 暂不支持 Codex app-server 请求: ${method}`,
});
}
function handleMessage(line) {
let message;
try {
message = JSON.parse(line);
} catch {
onLog('WARN', 'codex_app_server_invalid_json', { line: String(line || '').slice(0, 200) });
return;
}
if (Object.prototype.hasOwnProperty.call(message, 'id')) {
const pendingRequest = pending.get(message.id);
if (pendingRequest) {
pending.delete(message.id);
clearTimeout(pendingRequest.timer);
if (message.error) {
const err = new Error(message.error.message || 'Codex app-server 请求失败。');
err.code = message.error.code;
err.data = message.error.data;
pendingRequest.reject(err);
} else {
pendingRequest.resolve(message.result || {});
}
return;
}
if (message.method) {
handleServerRequest(message);
return;
}
}
if (message.method) {
onNotification(message);
}
}
function request(method, params = {}, timeoutMs = 300000) {
const id = nextId++;
const message = { id, method, params };
return new Promise((resolve, reject) => {
const timer = setTimeout(() => {
pending.delete(id);
reject(new Error(`Codex app-server 请求超时: ${method}`));
}, timeoutMs);
pending.set(id, { resolve, reject, timer, method });
try {
sendRaw(message);
} catch (err) {
clearTimeout(timer);
pending.delete(id);
reject(err);
}
});
}
function notification(method, params = {}) {
sendRaw({ method, params });
}
function reloadMcpServers() {
return request('config/mcpServer/reload', {}, 30000);
}
function releaseLock() {
const current = lockHandle;
lockHandle = null;
if (current) Promise.resolve(current.release()).catch(() => {});
}
function spawnAndInitialize(lock) {
lockHandle = lock;
exited = false;
proc = spawn(command, args, {
env,
cwd,
stdio: ['pipe', 'pipe', 'pipe'],
windowsHide: true,
});
processExitPromise = new Promise((resolve) => { resolveProcessExit = resolve; });
let stderr = '';
proc.stderr.on('data', (chunk) => {
stderr += chunk.toString();
if (stderr.length > 4000) stderr = stderr.slice(-4000);
});
rl = readline.createInterface({ input: proc.stdout });
rl.on('line', handleMessage);
proc.on('exit', (code, signal) => {
exited = true;
if (rl) rl.close();
const err = new Error(`Codex app-server 已退出: code=${code ?? 'null'} signal=${signal || 'null'}`);
err.exitCode = code;
err.signal = signal;
err.stderr = stderr;
rejectAllPending(err);
releaseLock();
if (resolveProcessExit) resolveProcessExit({ code, signal, stderr });
resolveProcessExit = null;
onExit({ code, signal, stderr });
});
proc.on('error', (err) => {
rejectAllPending(err);
releaseLock();
if (resolveProcessExit) resolveProcessExit({ code: null, signal: null, stderr: err.message });
resolveProcessExit = null;
onExit({ code: null, signal: null, stderr: err.message });
});
return request('initialize', {
clientInfo,
capabilities: { experimentalApi: true },
}, 30000)
.then(async (result) => {
notification('initialized', {});
if (postInitialize) await postInitialize({ request, notification, onLog });
return result;
})
.catch((err) => {
throw err;
});
}
async function terminateProcessAndWait() {
const child = proc;
if (!child || exited) {
releaseLock();
return;
}
try { child.kill('SIGTERM'); } catch {}
const timer = setTimeout(() => {
try { if (!child.killed) child.kill('SIGKILL'); } catch {}
}, 3000);
try {
await processExitPromise;
} finally {
clearTimeout(timer);
}
}
async function startWithRetry() {
const env = options.env || process.env;
const attempts = positiveInt(
options.startRetryAttempts ?? env.CC_WEB_CODEX_START_RETRY_ATTEMPTS,
DEFAULT_START_RETRY_ATTEMPTS,
{ min: 1, max: 8 },
);
const baseDelay = positiveInt(
options.startRetryDelayMs ?? env.CC_WEB_CODEX_START_RETRY_DELAY_MS,
DEFAULT_START_RETRY_DELAY_MS,
{ min: 0, max: 30_000 },
);
const maxDelay = positiveInt(
options.startRetryMaxDelayMs ?? env.CC_WEB_CODEX_START_RETRY_MAX_DELAY_MS,
DEFAULT_START_RETRY_MAX_DELAY_MS,
{ min: baseDelay || 1, max: 60_000 },
);
for (let attempt = 1; attempt <= attempts; attempt += 1) {
if (stopRequested) {
const aborted = new Error('Codex app-server 启动已取消。');
aborted.code = 'CODEX_APP_START_CANCELLED';
throw aborted;
}
const lock = await acquireCodexHomeLock({
env,
waitMs: options.lockWaitMs,
lockPath: options.lockPath,
});
if (stopRequested) {
await lock.release();
const aborted = new Error('Codex app-server 启动已取消。');
aborted.code = 'CODEX_APP_START_CANCELLED';
throw aborted;
}
try {
const result = await spawnAndInitialize(lock);
lastStartError = null;
return result;
} catch (error) {
lastStartError = error;
await terminateProcessAndWait();
if (!isCodexHomeBusyError(error) || attempt >= attempts) {
const detail = error?.stderr || error?.message || String(error || '');
const finalError = new Error(`Codex app-server 启动失败(尝试 ${attempt}/${attempts}):${detail}`);
finalError.code = error?.code || 'CODEX_APP_START_FAILED';
finalError.cause = error;
finalError.stderr = error?.stderr || '';
finalError.lockPath = lock.path;
throw finalError;
}
const delay = retryDelayMs(attempt, baseDelay, maxDelay);
onLog('WARN', 'codex_app_server_start_retry', {
attempt,
nextAttempt: attempt + 1,
delayMs: delay,
reason: 'sqlite_busy',
lockPath: lock.path,
});
await new Promise((resolve) => setTimeout(resolve, delay));
}
}
throw lastStartError || new Error('Codex app-server 启动失败。');
}
function start() {
if (initPromise) return initPromise;
stopRequested = false;
initPromise = startWithRetry();
return initPromise;
}
function stop() {
stopRequested = true;
initPromise = null;
if (rl) {
try { rl.close(); } catch {}
rl = null;
}
terminateProcessAndWait().catch(() => {});
rejectAllPending(new Error('Codex app-server 已停止。'));
}
function isRunning() {
return !!proc && !exited;
}
return {
start,
stop,
request,
notification,
reloadMcpServers,
isRunning,
pid: () => proc?.pid || null,
};
}
module.exports = { createCodexAppServerClient };