diff --git a/conversation-repair.mjs b/conversation-repair.mjs index 2058f34..0f62f1f 100644 --- a/conversation-repair.mjs +++ b/conversation-repair.mjs @@ -47,6 +47,117 @@ export function filterUserVisibleConversation(messages) { return messages.filter((message) => message?.metadata?.userVisible !== false); } +export function parseAgentRunUserMessage(row) { + const raw = row?.user_message_json ?? row?.userMessageJson; + let parsed = raw; + if (typeof raw === 'string') { + try { + parsed = JSON.parse(raw); + } catch { + return null; + } + } + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) return null; + const message = { + ...parsed, + role: 'user', + metadata: { + ...(parsed.metadata && typeof parsed.metadata === 'object' ? parsed.metadata : {}), + userVisible: true, + agentVisible: true, + memoryInputSource: 'agent-run-original', + }, + }; + return extractConversationMessageText(message) ? message : null; +} + +export function restoreConversationUserMessagesFromAgentRunRows(messages, runRows) { + const conversation = Array.isArray(messages) ? messages : []; + const originals = (Array.isArray(runRows) ? runRows : []) + .map((row) => parseAgentRunUserMessage(row)) + .filter(Boolean); + if (originals.length === 0) return conversation; + + const userIndexes = conversation + .map((message, index) => (message?.role === 'user' ? index : -1)) + .filter((index) => index >= 0); + if (userIndexes.length === 0) return conversation; + + const restored = [...conversation]; + let userOffset = userIndexes.length - 1; + let originalOffset = originals.length - 1; + while (userOffset >= 0 && originalOffset >= 0) { + const messageIndex = userIndexes[userOffset]; + const current = restored[messageIndex] ?? {}; + const original = originals[originalOffset]; + restored[messageIndex] = { + ...current, + role: 'user', + content: original.content, + metadata: { + ...(current.metadata && typeof current.metadata === 'object' ? current.metadata : {}), + ...(original.metadata && typeof original.metadata === 'object' ? original.metadata : {}), + userVisible: current?.metadata?.userVisible ?? true, + agentVisible: current?.metadata?.agentVisible ?? true, + memoryInputSource: 'agent-run-original', + }, + }; + userOffset -= 1; + originalOffset -= 1; + } + return restored; +} + +export async function loadSuccessfulAgentRunUserMessageRows( + pool, + sessionId, + userId, + { limit = 200 } = {}, +) { + if (!pool || !sessionId || !userId) return []; + const resolvedLimit = Math.max(1, Math.min(200, Number(limit) || 200)); + const [rows] = await pool.query( + `SELECT id, user_message_json, created_at + FROM h5_agent_runs + WHERE agent_session_id = ? + AND user_id = ? + AND status = 'succeeded' + ORDER BY created_at DESC, id DESC + LIMIT ?`, + [sessionId, userId, resolvedLimit], + ); + return [...rows].reverse(); +} + +export async function restoreConversationUserMessagesFromAgentRuns(pool, messages, sessionId, userId) { + const userMessageCount = Array.isArray(messages) + ? messages.filter((message) => message?.role === 'user').length + : 0; + if (userMessageCount === 0) return Array.isArray(messages) ? messages : []; + const rows = await loadSuccessfulAgentRunUserMessageRows(pool, sessionId, userId, { + limit: userMessageCount, + }); + return restoreConversationUserMessagesFromAgentRunRows(messages, rows); +} + +export async function restoreConversationUserMessagesFromAgentRunsFailOpen( + pool, + messages, + sessionId, + userId, + { logger = console } = {}, +) { + try { + return await restoreConversationUserMessagesFromAgentRuns(pool, messages, sessionId, userId); + } catch (err) { + logger?.warn?.( + '[memory-v2] original agent-run transcript unavailable; using visible session transcript:', + err instanceof Error ? err.message : err, + ); + return Array.isArray(messages) ? messages : []; + } +} + export function countNonEmptyConversationMessages(messages) { if (!Array.isArray(messages)) return 0; return messages.filter((message) => { diff --git a/conversation-repair.test.mjs b/conversation-repair.test.mjs index 7d18d19..a4b1070 100644 --- a/conversation-repair.test.mjs +++ b/conversation-repair.test.mjs @@ -5,8 +5,12 @@ import { buildConversationFromDbRows, countNonEmptyConversationMessages, filterUserVisibleConversation, + parseAgentRunUserMessage, parseStoredConversationRow, repairConversationFromDbRows, + restoreConversationUserMessagesFromAgentRunRows, + restoreConversationUserMessagesFromAgentRuns, + restoreConversationUserMessagesFromAgentRunsFailOpen, shouldRepairConversationFromDb, } from './conversation-repair.mjs'; @@ -97,3 +101,78 @@ test('filterUserVisibleConversation keeps messages without explicit userVisible assert.equal(visible[0].role, 'user'); assert.equal(visible[1].role, 'assistant'); }); + +test('parseAgentRunUserMessage keeps the original user text for memory extraction', () => { + const parsed = parseAgentRunUserMessage({ + user_message_json: JSON.stringify({ + role: 'user', + content: [{ type: 'text', text: '请记住灰度代号 MEM-NEW' }], + metadata: { userVisible: true }, + }), + }); + assert.equal(parsed.content[0].text, '请记住灰度代号 MEM-NEW'); + assert.equal(parsed.metadata.memoryInputSource, 'agent-run-original'); +}); + +test('restoreConversationUserMessagesFromAgentRunRows removes injected memory context', () => { + const messages = [ + { + id: 'upstream-user-1', + role: 'user', + content: [{ type: 'text', text: '【Memind 任务编排】\n[Memory Context]\n- 旧记忆\n用户任务:\n请记住灰度代号 MEM-NEW' }], + metadata: { userVisible: true, agentVisible: true }, + }, + gooseMsg('assistant-1', 'assistant', '已记住'), + ]; + const restored = restoreConversationUserMessagesFromAgentRunRows(messages, [{ + user_message_json: JSON.stringify({ + role: 'user', + content: [{ type: 'text', text: '请记住灰度代号 MEM-NEW' }], + metadata: { userVisible: true, agentVisible: true }, + }), + }]); + assert.equal(restored[0].id, 'upstream-user-1'); + assert.equal(restored[0].content[0].text, '请记住灰度代号 MEM-NEW'); + assert.doesNotMatch(restored[0].content[0].text, /旧记忆|Memory Context/); + assert.equal(restored[1].content[0].text, '已记住'); +}); + +test('restoreConversationUserMessagesFromAgentRuns scopes successful runs by session and user', async () => { + const calls = []; + const pool = { + async query(sql, params) { + calls.push({ sql, params }); + return [[{ + user_message_json: JSON.stringify({ + role: 'user', + content: [{ type: 'text', text: '原始用户消息' }], + }), + }]]; + }, + }; + const restored = await restoreConversationUserMessagesFromAgentRuns( + pool, + [gooseMsg('user-1', 'user', '编排后的消息')], + 'session-1', + 'user-1', + ); + assert.match(calls[0].sql, /status = 'succeeded'/); + assert.match(calls[0].sql, /ORDER BY created_at DESC/); + assert.match(calls[0].sql, /LIMIT \?/); + assert.deepEqual(calls[0].params, ['session-1', 'user-1', 1]); + assert.equal(restored[0].content[0].text, '原始用户消息'); +}); + +test('restoreConversationUserMessagesFromAgentRunsFailOpen preserves the visible transcript', async () => { + const warnings = []; + const messages = [gooseMsg('user-1', 'user', '可见会话原文')]; + const restored = await restoreConversationUserMessagesFromAgentRunsFailOpen( + { async query() { throw new Error('agent run query unavailable'); } }, + messages, + 'session-1', + 'user-1', + { logger: { warn: (...items) => warnings.push(items.join(' ')) } }, + ); + assert.equal(restored, messages); + assert.match(warnings[0], /agent run query unavailable/); +}); diff --git a/docs/regression-guards/memory-v2-candidate-and-lifecycle.md b/docs/regression-guards/memory-v2-candidate-and-lifecycle.md index 3feba84..db7a5c6 100644 --- a/docs/regression-guards/memory-v2-candidate-and-lifecycle.md +++ b/docs/regression-guards/memory-v2-candidate-and-lifecycle.md @@ -13,6 +13,10 @@ `updated_at=0` 游标开始做有限批次 backfill,新写入行可能长期排在批次之外。结果是 `agent_memory_resolved` 显示已注入,但回答只拿到旧的无关记忆。 +首次修复同步后,灰度又发现 `remember-recent` 读取的是 Agent 编排后的会话文本,其中 +包含已注入的 `[Memory Context]`。提取器因此可能把旧记忆再次沉淀,而忽略用户本轮明确 +要求保存的内容,形成旧记忆自我复制。 + ## 必须保留的行为 1. `MEMORY_CANDIDATE_PERSISTENCE_ENABLED=1` 且 MySQL 可用时,Portal 必须先执行 @@ -28,11 +32,14 @@ 6. 用户记忆 `write/compact` 成功后,必须在返回前按 `userId + sessionId` 将本次活跃记忆 幂等 upsert 到 pgvector;候选晋升成功后也必须按实际晋升用户同步。不得依赖从零开始 的全局有限批次 backfill 来保证新记忆可立即召回。 +7. 显式记忆提取必须优先使用同一用户、同一会话中已成功 `h5_agent_runs.user_message_json` + 保存的原始用户消息。不得把 Agent 编排提示或 `[Memory Context]` 当作用户的新记忆; + 原始消息查询失败时必须 fail-open 到现有可见会话,不能阻塞保存接口。 ## 回归检查 ```bash -node --test memory-v2-personal-store.test.mjs memory-v2-lifecycle.test.mjs \ +node --test conversation-repair.test.mjs memory-v2-personal-store.test.mjs memory-v2-lifecycle.test.mjs \ memory-v2-pgvector-backfill.test.mjs memory-v2-runtime.test.mjs npm test ``` diff --git a/server.mjs b/server.mjs index 09001d9..fc151f5 100644 --- a/server.mjs +++ b/server.mjs @@ -199,7 +199,11 @@ import { createFeedbackService } from './user-feedback.mjs'; import { startScheduleReminderWorker } from './schedule-reminder-worker.mjs'; import { createLlmProviderService, RELAY_BOOTSTRAP } from './llm-providers.mjs'; import { createDirectChatService, isDirectChatSessionId, isPortalDirectChatSnapshot, sendDirectChatSessionEvents, shouldExpirePortalDirectChatSnapshot } from './direct-chat-service.mjs'; -import { filterUserVisibleConversation, repairSessionConversationFromDb } from './conversation-repair.mjs'; +import { + filterUserVisibleConversation, + repairSessionConversationFromDb, + restoreConversationUserMessagesFromAgentRunsFailOpen, +} from './conversation-repair.mjs'; import { filterNonemptyUserVisibleMessages } from './conversation-transcript-persist.mjs'; import { createSessionStreamStore } from './session-stream-store.mjs'; import { isSessionStreamReplayEnabled } from './session-stream.mjs'; @@ -2275,7 +2279,14 @@ async function loadUserVisibleConversation(sessionId, userId) { if (authPool && userId) { session = await repairSessionConversationFromDb(authPool, session, sessionId, userId); } - return filterUserVisibleConversation(session?.conversation ?? []); + const visible = filterUserVisibleConversation(session?.conversation ?? []); + if (!authPool || !userId) return visible; + return restoreConversationUserMessagesFromAgentRunsFailOpen( + authPool, + visible, + sessionId, + userId, + ); } async function syncUserMemoriesIntoSession(userId, sessionId) {