fix(memory): extract from original user messages
Memind CI / Test, build, and release guards (push) Successful in 2m59s
Memind CI / Test, build, and release guards (push) Successful in 2m59s
This commit is contained in:
@@ -47,6 +47,117 @@ export function filterUserVisibleConversation(messages) {
|
|||||||
return messages.filter((message) => message?.metadata?.userVisible !== false);
|
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) {
|
export function countNonEmptyConversationMessages(messages) {
|
||||||
if (!Array.isArray(messages)) return 0;
|
if (!Array.isArray(messages)) return 0;
|
||||||
return messages.filter((message) => {
|
return messages.filter((message) => {
|
||||||
|
|||||||
@@ -5,8 +5,12 @@ import {
|
|||||||
buildConversationFromDbRows,
|
buildConversationFromDbRows,
|
||||||
countNonEmptyConversationMessages,
|
countNonEmptyConversationMessages,
|
||||||
filterUserVisibleConversation,
|
filterUserVisibleConversation,
|
||||||
|
parseAgentRunUserMessage,
|
||||||
parseStoredConversationRow,
|
parseStoredConversationRow,
|
||||||
repairConversationFromDbRows,
|
repairConversationFromDbRows,
|
||||||
|
restoreConversationUserMessagesFromAgentRunRows,
|
||||||
|
restoreConversationUserMessagesFromAgentRuns,
|
||||||
|
restoreConversationUserMessagesFromAgentRunsFailOpen,
|
||||||
shouldRepairConversationFromDb,
|
shouldRepairConversationFromDb,
|
||||||
} from './conversation-repair.mjs';
|
} from './conversation-repair.mjs';
|
||||||
|
|
||||||
@@ -97,3 +101,78 @@ test('filterUserVisibleConversation keeps messages without explicit userVisible
|
|||||||
assert.equal(visible[0].role, 'user');
|
assert.equal(visible[0].role, 'user');
|
||||||
assert.equal(visible[1].role, 'assistant');
|
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/);
|
||||||
|
});
|
||||||
|
|||||||
@@ -13,6 +13,10 @@
|
|||||||
`updated_at=0` 游标开始做有限批次 backfill,新写入行可能长期排在批次之外。结果是
|
`updated_at=0` 游标开始做有限批次 backfill,新写入行可能长期排在批次之外。结果是
|
||||||
`agent_memory_resolved` 显示已注入,但回答只拿到旧的无关记忆。
|
`agent_memory_resolved` 显示已注入,但回答只拿到旧的无关记忆。
|
||||||
|
|
||||||
|
首次修复同步后,灰度又发现 `remember-recent` 读取的是 Agent 编排后的会话文本,其中
|
||||||
|
包含已注入的 `[Memory Context]`。提取器因此可能把旧记忆再次沉淀,而忽略用户本轮明确
|
||||||
|
要求保存的内容,形成旧记忆自我复制。
|
||||||
|
|
||||||
## 必须保留的行为
|
## 必须保留的行为
|
||||||
|
|
||||||
1. `MEMORY_CANDIDATE_PERSISTENCE_ENABLED=1` 且 MySQL 可用时,Portal 必须先执行
|
1. `MEMORY_CANDIDATE_PERSISTENCE_ENABLED=1` 且 MySQL 可用时,Portal 必须先执行
|
||||||
@@ -28,11 +32,14 @@
|
|||||||
6. 用户记忆 `write/compact` 成功后,必须在返回前按 `userId + sessionId` 将本次活跃记忆
|
6. 用户记忆 `write/compact` 成功后,必须在返回前按 `userId + sessionId` 将本次活跃记忆
|
||||||
幂等 upsert 到 pgvector;候选晋升成功后也必须按实际晋升用户同步。不得依赖从零开始
|
幂等 upsert 到 pgvector;候选晋升成功后也必须按实际晋升用户同步。不得依赖从零开始
|
||||||
的全局有限批次 backfill 来保证新记忆可立即召回。
|
的全局有限批次 backfill 来保证新记忆可立即召回。
|
||||||
|
7. 显式记忆提取必须优先使用同一用户、同一会话中已成功 `h5_agent_runs.user_message_json`
|
||||||
|
保存的原始用户消息。不得把 Agent 编排提示或 `[Memory Context]` 当作用户的新记忆;
|
||||||
|
原始消息查询失败时必须 fail-open 到现有可见会话,不能阻塞保存接口。
|
||||||
|
|
||||||
## 回归检查
|
## 回归检查
|
||||||
|
|
||||||
```bash
|
```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
|
memory-v2-pgvector-backfill.test.mjs memory-v2-runtime.test.mjs
|
||||||
npm test
|
npm test
|
||||||
```
|
```
|
||||||
|
|||||||
+13
-2
@@ -199,7 +199,11 @@ import { createFeedbackService } from './user-feedback.mjs';
|
|||||||
import { startScheduleReminderWorker } from './schedule-reminder-worker.mjs';
|
import { startScheduleReminderWorker } from './schedule-reminder-worker.mjs';
|
||||||
import { createLlmProviderService, RELAY_BOOTSTRAP } from './llm-providers.mjs';
|
import { createLlmProviderService, RELAY_BOOTSTRAP } from './llm-providers.mjs';
|
||||||
import { createDirectChatService, isDirectChatSessionId, isPortalDirectChatSnapshot, sendDirectChatSessionEvents, shouldExpirePortalDirectChatSnapshot } from './direct-chat-service.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 { filterNonemptyUserVisibleMessages } from './conversation-transcript-persist.mjs';
|
||||||
import { createSessionStreamStore } from './session-stream-store.mjs';
|
import { createSessionStreamStore } from './session-stream-store.mjs';
|
||||||
import { isSessionStreamReplayEnabled } from './session-stream.mjs';
|
import { isSessionStreamReplayEnabled } from './session-stream.mjs';
|
||||||
@@ -2275,7 +2279,14 @@ async function loadUserVisibleConversation(sessionId, userId) {
|
|||||||
if (authPool && userId) {
|
if (authPool && userId) {
|
||||||
session = await repairSessionConversationFromDb(authPool, session, sessionId, 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) {
|
async function syncUserMemoriesIntoSession(userId, sessionId) {
|
||||||
|
|||||||
Reference in New Issue
Block a user