修复消息渲染并更新发布包
This commit is contained in:
@@ -761,9 +761,17 @@ function createCodexAppRuntime(deps = {}) {
|
||||
if (!nextText) return '';
|
||||
if (!entry.agentMessageItems) entry.agentMessageItems = new Map();
|
||||
const currentItemText = entry.agentMessageItems.get(itemId) || '';
|
||||
if (!entry.agentMessagePendingPrefixes) entry.agentMessagePendingPrefixes = new Map();
|
||||
const pendingPrefix = entry.agentMessagePendingPrefixes.get(itemId) || '';
|
||||
if (!currentItemText && !nextText.trim()) {
|
||||
entry.agentMessagePendingPrefixes.set(itemId, `${pendingPrefix}${nextText}`);
|
||||
return '';
|
||||
}
|
||||
entry.agentMessagePendingPrefixes.delete(itemId);
|
||||
const separator = agentMessageSeparator(entry, itemId, nextText);
|
||||
const appended = separator + nextText;
|
||||
entry.agentMessageItems.set(itemId, appendCappedText(currentItemText, nextText, RUNTIME_AGENT_ITEM_MAX_CHARS));
|
||||
const firstText = currentItemText ? nextText : `${pendingPrefix}${nextText}`;
|
||||
const appended = separator + firstText;
|
||||
entry.agentMessageItems.set(itemId, appendCappedText(currentItemText, firstText, RUNTIME_AGENT_ITEM_MAX_CHARS));
|
||||
entry.fullText = appendCappedText(entry.fullText || '', appended, RUNTIME_FULL_TEXT_MAX_CHARS);
|
||||
return capStreamDelta(appended);
|
||||
}
|
||||
@@ -773,6 +781,10 @@ function createCodexAppRuntime(deps = {}) {
|
||||
if (!text) return '';
|
||||
if (!entry.agentMessageItems) entry.agentMessageItems = new Map();
|
||||
const currentItemText = entry.agentMessageItems.get(item.id) || '';
|
||||
if (!entry.agentMessagePendingPrefixes) entry.agentMessagePendingPrefixes = new Map();
|
||||
const pendingPrefix = entry.agentMessagePendingPrefixes.get(item.id) || '';
|
||||
entry.agentMessagePendingPrefixes.delete(item.id);
|
||||
if (!currentItemText && !text.trim()) return '';
|
||||
if (currentItemText && text.startsWith(currentItemText)) {
|
||||
const remainder = text.slice(currentItemText.length);
|
||||
entry.agentMessageItems.set(item.id, keepTail(text, RUNTIME_AGENT_ITEM_MAX_CHARS));
|
||||
@@ -780,9 +792,12 @@ function createCodexAppRuntime(deps = {}) {
|
||||
return capStreamDelta(remainder);
|
||||
}
|
||||
if (currentItemText === text) return '';
|
||||
const separator = agentMessageSeparator(entry, item.id, text);
|
||||
const appended = separator + text;
|
||||
entry.agentMessageItems.set(item.id, keepTail(text, RUNTIME_AGENT_ITEM_MAX_CHARS));
|
||||
const completedText = !currentItemText && pendingPrefix && !text.startsWith(pendingPrefix)
|
||||
? `${pendingPrefix}${text}`
|
||||
: text;
|
||||
const separator = agentMessageSeparator(entry, item.id, completedText);
|
||||
const appended = separator + completedText;
|
||||
entry.agentMessageItems.set(item.id, keepTail(completedText, RUNTIME_AGENT_ITEM_MAX_CHARS));
|
||||
entry.fullText = appendCappedText(entry.fullText || '', appended, RUNTIME_FULL_TEXT_MAX_CHARS);
|
||||
return capStreamDelta(appended);
|
||||
}
|
||||
|
||||
@@ -2,6 +2,39 @@
|
||||
|
||||
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';
|
||||
@@ -26,6 +59,11 @@ function createCodexAppServerClient(options = {}) {
|
||||
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) {
|
||||
@@ -137,8 +175,14 @@ function createCodexAppServerClient(options = {}) {
|
||||
return request('config/mcpServer/reload', {}, 30000);
|
||||
}
|
||||
|
||||
function start() {
|
||||
if (initPromise) return initPromise;
|
||||
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,
|
||||
@@ -147,6 +191,8 @@ function createCodexAppServerClient(options = {}) {
|
||||
windowsHide: true,
|
||||
});
|
||||
|
||||
processExitPromise = new Promise((resolve) => { resolveProcessExit = resolve; });
|
||||
|
||||
let stderr = '';
|
||||
proc.stderr.on('data', (chunk) => {
|
||||
stderr += chunk.toString();
|
||||
@@ -164,15 +210,21 @@ function createCodexAppServerClient(options = {}) {
|
||||
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 });
|
||||
});
|
||||
|
||||
initPromise = request('initialize', {
|
||||
return request('initialize', {
|
||||
clientInfo,
|
||||
capabilities: { experimentalApi: true },
|
||||
}, 30000)
|
||||
@@ -182,28 +234,108 @@ function createCodexAppServerClient(options = {}) {
|
||||
return result;
|
||||
})
|
||||
.catch((err) => {
|
||||
stop();
|
||||
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;
|
||||
}
|
||||
if (proc && !exited) {
|
||||
try { proc.kill('SIGTERM'); } catch {}
|
||||
setTimeout(() => {
|
||||
try {
|
||||
if (proc && !proc.killed) proc.kill('SIGKILL');
|
||||
} catch {}
|
||||
}, 3000);
|
||||
}
|
||||
proc = null;
|
||||
terminateProcessAndWait().catch(() => {});
|
||||
rejectAllPending(new Error('Codex app-server 已停止。'));
|
||||
}
|
||||
|
||||
|
||||
189
lib/codex-home-lock.js
Normal file
189
lib/codex-home-lock.js
Normal file
@@ -0,0 +1,189 @@
|
||||
'use strict';
|
||||
|
||||
const crypto = require('crypto');
|
||||
const fs = require('fs');
|
||||
const os = require('os');
|
||||
const path = require('path');
|
||||
|
||||
const DEFAULT_LOCK_WAIT_MS = 30_000;
|
||||
const DEFAULT_LOCK_POLL_MS = 200;
|
||||
const DEFAULT_LOCK_POLL_MAX_MS = 2_000;
|
||||
const INVALID_LOCK_GRACE_MS = 2_000;
|
||||
|
||||
|
||||
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 resolveCodexHome(env = process.env) {
|
||||
const explicit = String(env?.CODEX_HOME || '').trim();
|
||||
if (explicit) return path.resolve(explicit);
|
||||
const home = String(env?.HOME || env?.USERPROFILE || os.homedir() || '').trim();
|
||||
return home ? path.join(home, '.codex') : '';
|
||||
}
|
||||
|
||||
function readProcessStartTime(pid) {
|
||||
if (!Number.isInteger(pid) || pid <= 0 || process.platform === 'win32') return null;
|
||||
try {
|
||||
const stat = fs.readFileSync(`/proc/${pid}/stat`, 'utf8');
|
||||
const closingParen = stat.lastIndexOf(')');
|
||||
if (closingParen < 0) return null;
|
||||
const fields = stat.slice(closingParen + 1).trim().split(/\s+/);
|
||||
// The remainder starts at stat field 3; field 22 is index 19 here.
|
||||
return fields[19] || null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function currentProcessOwner() {
|
||||
return {
|
||||
pid: process.pid,
|
||||
startTime: readProcessStartTime(process.pid),
|
||||
token: crypto.randomBytes(16).toString('hex'),
|
||||
createdAt: new Date().toISOString(),
|
||||
};
|
||||
}
|
||||
|
||||
function readLock(lockPath) {
|
||||
try {
|
||||
const stat = fs.statSync(lockPath);
|
||||
let owner = null;
|
||||
try {
|
||||
const parsed = JSON.parse(fs.readFileSync(lockPath, 'utf8'));
|
||||
if (parsed && typeof parsed === 'object') owner = parsed;
|
||||
} catch {}
|
||||
return { owner, mtimeMs: stat.mtimeMs };
|
||||
} catch (error) {
|
||||
if (error?.code === 'ENOENT') return null;
|
||||
return { owner: null, mtimeMs: 0 };
|
||||
}
|
||||
}
|
||||
|
||||
function processAlive(pid) {
|
||||
if (!Number.isInteger(pid) || pid <= 0) return false;
|
||||
try {
|
||||
process.kill(pid, 0);
|
||||
return true;
|
||||
} catch (error) {
|
||||
return error?.code === 'EPERM';
|
||||
}
|
||||
}
|
||||
|
||||
function lockOwnerActive(owner) {
|
||||
if (!owner || !Number.isInteger(Number(owner.pid))) return false;
|
||||
const pid = Number(owner.pid);
|
||||
if (!processAlive(pid)) return false;
|
||||
const recordedStart = owner.startTime == null ? null : String(owner.startTime);
|
||||
const currentStart = readProcessStartTime(pid);
|
||||
if (recordedStart && currentStart && recordedStart !== currentStart) return false;
|
||||
// On platforms without /proc, a live PID is the strongest available signal.
|
||||
return true;
|
||||
}
|
||||
|
||||
function lockOwnerSummary(owner) {
|
||||
if (!owner || typeof owner !== 'object') return null;
|
||||
return {
|
||||
pid: Number.isInteger(Number(owner.pid)) ? Number(owner.pid) : null,
|
||||
startTime: owner.startTime || null,
|
||||
createdAt: owner.createdAt || null,
|
||||
};
|
||||
}
|
||||
|
||||
function wait(ms) {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
async function acquireCodexHomeLock(options = {}) {
|
||||
const lockPath = options.lockPath
|
||||
? path.resolve(String(options.lockPath))
|
||||
: (() => {
|
||||
const home = resolveCodexHome(options.env || process.env);
|
||||
return home ? path.join(home, '.cc-web-codex.lock') : '';
|
||||
})();
|
||||
if (!lockPath) {
|
||||
return { path: '', owner: null, release: async () => {} };
|
||||
}
|
||||
|
||||
const waitMs = positiveInt(
|
||||
options.waitMs ?? options.env?.CC_WEB_CODEX_LOCK_WAIT_MS,
|
||||
DEFAULT_LOCK_WAIT_MS,
|
||||
{ min: 0, max: 10 * 60 * 1000 },
|
||||
);
|
||||
const pollMs = positiveInt(options.pollMs, DEFAULT_LOCK_POLL_MS, { min: 10, max: 5000 });
|
||||
const pollMaxMs = positiveInt(options.pollMaxMs, DEFAULT_LOCK_POLL_MAX_MS, { min: pollMs, max: 10_000 });
|
||||
const owner = currentProcessOwner();
|
||||
const startedAt = Date.now();
|
||||
let lastOwner = null;
|
||||
|
||||
fs.mkdirSync(path.dirname(lockPath), { recursive: true });
|
||||
while (true) {
|
||||
try {
|
||||
const fd = fs.openSync(lockPath, 'wx', 0o600);
|
||||
try {
|
||||
fs.writeFileSync(fd, `${JSON.stringify(owner)}\n`, 'utf8');
|
||||
} finally {
|
||||
fs.closeSync(fd);
|
||||
}
|
||||
return createLockHandle(lockPath, owner);
|
||||
} catch (error) {
|
||||
if (error?.code !== 'EEXIST') throw error;
|
||||
const existing = readLock(lockPath);
|
||||
lastOwner = existing?.owner || null;
|
||||
const ageMs = existing ? Math.max(0, Date.now() - Number(existing.mtimeMs || 0)) : 0;
|
||||
const stale = existing && (
|
||||
(existing.owner && !lockOwnerActive(existing.owner))
|
||||
|| (!existing.owner && ageMs >= INVALID_LOCK_GRACE_MS)
|
||||
);
|
||||
if (stale) {
|
||||
try {
|
||||
const latest = readLock(lockPath);
|
||||
if (latest?.owner?.token === existing.owner?.token || (!latest?.owner && ageMs >= INVALID_LOCK_GRACE_MS)) {
|
||||
fs.unlinkSync(lockPath);
|
||||
continue;
|
||||
}
|
||||
} catch (unlinkError) {
|
||||
if (unlinkError?.code === 'ENOENT') continue;
|
||||
}
|
||||
}
|
||||
|
||||
if (Date.now() - startedAt >= waitMs) {
|
||||
const err = new Error(`Codex CODEX_HOME 已被占用: ${lockPath}`);
|
||||
err.code = 'CODEX_HOME_LOCK_BUSY';
|
||||
err.lockPath = lockPath;
|
||||
err.owner = lockOwnerSummary(lastOwner);
|
||||
throw err;
|
||||
}
|
||||
const elapsed = Date.now() - startedAt;
|
||||
const base = Math.min(pollMaxMs, pollMs * (2 ** Math.min(4, Math.floor(elapsed / 1000))));
|
||||
const jitter = Math.floor(Math.random() * Math.max(1, Math.min(100, base * 0.2)));
|
||||
await wait(Math.min(pollMaxMs, base + jitter));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function createLockHandle(lockPath, owner) {
|
||||
let released = false;
|
||||
return {
|
||||
path: lockPath,
|
||||
owner,
|
||||
release: async () => {
|
||||
if (released) return;
|
||||
released = true;
|
||||
try {
|
||||
const onDisk = readLock(lockPath)?.owner;
|
||||
if (onDisk?.token === owner.token) fs.unlinkSync(lockPath);
|
||||
} catch {}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
DEFAULT_LOCK_WAIT_MS,
|
||||
acquireCodexHomeLock,
|
||||
lockOwnerActive,
|
||||
readProcessStartTime,
|
||||
resolveCodexHome,
|
||||
};
|
||||
Reference in New Issue
Block a user