410 lines
18 KiB
JavaScript
410 lines
18 KiB
JavaScript
const fs = require('fs');
|
||
const path = require('path');
|
||
const crypto = require('crypto');
|
||
|
||
function createCodexRolloutStore(deps) {
|
||
const { codexSessionsDir, sessionsDir, normalizeSession, sanitizeToolInput } = deps;
|
||
|
||
function extractCodexMessageText(content) {
|
||
if (!Array.isArray(content)) return '';
|
||
return content
|
||
.filter((item) => item && (item.type === 'input_text' || item.type === 'output_text'))
|
||
.map((item) => item.text || '')
|
||
.join('');
|
||
}
|
||
|
||
function appendAssistantContent(turn, text) {
|
||
if (!turn || !text || !text.trim()) return;
|
||
turn.content = turn.content ? `${turn.content}\n\n${text}` : text;
|
||
}
|
||
|
||
function stableHash(value) {
|
||
return crypto.createHash('sha256').update(String(value || '')).digest('hex').slice(0, 24);
|
||
}
|
||
|
||
function getRolloutTurnId(entry) {
|
||
const payload = entry?.payload || {};
|
||
const context = payload.turn_context || entry?.turn_context || {};
|
||
const value = payload.turn_id
|
||
|| payload.internal_chat_message_metadata_passthrough?.turn_id
|
||
|| entry?.turn_id
|
||
|| context.turn_id
|
||
|| context.id
|
||
|| (entry?.type === 'turn_context' ? payload.id : null);
|
||
return value ? String(value) : null;
|
||
}
|
||
|
||
function extractCcwebSourceConversation(text) {
|
||
const value = String(text || '').replace(/\r\n/g, '\n');
|
||
const match = value.match(/^来自「([^」]+)」对话(ID:\s*([0-9a-fA-F-]{36}))的消息:/);
|
||
if (match) return { title: match[1], id: match[2].toLowerCase() };
|
||
// 自动继续包装必须同时满足固定开头、请求标识、目标会话和正文结构。
|
||
// 普通用户只是在正文中提到“子对话回传”不能被当作内部消息过滤。
|
||
const returned = value.match(/^子对话回传已返回,但已返回不等于已完成。请先检查返回内容是否完整满足原始请求,再决定继续推进、补问目标对话或向用户汇报。\n\nrequestId:([^\n]+)\n目标对话:「([^」]+)」(ID: ([0-9a-fA-F-]{36}|未知))\n\n原始请求:\n[\s\S]*?\n\n返回正文:\n/);
|
||
if (!returned) return null;
|
||
return { title: returned[2], id: returned[3] === '未知' ? null : returned[3].toLowerCase() };
|
||
}
|
||
|
||
function parseCodexRolloutLines(lines) {
|
||
const entries = [];
|
||
for (const line of lines) {
|
||
try {
|
||
const entry = JSON.parse(line);
|
||
if (entry && typeof entry === 'object') entries.push(entry);
|
||
} catch {}
|
||
}
|
||
const messages = [];
|
||
const pendingToolCalls = new Map();
|
||
const meta = {
|
||
threadId: null,
|
||
cwd: null,
|
||
title: '',
|
||
updatedAt: null,
|
||
cliVersion: null,
|
||
source: null,
|
||
sourceConversationId: null,
|
||
sourceConversationTitle: '',
|
||
};
|
||
const totalUsage = { inputTokens: 0, cachedInputTokens: 0, outputTokens: 0 };
|
||
let currentAssistant = null;
|
||
let currentTurnId = null;
|
||
let currentTurnKey = null;
|
||
let currentSegment = null;
|
||
let implicitTurnSequence = 0;
|
||
let lastUser = null;
|
||
const userMessagesById = new Map();
|
||
const assistantsByTurn = new Map();
|
||
const originalAssistantIds = new WeakMap();
|
||
|
||
function rememberSourceConversation(text) {
|
||
if (meta.sourceConversationId) return;
|
||
const sourceConversation = extractCcwebSourceConversation(text);
|
||
if (!sourceConversation) return;
|
||
meta.sourceConversationId = sourceConversation.id;
|
||
meta.sourceConversationTitle = sourceConversation.title;
|
||
}
|
||
|
||
function messageIdentity(entry) {
|
||
const payload = entry.payload || {};
|
||
const passthrough = payload.internal_chat_message_metadata_passthrough || {};
|
||
const clientMessageId = payload.clientMessageId || payload.client_user_message_id
|
||
|| passthrough.clientMessageId || passthrough.client_user_message_id || null;
|
||
const id = payload.id || payload.message_id || entry.id || entry.message_id || null;
|
||
return { id: id ? String(id) : null, clientMessageId: clientMessageId ? String(clientMessageId) : null };
|
||
}
|
||
|
||
// 有 user_message 事件时,它才是用户输入记录;response_item 中还混有
|
||
// AGENTS、environment 等运行上下文,只有配对成功的记录才可补充身份字段。
|
||
const userEventsByText = new Map();
|
||
const responseUsers = [];
|
||
const responseForUserEvent = new Map();
|
||
let scannedTurnId = null;
|
||
let hasUserEvents = false;
|
||
for (let index = 0; index < entries.length; index += 1) {
|
||
const entry = entries[index];
|
||
scannedTurnId = getRolloutTurnId(entry) || scannedTurnId;
|
||
if (entry.type === 'event_msg' && entry.payload?.type === 'user_message') {
|
||
hasUserEvents = true;
|
||
const text = String(entry.payload.message || '').trim();
|
||
const candidates = userEventsByText.get(text) || [];
|
||
candidates.push({ entry, index, turnId: scannedTurnId });
|
||
userEventsByText.set(text, candidates);
|
||
} else if (entry.type === 'response_item' && entry.payload?.type === 'message' && entry.payload.role === 'user') {
|
||
responseUsers.push({ entry, index, turnId: scannedTurnId });
|
||
}
|
||
if (entry.type === 'event_msg' && ['task_complete', 'task_completed', 'turn_complete', 'turn_completed', 'turn_failed', 'turn_aborted'].includes(entry.payload?.type)) {
|
||
scannedTurnId = null;
|
||
}
|
||
}
|
||
for (const response of responseUsers) {
|
||
const text = extractCodexMessageText(response.entry.payload.content).trim();
|
||
let matched = null;
|
||
for (const candidate of userEventsByText.get(text) || []) {
|
||
if (responseForUserEvent.has(candidate.entry)) continue;
|
||
if (response.turnId && candidate.turnId && response.turnId !== candidate.turnId) continue;
|
||
if (!matched || Math.abs(candidate.index - response.index) < Math.abs(matched.index - response.index)) matched = candidate;
|
||
}
|
||
if (matched) responseForUserEvent.set(matched.entry, response.entry);
|
||
}
|
||
|
||
function assignAssistantIdentity(assistant) {
|
||
const segmentSuffix = currentSegment ? `:segment:${currentSegment}` : '';
|
||
assistant.turnId = currentTurnId;
|
||
assistant.nativeThreadId = meta.threadId;
|
||
assistant.nativeTurnKey = `native:${meta.threadId || 'unknown-thread'}:${currentTurnKey}${segmentSuffix}`;
|
||
assistant.id = originalAssistantIds.get(assistant)
|
||
|| `history-assistant:${stableHash(`${meta.threadId || ''}:${currentTurnKey}${segmentSuffix}`)}`;
|
||
if (currentSegment) assistant.nativeTurnSegment = currentSegment;
|
||
}
|
||
|
||
function flushAssistant() {
|
||
if (currentAssistant && ((currentAssistant.content || '').trim() || currentAssistant.toolCalls.length > 0)) {
|
||
if (!currentAssistant.turnId && !originalAssistantIds.has(currentAssistant)) {
|
||
const content = currentAssistant.content || JSON.stringify(currentAssistant.toolCalls);
|
||
currentAssistant.id = `history-assistant:${stableHash(`${meta.threadId || ''}:${currentAssistant.timestamp || ''}:${content}`)}`;
|
||
}
|
||
messages.push(currentAssistant);
|
||
}
|
||
currentAssistant = null;
|
||
pendingToolCalls.clear();
|
||
}
|
||
|
||
function selectTurn(turnId) {
|
||
if (!turnId || turnId === currentTurnId) return;
|
||
const previousTurnKey = currentTurnKey;
|
||
if (currentTurnId) flushAssistant();
|
||
currentTurnId = turnId;
|
||
currentTurnKey = turnId;
|
||
currentSegment = null;
|
||
// 老格式可能在首段输出后才提供 turn_context,补齐身份而不凭空拆出气泡。
|
||
if (currentAssistant) {
|
||
assistantsByTurn.delete(previousTurnKey);
|
||
assistantsByTurn.set(currentTurnKey, [currentAssistant]);
|
||
assignAssistantIdentity(currentAssistant);
|
||
}
|
||
if (lastUser && !lastUser.assistantStarted && lastUser.turnKey === previousTurnKey) {
|
||
lastUser.turnKey = currentTurnKey;
|
||
lastUser.message.turnId = currentTurnId;
|
||
}
|
||
}
|
||
|
||
function ensureAssistant(entry) {
|
||
selectTurn(getRolloutTurnId(entry));
|
||
if (!currentTurnKey) currentTurnKey = `implicit:${implicitTurnSequence += 1}`;
|
||
if (!currentAssistant) {
|
||
currentAssistant = { role: 'assistant', content: '', toolCalls: [], timestamp: entry.timestamp || null };
|
||
assignAssistantIdentity(currentAssistant);
|
||
const assistants = assistantsByTurn.get(currentTurnKey) || [];
|
||
assistants.push(currentAssistant);
|
||
assistantsByTurn.set(currentTurnKey, assistants);
|
||
}
|
||
if (lastUser) lastUser.assistantStarted = true;
|
||
return currentAssistant;
|
||
}
|
||
|
||
function rememberAssistantMessageId(assistant, entry) {
|
||
const { id } = messageIdentity(entry);
|
||
if (!id) return;
|
||
if (!originalAssistantIds.has(assistant)) {
|
||
originalAssistantIds.set(assistant, id);
|
||
assistant.id = id;
|
||
}
|
||
if (!assistant.nativeMessageIds) assistant.nativeMessageIds = [];
|
||
if (!assistant.nativeMessageIds.includes(id)) assistant.nativeMessageIds.push(id);
|
||
}
|
||
|
||
function appendUser(entry, text) {
|
||
text = String(text || '').trim();
|
||
if (!text) return;
|
||
rememberSourceConversation(text);
|
||
// 内部回传只是 native 的输入包装,不能冒充用户,也不能打断正在聚合的助手。
|
||
if (extractCcwebSourceConversation(text)) return;
|
||
const identity = messageIdentity(entry);
|
||
const turnId = getRolloutTurnId(entry);
|
||
const ts = entry.timestamp || null;
|
||
const known = (identity.id && userMessagesById.get(identity.id))
|
||
|| (identity.clientMessageId && userMessagesById.get(identity.clientMessageId));
|
||
const existing = known;
|
||
if (existing) {
|
||
if (identity.clientMessageId) existing.message.clientMessageId = identity.clientMessageId;
|
||
if (identity.id) {
|
||
if (!existing.hasOriginalId) {
|
||
existing.message.id = identity.id;
|
||
existing.hasOriginalId = true;
|
||
}
|
||
userMessagesById.set(identity.id, existing);
|
||
}
|
||
if (identity.clientMessageId) userMessagesById.set(identity.clientMessageId, existing);
|
||
return;
|
||
}
|
||
|
||
selectTurn(turnId);
|
||
const hadAssistant = (assistantsByTurn.get(currentTurnKey) || []).length > 0;
|
||
flushAssistant();
|
||
if (!currentTurnId) {
|
||
currentTurnKey = `implicit:${implicitTurnSequence += 1}`;
|
||
currentSegment = null;
|
||
}
|
||
const id = identity.id || (identity.clientMessageId ? `client:${identity.clientMessageId}` : null)
|
||
|| `history-user:${stableHash(`${meta.threadId || ''}:${currentTurnId || ''}:${ts || ''}:${text}`)}`;
|
||
if (hadAssistant && currentTurnId) {
|
||
// steer 是真实用户输入:保持前后邻接,同时避免同 turn 的两个助手段互相去重。
|
||
for (const assistant of assistantsByTurn.get(currentTurnKey) || []) {
|
||
if (!assistant.nativeTurnSegment) assistant.nativeTurnSegment = 'initial';
|
||
}
|
||
currentSegment = stableHash(id);
|
||
}
|
||
const message = { role: 'user', content: text, timestamp: ts, id, turnId: currentTurnId, nativeThreadId: meta.threadId };
|
||
if (identity.clientMessageId) message.clientMessageId = identity.clientMessageId;
|
||
messages.push(message);
|
||
lastUser = { message, turnKey: currentTurnKey, assistantStarted: false, hasOriginalId: !!identity.id };
|
||
if (identity.id) userMessagesById.set(identity.id, lastUser);
|
||
if (identity.clientMessageId) userMessagesById.set(identity.clientMessageId, lastUser);
|
||
if (!meta.title) meta.title = text.slice(0, 80).replace(/\n/g, ' ');
|
||
}
|
||
|
||
function appendTool(entry, output = false) {
|
||
const payload = entry.payload || {};
|
||
const assistant = ensureAssistant(entry);
|
||
const fallbackName = payload.type.startsWith('custom_') ? 'CustomToolCall'
|
||
: payload.type.startsWith('mcp_') ? 'McpToolCall' : 'FunctionCall';
|
||
const toolUseId = payload.call_id || payload.id
|
||
|| `native-tool:${stableHash(`${assistant.nativeTurnKey}:${payload.type}:${entry.timestamp || ''}:${JSON.stringify(payload)}`)}`;
|
||
let tool = pendingToolCalls.get(toolUseId);
|
||
if (!tool) {
|
||
tool = { name: payload.name || payload.tool || fallbackName, id: toolUseId, input: null, done: false };
|
||
assistant.toolCalls.push(tool);
|
||
pendingToolCalls.set(toolUseId, tool);
|
||
}
|
||
if (!output) {
|
||
tool.name = payload.name || payload.tool || tool.name;
|
||
tool.input = sanitizeToolInput(tool.name, payload.input ?? payload.arguments ?? '');
|
||
}
|
||
if (output || payload.status === 'completed' || payload.status === 'failed' || payload.result !== undefined) {
|
||
tool.done = true;
|
||
const result = payload.output ?? payload.result;
|
||
if (result !== undefined) tool.result = (typeof result === 'string' ? result : JSON.stringify(result)).slice(0, 2000);
|
||
}
|
||
}
|
||
|
||
for (const entry of entries) {
|
||
const ts = entry.timestamp || null;
|
||
if (ts) meta.updatedAt = ts;
|
||
|
||
if (entry.type === 'session_meta') {
|
||
meta.threadId = entry.payload?.id || meta.threadId;
|
||
meta.cwd = entry.payload?.cwd || meta.cwd;
|
||
meta.cliVersion = entry.payload?.cli_version || meta.cliVersion;
|
||
meta.source = entry.payload?.source || meta.source;
|
||
continue;
|
||
}
|
||
|
||
if ((entry.type === 'event_msg' && ['task_started', 'turn_started'].includes(entry.payload?.type))
|
||
|| entry.type === 'turn_context') {
|
||
selectTurn(getRolloutTurnId(entry));
|
||
if (!currentTurnKey) currentTurnKey = `implicit:${implicitTurnSequence += 1}`;
|
||
continue;
|
||
}
|
||
|
||
if (entry.type === 'event_msg' && ['task_complete', 'task_completed', 'turn_complete', 'turn_completed', 'turn_failed', 'turn_aborted'].includes(entry.payload?.type)) {
|
||
const turnId = getRolloutTurnId(entry);
|
||
// 迟到的上一轮结束事件不能截断当前轮。
|
||
if (!turnId || !currentTurnId || turnId === currentTurnId) {
|
||
flushAssistant();
|
||
currentTurnId = null;
|
||
currentTurnKey = null;
|
||
currentSegment = null;
|
||
}
|
||
continue;
|
||
}
|
||
|
||
if (entry.type === 'event_msg' && entry.payload?.type === 'token_count') {
|
||
const total = entry.payload?.info?.total_token_usage || null;
|
||
const usage = entry.payload?.info?.last_token_usage || null;
|
||
if (total) {
|
||
totalUsage.inputTokens = Math.max(totalUsage.inputTokens, total.input_tokens || 0);
|
||
totalUsage.cachedInputTokens = Math.max(totalUsage.cachedInputTokens, total.cached_input_tokens || 0);
|
||
totalUsage.outputTokens = Math.max(totalUsage.outputTokens, total.output_tokens || 0);
|
||
} else if (usage) {
|
||
totalUsage.inputTokens += usage.input_tokens || 0;
|
||
totalUsage.cachedInputTokens += usage.cached_input_tokens || 0;
|
||
totalUsage.outputTokens += usage.output_tokens || 0;
|
||
}
|
||
continue;
|
||
}
|
||
|
||
if (entry.type === 'event_msg' && entry.payload?.type === 'user_message') {
|
||
const response = responseForUserEvent.get(entry);
|
||
const responseIdentity = response ? messageIdentity(response) : {};
|
||
const eventIdentity = messageIdentity(entry);
|
||
const enrichedEntry = response ? {
|
||
...entry,
|
||
payload: {
|
||
...entry.payload,
|
||
id: eventIdentity.id || responseIdentity.id,
|
||
clientMessageId: eventIdentity.clientMessageId || responseIdentity.clientMessageId,
|
||
turn_id: getRolloutTurnId(entry) || getRolloutTurnId(response),
|
||
},
|
||
} : entry;
|
||
appendUser(enrichedEntry, entry.payload?.message);
|
||
continue;
|
||
}
|
||
|
||
if (entry.type !== 'response_item') continue;
|
||
const payload = entry.payload || {};
|
||
if (payload.type === 'message') {
|
||
const text = extractCodexMessageText(payload.content);
|
||
if (payload.role === 'assistant' && text.trim()) {
|
||
const assistant = ensureAssistant(entry);
|
||
rememberAssistantMessageId(assistant, entry);
|
||
appendAssistantContent(assistant, text);
|
||
} else if (payload.role === 'user' && !hasUserEvents) {
|
||
appendUser(entry, text);
|
||
}
|
||
} else if (['function_call', 'custom_tool_call', 'mcp_tool_call', 'mcp_call'].includes(payload.type)) {
|
||
appendTool(entry);
|
||
} else if (['function_call_output', 'custom_tool_call_output', 'mcp_tool_call_output', 'mcp_call_output'].includes(payload.type)) {
|
||
appendTool(entry, true);
|
||
}
|
||
}
|
||
|
||
flushAssistant();
|
||
return { meta, messages, totalUsage };
|
||
}
|
||
|
||
function walkFiles(dir, files = []) {
|
||
let entries;
|
||
try {
|
||
entries = fs.readdirSync(dir, { withFileTypes: true });
|
||
} catch {
|
||
return files;
|
||
}
|
||
for (const entry of entries) {
|
||
const fullPath = path.join(dir, entry.name);
|
||
if (entry.isDirectory()) walkFiles(fullPath, files);
|
||
else if (entry.isFile()) files.push(fullPath);
|
||
}
|
||
return files;
|
||
}
|
||
|
||
function getCodexRolloutFiles() {
|
||
if (!fs.existsSync(codexSessionsDir)) return [];
|
||
return walkFiles(codexSessionsDir, []).filter((filePath) => filePath.endsWith('.jsonl')).sort().reverse();
|
||
}
|
||
|
||
function getImportedCodexThreadIds(agent = 'codex') {
|
||
const field = agent === 'codexapp' ? 'codexAppThreadId' : 'codexThreadId';
|
||
const imported = new Set();
|
||
try {
|
||
for (const f of fs.readdirSync(sessionsDir).filter((name) => name.endsWith('.json'))) {
|
||
try {
|
||
const session = normalizeSession(JSON.parse(fs.readFileSync(path.join(sessionsDir, f), 'utf8')));
|
||
if (session[field]) imported.add(session[field]);
|
||
} catch {}
|
||
}
|
||
} catch {}
|
||
return imported;
|
||
}
|
||
|
||
function parseCodexRolloutFile(filePath) {
|
||
try {
|
||
const content = fs.readFileSync(filePath, 'utf8');
|
||
const parsed = parseCodexRolloutLines(content.split('\n'));
|
||
parsed.filePath = filePath;
|
||
return parsed;
|
||
} catch {
|
||
return null;
|
||
}
|
||
}
|
||
|
||
return {
|
||
getRolloutTurnId,
|
||
parseCodexRolloutLines,
|
||
getCodexRolloutFiles,
|
||
getImportedCodexThreadIds,
|
||
parseCodexRolloutFile,
|
||
};
|
||
}
|
||
|
||
module.exports = { createCodexRolloutStore };
|