339 lines
14 KiB
JavaScript
339 lines
14 KiB
JavaScript
'use strict';
|
||
|
||
const nativeFs = require('fs');
|
||
const path = require('path');
|
||
const crypto = require('crypto');
|
||
|
||
const SESSION_RUNTIME_KEYS = new Set([
|
||
'ws', 'tailer', 'codexAppStateTimer', 'codexAppStateDirty', 'codexAppStateCleaned',
|
||
'toolOutputDeltas', 'agentMessageItems',
|
||
]);
|
||
|
||
function createSessionHistoryStore(options = {}) {
|
||
const fs = options.fs || nativeFs;
|
||
const historyDir = path.join(options.sessionsDir, '_history');
|
||
const turnsDir = path.join(historyDir, '_turns');
|
||
// 身份只留在当前存储实例,不能通过 JSON 快照伪造可信完整加载或修订。
|
||
const loadedStates = new WeakMap();
|
||
const removedIds = new Set();
|
||
|
||
function validId(value) {
|
||
const id = String(value || '');
|
||
if (!/^[a-zA-Z0-9_-]+$/.test(id)) throw Object.assign(new Error('会话归档 ID 无效'), { code: 'history_invalid_id' });
|
||
return id;
|
||
}
|
||
|
||
function filePath(id) {
|
||
return path.join(historyDir, `${validId(id)}.json`);
|
||
}
|
||
|
||
function readRecord(id) {
|
||
let json;
|
||
try {
|
||
json = fs.readFileSync(filePath(id), 'utf8');
|
||
} catch (error) {
|
||
if (error?.code === 'ENOENT') return null;
|
||
throw error;
|
||
}
|
||
const record = JSON.parse(json);
|
||
if (record?.version !== 1 || typeof record.revision !== 'string' || !record.revision
|
||
|| record.session?.id !== id || !Array.isArray(record.session.messages)
|
||
|| typeof record.complete !== 'boolean') {
|
||
throw Object.assign(new Error('完整会话归档格式无效'), { code: 'history_invalid_archive' });
|
||
}
|
||
return record;
|
||
}
|
||
|
||
function state(session) {
|
||
const value = session && loadedStates.get(session);
|
||
if (!value) return null;
|
||
const { baseline, ...metadata } = value;
|
||
return metadata;
|
||
}
|
||
|
||
function load(id) {
|
||
const record = readRecord(validId(id));
|
||
if (!record) return null;
|
||
loadedStates.set(record.session, {
|
||
id, revision: record.revision, complete: record.complete, source: 'archive',
|
||
resetRevision: record.resetRevision || null, baseline: JSON.stringify(record.session),
|
||
});
|
||
return record.session;
|
||
}
|
||
|
||
function markLegacy(session) {
|
||
validId(session?.id);
|
||
loadedStates.set(session, {
|
||
id: session.id, revision: null, source: 'legacy',
|
||
complete: !session.messages?.some(message => message?.ccwebPersistenceNotice === true)
|
||
&& !(Number(session.historySnapshotBaseIndex) > 0),
|
||
});
|
||
return session;
|
||
}
|
||
|
||
function markComplete(session, options = {}) {
|
||
const previous = loadedStates.get(session);
|
||
if (!previous || previous.id !== session.id) {
|
||
throw Object.assign(new Error('只能确认可信加载会话的历史状态'), { code: 'history_untrusted_source' });
|
||
}
|
||
loadedStates.set(session, { ...previous, complete: true, reset: options.reset === true });
|
||
session.historySnapshotBaseIndex = 0;
|
||
session.historySnapshotCount = session.messages.length;
|
||
return () => loadedStates.set(session, previous);
|
||
}
|
||
|
||
function serializeSession(session) {
|
||
// 完整档案没有文本、消息数、数组长度或深度裁剪;不可序列化时整次提交失败。
|
||
const json = JSON.stringify(session, function replacer(key, value) {
|
||
if (this === session && SESSION_RUNTIME_KEYS.has(key)) return undefined;
|
||
if (typeof value === 'bigint') return String(value);
|
||
if (value instanceof Map) return Object.fromEntries(value);
|
||
if (value instanceof Set) return [...value];
|
||
return value;
|
||
});
|
||
const persisted = JSON.parse(json);
|
||
if (!Array.isArray(persisted.messages)) throw new Error('完整会话消息必须为数组');
|
||
return persisted;
|
||
}
|
||
|
||
function equal(left, right) {
|
||
return JSON.stringify(left) === JSON.stringify(right);
|
||
}
|
||
|
||
function plainObject(value) {
|
||
return value && typeof value === 'object' && !Array.isArray(value);
|
||
}
|
||
|
||
function identifiedArray(value) {
|
||
return Array.isArray(value) && value.every(item => item && typeof item.id === 'string' && item.id);
|
||
}
|
||
|
||
function mergeIdentifiedArray(baseline, incoming, current) {
|
||
const baseById = new Map(baseline.map(item => [item.id, item]));
|
||
const incomingById = new Map(incoming.map(item => [item.id, item]));
|
||
const currentById = new Map(current.map(item => [item.id, item]));
|
||
const output = [];
|
||
const seen = new Set();
|
||
for (const item of current) {
|
||
if (seen.has(item.id)) continue;
|
||
seen.add(item.id);
|
||
const base = baseById.get(item.id);
|
||
const next = incomingById.get(item.id);
|
||
if (base && !next && equal(item, base)) continue;
|
||
output.push(next ? mergeValue(base, next, item) : item);
|
||
}
|
||
// 只补本次相对加载基线新增的气泡;后来已经删除的旧气泡不能复活。
|
||
for (let index = 0; index < incoming.length; index += 1) {
|
||
const item = incoming[index];
|
||
if (baseById.has(item.id) || currentById.has(item.id) || output.some(existing => existing.id === item.id)) continue;
|
||
let position = -1;
|
||
for (let next = index + 1; next < incoming.length; next += 1) {
|
||
position = output.findIndex(existing => existing.id === incoming[next].id);
|
||
if (position >= 0) break;
|
||
}
|
||
if (position < 0) {
|
||
for (let previous = index - 1; previous >= 0; previous -= 1) {
|
||
const anchor = output.findIndex(existing => existing.id === incoming[previous].id);
|
||
if (anchor >= 0) { position = anchor + 1; break; }
|
||
}
|
||
}
|
||
output.splice(position < 0 ? output.length : position, 0, item);
|
||
}
|
||
return output;
|
||
}
|
||
|
||
function mergeValue(baseline, incoming, current) {
|
||
if (equal(incoming, baseline)) return current;
|
||
if (equal(current, baseline) || equal(incoming, current)) return incoming;
|
||
if ((baseline === undefined || identifiedArray(baseline)) && identifiedArray(incoming) && identifiedArray(current)) {
|
||
return mergeIdentifiedArray(baseline || [], incoming, current);
|
||
}
|
||
if (plainObject(incoming) && plainObject(current) && (baseline === undefined || plainObject(baseline))) {
|
||
const output = Object.create(null);
|
||
const base = baseline || {};
|
||
for (const key of new Set([...Object.keys(base), ...Object.keys(incoming), ...Object.keys(current)])) {
|
||
const own = (value) => Object.prototype.hasOwnProperty.call(value, key) ? value[key] : undefined;
|
||
const value = mergeValue(own(base), own(incoming), own(current));
|
||
if (value !== undefined) output[key] = value;
|
||
}
|
||
return output;
|
||
}
|
||
// 同字段冲突保留已经提交的值;互不重叠的正文、工具与元数据修改会逐字段合并。
|
||
return current;
|
||
}
|
||
|
||
function synchronizeValue(target, persisted, root = false) {
|
||
if (Array.isArray(target) && Array.isArray(persisted)) {
|
||
const identified = identifiedArray(target) && identifiedArray(persisted);
|
||
const byId = identified ? new Map(target.map(item => [item.id, item])) : null;
|
||
const updated = persisted.map((item, index) => synchronizeValue(byId ? byId.get(item.id) : target[index], item));
|
||
target.length = updated.length;
|
||
for (let index = 0; index < updated.length; index += 1) target[index] = updated[index];
|
||
return target;
|
||
}
|
||
const prototype = target && typeof target === 'object' ? Object.getPrototypeOf(target) : undefined;
|
||
if ((prototype === Object.prototype || prototype === null) && plainObject(persisted)) {
|
||
for (const key of Object.keys(target)) {
|
||
if ((!root || !SESSION_RUNTIME_KEYS.has(key)) && !Object.prototype.hasOwnProperty.call(persisted, key)) delete target[key];
|
||
}
|
||
for (const [key, value] of Object.entries(persisted)) {
|
||
const existing = Object.prototype.hasOwnProperty.call(target, key) ? target[key] : undefined;
|
||
Object.defineProperty(target, key, {
|
||
value: synchronizeValue(existing, value), enumerable: true, writable: true, configurable: true,
|
||
});
|
||
}
|
||
return target;
|
||
}
|
||
return persisted;
|
||
}
|
||
|
||
function writeAtomicJson(target, record) {
|
||
fs.mkdirSync(path.dirname(target), { recursive: true });
|
||
const temporary = `${target}.${process.pid}.${crypto.randomUUID()}.tmp`;
|
||
let fd;
|
||
try {
|
||
fd = fs.openSync(temporary, 'wx', 0o600);
|
||
fs.writeFileSync(fd, JSON.stringify(record), 'utf8');
|
||
fs.fsyncSync(fd);
|
||
fs.closeSync(fd);
|
||
fd = undefined;
|
||
fs.renameSync(temporary, target);
|
||
} finally {
|
||
if (fd !== undefined) {
|
||
try { fs.closeSync(fd); } catch {}
|
||
}
|
||
try { fs.unlinkSync(temporary); } catch {}
|
||
}
|
||
}
|
||
|
||
function writeRecord(id, record) {
|
||
writeAtomicJson(filePath(id), record);
|
||
}
|
||
|
||
function turnPath(id) {
|
||
return path.join(turnsDir, `${validId(id)}.json`);
|
||
}
|
||
|
||
function saveTurn(id, state) {
|
||
writeAtomicJson(turnPath(id), state);
|
||
}
|
||
|
||
function loadTurn(id) {
|
||
try {
|
||
const state = JSON.parse(fs.readFileSync(turnPath(id), 'utf8'));
|
||
if (state?.sessionId !== id || state.agent !== 'codexapp') throw new Error('完整运行态归档无效');
|
||
return state;
|
||
} catch (error) {
|
||
if (error?.code === 'ENOENT') return null;
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
function removeTurn(id) {
|
||
try { fs.unlinkSync(turnPath(id)); } catch (error) {
|
||
if (error?.code !== 'ENOENT') throw error;
|
||
}
|
||
}
|
||
|
||
function listTurnIds() {
|
||
try {
|
||
return fs.readdirSync(turnsDir).filter(name => /^[a-zA-Z0-9_-]+\.json$/.test(name)).map(name => name.slice(0, -5));
|
||
} catch (error) {
|
||
if (error?.code === 'ENOENT') return [];
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
function save(session) {
|
||
const id = validId(session?.id);
|
||
const previous = loadedStates.get(session);
|
||
if (removedIds.has(id)) throw Object.assign(new Error('已删除会话不能由旧对象重新保存'), { code: 'history_removed' });
|
||
const current = readRecord(id);
|
||
if (current) {
|
||
if (previous?.source !== 'archive' || previous.id !== id) {
|
||
throw Object.assign(new Error('受限预览不能覆盖完整会话归档'), { code: 'history_untrusted_source' });
|
||
}
|
||
} else if (previous?.source === 'archive') {
|
||
throw Object.assign(new Error('完整会话归档已移除,拒绝旧对象重新创建'), { code: 'history_removed' });
|
||
}
|
||
const complete = current?.complete || (previous ? previous.complete
|
||
: !session.messages?.some(message => message?.ccwebPersistenceNotice === true)
|
||
&& !(Number(session.historySnapshotBaseIndex) > 0));
|
||
let persisted = serializeSession(session);
|
||
const rebased = !!current && previous.revision !== current.revision;
|
||
if (rebased) {
|
||
const baseline = JSON.parse(previous.baseline);
|
||
if (!identifiedArray(baseline.messages) || !identifiedArray(persisted.messages) || !identifiedArray(current.session.messages)) {
|
||
throw Object.assign(new Error('旧修订缺少可靠消息 ID,拒绝覆盖新历史'), { code: 'history_stale_unidentified' });
|
||
}
|
||
if (baseline.messages.length > 0 && persisted.messages.length === 0 && !equal(current.session.messages, baseline.messages)) {
|
||
throw Object.assign(new Error('清空前历史已经更新,请重新加载'), { code: 'history_stale_clear' });
|
||
}
|
||
if ((current.resetRevision || null) !== (previous.resetRevision || null)) {
|
||
// 清空提交后旧运行对象只能同步元数据,旧线程的历史和输出不能复活。
|
||
persisted.messages = baseline.messages;
|
||
for (const key of ['claudeSessionId', 'codexThreadId', 'codexAppThreadId', 'historySnapshotBaseIndex', 'historySnapshotCount']) {
|
||
if (Object.prototype.hasOwnProperty.call(baseline, key)) persisted[key] = baseline[key];
|
||
else delete persisted[key];
|
||
}
|
||
}
|
||
persisted = mergeValue(baseline, persisted, current.session);
|
||
}
|
||
if (complete) {
|
||
persisted.historySnapshotBaseIndex = 0;
|
||
persisted.historySnapshotCount = persisted.messages.length;
|
||
}
|
||
const revision = crypto.randomUUID();
|
||
const resetRevision = previous?.reset ? revision : current?.resetRevision || null;
|
||
const record = { version: 1, revision, complete, resetRevision, session: persisted };
|
||
writeRecord(id, record);
|
||
// 只有原子提交后才推进修订,失败重试仍以原修订为基线。
|
||
loadedStates.set(session, { id, revision, complete, resetRevision, source: 'archive', baseline: JSON.stringify(persisted) });
|
||
if (rebased) {
|
||
// 流式处理闭包可能持有气泡和工具对象,合并时按 ID 原位协调,保持后续 delta 写入目标。
|
||
synchronizeValue(session, persisted, true);
|
||
}
|
||
return { revision, complete, session: persisted, rebased };
|
||
}
|
||
|
||
function listIds() {
|
||
try {
|
||
return fs.readdirSync(historyDir).filter(name => /^[a-zA-Z0-9_-]+\.json$/.test(name))
|
||
.map(name => name.slice(0, -5));
|
||
} catch (error) {
|
||
if (error?.code === 'ENOENT') return [];
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
function revision(id) {
|
||
let fd;
|
||
try {
|
||
fd = fs.openSync(filePath(id), 'r');
|
||
const head = Buffer.alloc(256);
|
||
const count = fs.readSync(fd, head, 0, head.length, 0);
|
||
const matched = /^\{"version":1,"revision":"([0-9a-f-]{36})",/.exec(head.toString('utf8', 0, count));
|
||
if (!matched) throw new Error('完整会话归档修订无效');
|
||
return matched[1];
|
||
} catch (error) {
|
||
if (error?.code === 'ENOENT') return null;
|
||
throw error;
|
||
} finally {
|
||
if (fd !== undefined) fs.closeSync(fd);
|
||
}
|
||
}
|
||
|
||
function remove(id) {
|
||
id = validId(id);
|
||
try { fs.unlinkSync(filePath(id)); } catch (error) {
|
||
if (error?.code !== 'ENOENT') throw error;
|
||
}
|
||
removedIds.add(id);
|
||
removeTurn(id);
|
||
}
|
||
|
||
return { load, save, state, markLegacy, markComplete, listIds, revision, remove, saveTurn, loadTurn, removeTurn, listTurnIds, historyDir };
|
||
}
|
||
|
||
module.exports = { createSessionHistoryStore };
|