修复历史边界与流式消息去重并更新发布包

This commit is contained in:
shiyue
2026-09-21 00:38:51 +08:00
parent db04e088bc
commit f8f5ef3a0b
14 changed files with 1733 additions and 423 deletions

290
server.js
View File

@@ -74,6 +74,11 @@ if (process.argv.includes('--ccweb-mcp-server')) {
return;
}
if (process.argv.includes('--migrate-session-history')) {
require('./scripts/migrate-session-history');
return;
}
if (process.argv.includes('--codex-app-worker')) {
require('./lib/codex-app-worker');
return;
@@ -4438,91 +4443,121 @@ function messageReplyToRequestId(message) {
function ensureStableMessageId(message) {
if (!message || typeof message !== 'object' || Array.isArray(message)) return message;
if (String(message.id || '').trim()) return message;
const explicitId = String(message.clientMessageId || message.messageId || '').trim();
const replyToRequestId = messageReplyToRequestId(message);
const turnId = String(message.turnId || message.codexAppTurnId || '').trim();
const output = turnId && !message.turnId ? { ...message, turnId } : message;
if (String(output.id || '').trim()) return output;
const explicitId = String(output.clientMessageId || output.messageId || '').trim();
const replyToRequestId = messageReplyToRequestId(output);
const threadId = output.nativeThreadId || output.codexAppThreadId || '';
const stableId = explicitId
|| (replyToRequestId ? `reply:${replyToRequestId}` : '')
|| (message.codexAppTurnKey ? `codexapp:${message.codexAppTurnKey}` : '')
|| (message.nativeTurnKey ? `native:${message.nativeTurnKey}` : '');
if (stableId) return { ...message, id: stableId };
const content = typeof message.content === 'string' ? message.content : JSON.stringify(message.content || '');
const identity = [message.role || '', message.timestamp || message.createdAt || '', content].join('\u001f');
return { ...message, id: `message:${stableMessageHash(identity)}` };
|| (turnId ? `history:${output.role || ''}:${threadId}:${turnId}` : '')
|| (output.codexAppTurnKey ? `codexapp:${output.codexAppTurnKey}` : '')
|| (output.nativeTurnKey ? `native:${output.nativeTurnKey}` : '');
if (stableId) return { ...output, id: stableId };
const content = typeof output.content === 'string' ? output.content : JSON.stringify(output.content || '');
const identity = [output.role || '', output.timestamp || output.createdAt || '', content].join('\u001f');
return { ...output, id: `message:${stableMessageHash(identity)}` };
}
function normalizeHistoryMessage(message) {
return ensureStableMessageId(message);
}
function messagesEquivalent(left, right) {
if (!left || !right || left.role !== right.role) return false;
const leftId = String(left.id || '').trim();
const rightId = String(right.id || '').trim();
if (leftId && rightId && leftId === rightId) return true;
if (!leftId.startsWith('message:') || !rightId.startsWith('message:')) return false;
const leftContent = typeof left.content === 'string' ? left.content.trim() : JSON.stringify(left.content || '');
const rightContent = typeof right.content === 'string' ? right.content.trim() : JSON.stringify(right.content || '');
return !!leftContent
&& leftContent === rightContent
&& String(left.timestamp || '') === String(right.timestamp || '');
function messageIdentityKeys(message, session = {}) {
if (!message || typeof message !== 'object') return [];
const role = message.role || '';
const keys = [message.id, message.clientMessageId, message.messageId]
.map((value) => String(value || '').trim()).filter(Boolean)
.map((value) => `${role}:id:${value}`);
const replyId = messageReplyToRequestId(message);
if (replyId) keys.push(`${role}:reply:${replyId}`);
if (message.codexAppTurnKey) keys.push(`${role}:codexapp:${message.codexAppTurnKey}`);
if (message.nativeTurnKey) keys.push(`${role}:native:${message.nativeTurnKey}`);
const turnId = String(message.turnId || message.codexAppTurnId || '').trim();
const threadId = message.nativeThreadId || message.codexAppThreadId || session.codexAppThreadId || session.codexThreadId || '';
// 同 turn 的引导输入和分段回复是不同气泡,不能仅靠 turnId 将它们吞并。
if (role === 'assistant' && turnId && !message.nativeTurnSegment) {
keys.push(`${role}:turn:${threadId}:${turnId}`);
}
return keys;
}
function mergeHistorySegments(nativePrefix, snapshot) {
function messageContentKey(message) {
if (!message || !message.role) return '';
const timestamp = String(message.timestamp || message.createdAt || '').trim();
const time = Date.parse(timestamp);
const content = typeof message.content === 'string' ? message.content.trim() : JSON.stringify(message.content || '');
if (!timestamp || !Number.isFinite(time) || !content) return '';
return `${message.role}:${time}:${stableMessageHash(content)}`;
}
function messagesEquivalent(left, right, session = {}) {
if (!left || !right || left.role !== right.role) return false;
const rightKeys = new Set(messageIdentityKeys(right, session));
if (messageIdentityKeys(left, session).some((key) => rightKeys.has(key))) return true;
const contentKey = messageContentKey(left);
return !!contentKey && contentKey === messageContentKey(right);
}
function mergeHistorySegments(nativePrefix, snapshot, session = {}) {
const prefix = Array.isArray(nativePrefix) ? nativePrefix : [];
const tail = Array.isArray(snapshot) ? snapshot : [];
const snapshotIds = new Set(tail.map((message) => String(message?.id || '').trim()).filter(Boolean));
const snapshotKeys = new Set(tail.flatMap((message) => messageIdentityKeys(message, session)));
const merged = [];
const prefixIds = new Set();
for (const message of prefix) {
const id = String(message?.id || '').trim();
if (id && snapshotIds.has(id)) continue;
if (id && prefixIds.has(id)) continue;
if (id) prefixIds.add(id);
merged.push(message);
}
const snapshotPositions = new Map();
for (const message of tail) {
const id = String(message?.id || '').trim();
if (id && snapshotPositions.has(id)) {
merged[snapshotPositions.get(id)] = message;
continue;
const positions = new Map();
const earlier = prefix.filter((message) => !messageIdentityKeys(message, session).some((key) => snapshotKeys.has(key)));
for (const message of [...earlier, ...tail]) {
const keys = messageIdentityKeys(message, session);
const position = keys.map((key) => positions.get(key)).find((value) => value !== undefined);
if (position !== undefined) {
merged[position] = message;
keys.forEach((key) => positions.set(key, position));
} else {
keys.forEach((key) => positions.set(key, merged.length));
merged.push(message);
}
if (id) snapshotPositions.set(id, merged.length);
merged.push(message);
}
return merged;
}
function mergeNativeHistoryWithSnapshot(session, persistedMessages, nativeMessages) {
const snapshot = (Array.isArray(persistedMessages) ? persistedMessages : []).map(normalizeHistoryMessage);
const persisted = (Array.isArray(persistedMessages) ? persistedMessages : []).map(normalizeHistoryMessage);
const snapshot = mergeHistorySegments([], persisted, session);
const native = (Array.isArray(nativeMessages) ? nativeMessages : []).map(normalizeHistoryMessage);
if (snapshot.length === 0) return native.length > 0 ? { messages: mergeHistorySegments(native, []), confirmed: true } : { messages: snapshot, confirmed: true };
if (native.length === 0) return { messages: snapshot, confirmed: false };
const fallback = { messages: snapshot, confirmed: false };
if (snapshot.length === 0 || native.length === 0) return fallback;
// 从尾部寻找快照与 native 的最长有序重叠,允许 native 中存在旧版本额外事件。
let nativeCursor = native.length - 1;
let snapshotCursor = snapshot.length - 1;
while (nativeCursor >= 0 && snapshotCursor >= 0) {
if (messagesEquivalent(native[nativeCursor], snapshot[snapshotCursor])) {
nativeCursor -= 1;
snapshotCursor -= 1;
continue;
}
nativeCursor -= 1;
}
if (snapshotCursor < 0) {
const overlapStart = nativeCursor + 1;
return { messages: mergeHistorySegments(native.slice(0, overlapStart), snapshot), confirmed: true };
const notice = (session?.messages || []).find((message) => message?.ccwebPersistenceNotice === true);
const threadId = session?.codexAppThreadId || session?.codexThreadId || session?.claudeSessionId || '';
if (notice?.nativeThreadId && notice.nativeThreadId !== threadId) return fallback;
if (notice?.snapshotFirstMessageId && notice.snapshotFirstMessageId !== snapshot[0]?.id) return fallback;
if (notice?.snapshotMessageCount !== undefined) {
const count = Number(notice.snapshotMessageCount);
if (!Number.isSafeInteger(count) || count <= 0 || count > persisted.length) return fallback;
if (notice.snapshotLastMessageId && notice.snapshotLastMessageId !== persisted[count - 1]?.id) return fallback;
}
const baseIndex = Number(session?.historySnapshotBaseIndex);
const hasSnapshotBoundary = Number.isSafeInteger(baseIndex) && baseIndex > 0;
if (hasSnapshotBoundary && native.length >= baseIndex) {
const prefix = native.slice(0, Math.min(baseIndex, native.length));
return { messages: mergeHistorySegments(prefix, snapshot), confirmed: true };
// 只定位快照首条的可靠边界。快照尾部可能包含尚未写入 rollout 的当前输入和回传。
const firstKeys = new Set(messageIdentityKeys(snapshot[0], session));
let matches = native.map((message, index) => ({ message, index }))
.filter(({ message }) => messageIdentityKeys(message, session).some((key) => firstKeys.has(key)));
if (matches.length === 0) {
const fingerprint = messageContentKey(snapshot[0]);
if (!fingerprint) return fallback;
matches = native.map((message, index) => ({ message, index }))
.filter(({ message }) => messageContentKey(message) === fingerprint);
if (matches.length !== 1) return fallback;
} else if (matches.length > 1) {
// 原生记录重复同一个 ID 可以取最早出现处;多个不同 ID 的别名命中不能猜测。
if (new Set(matches.map(({ message }) => message.id)).size !== 1) return fallback;
}
return { messages: snapshot, confirmed: false };
const boundaryIndex = matches[0].index;
return {
messages: mergeHistorySegments(native.slice(0, boundaryIndex), snapshot, session),
confirmed: true,
nativePrefixCount: boundaryIndex,
};
}
function sessionHistoryCacheKey(session) {
@@ -4594,7 +4629,7 @@ function resolveSessionHistory(session) {
const notice = persisted.some(isPersistenceNoticeMessage);
if (!notice) {
return {
messages: persistedMessages.map(normalizeHistoryMessage),
messages: mergeHistorySegments([], persistedMessages.map(normalizeHistoryMessage), session),
source: 'snapshot',
available: true,
recoverable: false,
@@ -4603,7 +4638,7 @@ function resolveSessionHistory(session) {
const native = loadNativeSessionHistory(session);
const merged = mergeNativeHistoryWithSnapshot(session, persistedMessages, native?.messages);
if (merged.confirmed && merged.messages.length >= persistedMessages.length) {
if (merged.confirmed) {
return {
messages: merged.messages,
source: merged.messages.length > persistedMessages.length ? 'merged' : 'snapshot',
@@ -4618,7 +4653,7 @@ function resolveSessionHistory(session) {
nativeCount: Array.isArray(native.messages) ? native.messages.length : 0,
});
}
return { messages: persistedMessages.map(normalizeHistoryMessage), source: 'snapshot', available: false, recoverable: false };
return { messages: mergeHistorySegments([], persistedMessages.map(normalizeHistoryMessage), session), source: 'snapshot', available: false, recoverable: false };
}
function normalizeAgent(agent) {
@@ -5526,17 +5561,27 @@ function sanitizeMessageForPersist(message, limits = {}) {
return output;
}
function sanitizeMessagesForPersist(messages, limits = {}) {
const list = Array.isArray(messages) ? messages : [];
function sanitizeMessagesForPersist(messages, limits = {}, session = {}) {
const source = Array.isArray(messages) ? messages : [];
const previousNotice = source.find(isPersistenceNoticeMessage);
const list = source.filter((message) => !isPersistenceNoticeMessage(message));
const maxMessages = limits.maxMessages || SESSION_PERSIST_MAX_MESSAGES;
const selected = list.length > maxMessages ? list.slice(-maxMessages) : list;
const output = selected.map((message) => sanitizeMessageForPersist(ensureStableMessageId(message), limits));
if (list.length > selected.length) {
if (previousNotice || list.length > selected.length) {
const truncatedAt = list.length > selected.length
? new Date().toISOString()
: previousNotice.truncatedAt || previousNotice.timestamp || new Date().toISOString();
output.unshift({
role: 'system',
content: `历史消息过多,cc-web 本地快照只保留最近 ${selected.length} 条;点击顶部“查看更早消息”可从原始会话记录加载省略的 ${list.length - selected.length} 条旧消息。`,
timestamp: new Date().toISOString(),
content: `cc-web 本地快照保留最近 ${selected.length} 条消息;更早历史仅在确认与原始会话的边界后加载。`,
timestamp: truncatedAt,
ccwebPersistenceNotice: true,
snapshotMessageCount: output.length,
snapshotFirstMessageId: output[0]?.id || null,
snapshotLastMessageId: output[output.length - 1]?.id || null,
nativeThreadId: session.codexAppThreadId || session.codexThreadId || session.claudeSessionId || previousNotice?.nativeThreadId || null,
truncatedAt,
});
}
return output;
@@ -5571,29 +5616,10 @@ function sanitizeSessionForPersist(session, limits = {}) {
for (const [key, value] of Object.entries(session || {})) {
if (skipKeys.has(key)) continue;
if (key === 'messages') {
output.messages = sanitizeMessagesForPersist(value, limits);
const sourceMessages = Array.isArray(value) ? value : [];
const persistedMessages = sourceMessages.filter((message) => !isPersistenceNoticeMessage(message));
const existingBaseIndex = Number(session?.historySnapshotBaseIndex);
const existingCount = Number(session?.historySnapshotCount);
if (sourceMessages.length > (limits.maxMessages || SESSION_PERSIST_MAX_MESSAGES)) {
const maxPersistedMessages = limits.maxMessages || SESSION_PERSIST_MAX_MESSAGES;
const logicalDroppedCount = Math.max(0, persistedMessages.length - maxPersistedMessages);
output.historySnapshotBaseIndex = Number.isSafeInteger(existingBaseIndex) && existingBaseIndex > 0
? existingBaseIndex + logicalDroppedCount
: logicalDroppedCount;
output.historySnapshotCount = Math.max(0, output.messages.length - (output.messages.some(isPersistenceNoticeMessage) ? 1 : 0));
} else if (Number.isSafeInteger(existingBaseIndex) && existingBaseIndex > 0) {
output.historySnapshotBaseIndex = existingBaseIndex;
output.historySnapshotCount = Number.isSafeInteger(existingCount) && existingCount > 0
? Math.max(existingCount, persistedMessages.length)
: persistedMessages.length;
} else {
output.historySnapshotBaseIndex = 0;
output.historySnapshotCount = persistedMessages.length;
}
output.messages = sanitizeMessagesForPersist(value, limits, session);
continue;
}
if (key === 'historySnapshotBaseIndex' || key === 'historySnapshotCount') continue;
output[key] = sanitizePersistValue(value, {
maxString: limits.topLevelMaxChars || 16 * 1024,
maxDepth: 4,
@@ -5602,6 +5628,13 @@ function sanitizeSessionForPersist(session, limits = {}) {
});
}
if (!Object.prototype.hasOwnProperty.call(output, 'messages')) output.messages = [];
const sourceMessages = (Array.isArray(session?.messages) ? session.messages : []).filter((message) => !isPersistenceNoticeMessage(message));
const snapshotMessages = output.messages.filter((message) => !isPersistenceNoticeMessage(message));
const existingBaseIndex = Number(session?.historySnapshotBaseIndex);
// 基线只表示 cc-web 的逻辑序列位置,不能拿它直接切割不同粒度的 native 数组。
output.historySnapshotBaseIndex = (Number.isSafeInteger(existingBaseIndex) && existingBaseIndex > 0 ? existingBaseIndex : 0)
+ Math.max(0, sourceMessages.length - snapshotMessages.length);
output.historySnapshotCount = snapshotMessages.length;
return normalizeSession(output);
}
@@ -9762,6 +9795,9 @@ wss.on('connection', (ws, req) => {
case 'load_history_page':
handleLoadHistoryPage(ws, msg);
break;
case 'migrate_session_history':
handleMigrateSessionHistory(ws, msg);
break;
case 'search_sessions':
handleSearchSessions(ws, msg);
break;
@@ -11313,6 +11349,61 @@ function handleLoadHistoryPage(ws, msg = {}) {
}, msg));
}
function handleMigrateSessionHistory(ws, msg = {}) {
const sessionId = sanitizeId(msg.sessionId || '');
const respond = (result) => wsSend(ws, attachClientRequestId({
type: 'migrate_session_history_result', sessionId, ...result,
}, msg));
// 检查、备份和写入在服务端同一同步调用内完成,避免运行任务保存覆盖迁移。
if (isSessionRunning(sessionId)) return respond({ ok: false, code: 'session_running' });
const filePath = sessionPath(sessionId);
if (!sessionId || !fs.existsSync(filePath)) return respond({ ok: false, code: 'session_not_found' });
try {
const original = fs.readFileSync(filePath, 'utf8');
const session = normalizeSession(JSON.parse(original));
if (!session.messages.some(isPersistenceNoticeMessage)) {
return respond({ ok: true, changed: false, messageCount: session.messages.length });
}
const snapshot = session.messages.filter((message) => !isPersistenceNoticeMessage(message));
const native = loadNativeSessionHistory(session);
const merged = mergeNativeHistoryWithSnapshot(session, snapshot, native?.messages);
if (!merged.confirmed) return respond({ ok: false, code: 'history_boundary_unconfirmed' });
const migrated = {
...session,
messages: merged.messages,
historySnapshotBaseIndex: 0,
historySnapshotCount: merged.messages.length,
};
// 迁移保留快照原文;不再走会立即截断的普通持久化路径。
const json = JSON.stringify(migrated, null, 2);
if (textByteLength(json) > SESSION_LOAD_MAX_BYTES) return respond({ ok: false, code: 'history_migration_too_large' });
const details = { snapshotCount: snapshot.length, messageCount: merged.messages.length, nativePrefixCount: merged.nativePrefixCount };
if (msg.apply !== true) return respond({ ok: true, preview: true, ...details });
if (isSessionRunning(sessionId) || fs.readFileSync(filePath, 'utf8') !== original) {
return respond({ ok: false, code: 'session_changed' });
}
const backupDir = path.join(SESSIONS_DIR, '_history-backups');
fs.mkdirSync(backupDir, { recursive: true });
const backupPath = path.join(backupDir, `${sessionId}.${Date.now()}.${crypto.randomUUID()}.json`);
fs.writeFileSync(backupPath, original, { flag: 'wx', mode: 0o600 });
writeFileAtomicSync(filePath, json);
const reloaded = loadSession(sessionId);
if (!reloaded || reloaded.messages.length !== migrated.messages.length) {
writeFileAtomicSync(filePath, original);
return respond({ ok: false, code: 'history_reload_failed', backupPath });
}
sessionHistoryCache.delete(sessionHistoryCacheKey(session));
sessionSearchIndex.scheduleUpsert(sessionId);
scheduleUsageStatisticsUpsert(sessionId);
respond({ ok: true, changed: true, backupPath, ...details });
const viewingWs = findViewingSessionWs(sessionId);
if (viewingWs) wsSend(viewingWs, buildSessionInfoPayload(reloaded));
} catch (error) {
plog('WARN', 'session_history_migration_failed', { sessionId, error: error?.message || String(error) });
respond({ ok: false, code: 'history_migration_failed' });
}
}
function attachActiveRuntimeToWs(ws, sessionId, source = {}) {
if (activeProcesses.has(sessionId)) {
const entry = activeProcesses.get(sessionId);
@@ -11345,6 +11436,11 @@ function attachActiveRuntimeToWs(ws, sessionId, source = {}) {
wsSend(ws, attachClientRequestId({
type: 'resume_generating',
sessionId,
assistantMessageId: ensureStableMessageId({
role: 'assistant', codexAppThreadId: entry.threadId, codexAppTurnId: entry.turnId,
codexAppTurnKey: codexAppTurnKey(sessionId, entry),
}).id,
turnId: entry.turnId || null,
text: truncateTextValue(entry.fullText || '', SESSION_MESSAGE_CONTENT_MAX_CHARS),
toolCalls: sanitizeToolCallsForPersist(entry.toolCalls || []),
}, source));
@@ -14268,7 +14364,7 @@ function handleCodexAppTurnComplete(sessionId, options = {}) {
}
if (session && (assistantContent.trim() || assistantToolCalls.length > 0) && !hasCodexAppTurnMessage(session, turnKey)) {
const assistantMessage = {
const assistantMessage = ensureStableMessageId({
role: 'assistant',
content: assistantContent,
toolCalls: assistantToolCalls,
@@ -14277,7 +14373,7 @@ function handleCodexAppTurnComplete(sessionId, options = {}) {
codexAppThreadId: entry.threadId || null,
codexAppTurnId: entry.turnId || null,
interrupted: !!options.interrupted,
};
});
const beforeUserMessage = options.beforeUserMessage;
const beforeUserMessageIndex = beforeUserMessage
? session.messages.findIndex((message) => (
@@ -14355,7 +14451,14 @@ function handleCodexAppTurnComplete(sessionId, options = {}) {
entry.errorSent = true;
wsSend(entry.ws, { type: 'error', sessionId, message: completionError });
}
wsSend(entry.ws, { type: 'done', sessionId, costUsd: null, goalActive });
const assistantMessage = session?.messages?.find((message) => message.codexAppTurnKey === turnKey);
const assistantMessageIndex = assistantMessage
? resolveSessionHistory(session).messages.findIndex((message) => message.id === assistantMessage.id)
: -1;
wsSend(entry.ws, {
type: 'done', sessionId, costUsd: null, goalActive,
...(assistantMessage ? { assistantMessage: sanitizeMessageForTransport(assistantMessage), assistantMessageIndex } : {}),
});
sendSessionList(entry.ws);
return;
}
@@ -14546,6 +14649,7 @@ function handleCodexAppSteerMessage(ws, msg, options = {}) {
}
persistedUserMessage = {
id: userMessageId,
...(clientMessageId ? { clientMessageId } : {}),
role: 'user',
content: textValue,
attachments: savedAttachments,