Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 57d2c14fca | |||
| d52aab8ab0 | |||
| c10788dfe1 | |||
| 1f9e147dea | |||
| 7b3acb6813 | |||
| aafda0cf64 |
@@ -109,6 +109,7 @@ function buildMemoryPrompt(messages) {
|
|||||||
'你是 TKMind 的用户长期记忆提取器。',
|
'你是 TKMind 的用户长期记忆提取器。',
|
||||||
'请从下面的用户对话中提取可以长期复用的个人记忆。',
|
'请从下面的用户对话中提取可以长期复用的个人记忆。',
|
||||||
'只记录用户明确表达或强证据支持的信息,不要猜测,不要记录临时闲聊。',
|
'只记录用户明确表达或强证据支持的信息,不要猜测,不要记录临时闲聊。',
|
||||||
|
'不要把问题、请求、指令或待办本身当作用户事实;问句没有提供答案时不要记录。',
|
||||||
'不要记录密码、密钥、身份证、手机号、银行卡等敏感信息。',
|
'不要记录密码、密钥、身份证、手机号、银行卡等敏感信息。',
|
||||||
'只返回 JSON 对象,不要 Markdown,不要解释。',
|
'只返回 JSON 对象,不要 Markdown,不要解释。',
|
||||||
'格式:{"memories":[{"label":"preference|habit|interest|goal|fact|experience|knowledge","text":"一句完整中文记忆","confidence":0.0-1.0}]}',
|
'格式:{"memories":[{"label":"preference|habit|interest|goal|fact|experience|knowledge","text":"一句完整中文记忆","confidence":0.0-1.0}]}',
|
||||||
@@ -382,7 +383,9 @@ export function createConversationMemoryService(pool, options = {}) {
|
|||||||
warnLlmExtractionFailed(err);
|
warnLlmExtractionFailed(err);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (!memories?.length) memories = fallbackMemoriesFromMessages(messages);
|
// An empty array is an authoritative LLM decision that the batch contains no
|
||||||
|
// durable memory. Only fall back when extraction was unavailable (`null`).
|
||||||
|
if (memories == null) memories = fallbackMemoriesFromMessages(messages);
|
||||||
const stored = await storeMemories(userId, messages, memories);
|
const stored = await storeMemories(userId, messages, memories);
|
||||||
await markAnalyzed(messages.map((message) => message.id));
|
await markAnalyzed(messages.map((message) => message.id));
|
||||||
return { ok: true, analyzed: messages.length, memories: stored };
|
return { ok: true, analyzed: messages.length, memories: stored };
|
||||||
|
|||||||
@@ -208,6 +208,35 @@ test('saveAndAnalyze marks messages analyzed when llm extraction fails and no me
|
|||||||
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('saveAndAnalyze still uses fallback when llm extraction is unavailable', async () => {
|
||||||
|
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||||
|
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '1';
|
||||||
|
const pool = createPool();
|
||||||
|
const service = createConversationMemoryService(pool, {
|
||||||
|
now: () => 2250,
|
||||||
|
llmProviderService: {
|
||||||
|
async createChatCompletion() {
|
||||||
|
throw new Error('upstream unavailable');
|
||||||
|
},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await service.saveAndAnalyze('session-fallback', 'user-fallback', [
|
||||||
|
{
|
||||||
|
id: 'm-fallback',
|
||||||
|
role: 'user',
|
||||||
|
content: [{ type: 'text', text: '我喜欢简洁的回答。' }],
|
||||||
|
metadata: { userVisible: true },
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
|
||||||
|
assert.equal(result.memories, 1);
|
||||||
|
assert.match(pool.state.memories[0].memory_text, /简洁/);
|
||||||
|
|
||||||
|
if (previous == null) delete process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||||
|
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
||||||
|
});
|
||||||
|
|
||||||
test('saveAndAnalyze uses admin effective env for memory extraction model', async () => {
|
test('saveAndAnalyze uses admin effective env for memory extraction model', async () => {
|
||||||
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||||
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '1';
|
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '1';
|
||||||
@@ -284,6 +313,40 @@ test('saveAndAnalyze extracts memories through llmProviderService', async () =>
|
|||||||
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('saveAndAnalyze respects an empty llm result instead of storing a question through fallback', async () => {
|
||||||
|
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||||
|
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '1';
|
||||||
|
const pool = createPool();
|
||||||
|
const service = createConversationMemoryService(pool, {
|
||||||
|
now: () => 2750,
|
||||||
|
llmProviderService: {
|
||||||
|
async createChatCompletion({ messages }) {
|
||||||
|
assert.match(String(messages?.[0]?.content ?? ''), /问句没有提供答案时不要记录/);
|
||||||
|
return {
|
||||||
|
ok: true,
|
||||||
|
reply: JSON.stringify({ memories: [] }),
|
||||||
|
};
|
||||||
|
},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await service.saveAndAnalyze('session-question', 'user-question', [
|
||||||
|
{
|
||||||
|
id: 'm-question',
|
||||||
|
role: 'user',
|
||||||
|
content: [{ type: 'text', text: '用户记住的记忆召回灰度测试代号是什么?' }],
|
||||||
|
metadata: { userVisible: true },
|
||||||
|
},
|
||||||
|
]);
|
||||||
|
|
||||||
|
assert.equal(result.analyzed, 1);
|
||||||
|
assert.equal(result.memories, 0);
|
||||||
|
assert.equal(pool.state.memories.length, 0);
|
||||||
|
|
||||||
|
if (previous == null) delete process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||||
|
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
||||||
|
});
|
||||||
|
|
||||||
test('saveAndAnalyze throttles repeated llm extraction warnings', async () => {
|
test('saveAndAnalyze throttles repeated llm extraction warnings', async () => {
|
||||||
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||||
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '1';
|
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '1';
|
||||||
|
|||||||
@@ -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,18 @@
|
|||||||
`updated_at=0` 游标开始做有限批次 backfill,新写入行可能长期排在批次之外。结果是
|
`updated_at=0` 游标开始做有限批次 backfill,新写入行可能长期排在批次之外。结果是
|
||||||
`agent_memory_resolved` 显示已注入,但回答只拿到旧的无关记忆。
|
`agent_memory_resolved` 显示已注入,但回答只拿到旧的无关记忆。
|
||||||
|
|
||||||
|
首次修复同步后,灰度又发现 `remember-recent` 读取的是 Agent 编排后的会话文本,其中
|
||||||
|
包含已注入的 `[Memory Context]`。提取器因此可能把旧记忆再次沉淀,而忽略用户本轮明确
|
||||||
|
要求保存的内容,形成旧记忆自我复制。
|
||||||
|
|
||||||
|
提取器返回 `{"memories":[]}` 时,旧逻辑仍会触发规则回退,可能把“……是什么”一类
|
||||||
|
问句误存为事实。空数组必须视为成功的“无需保存”判断;只有提取不可用并返回 `null`
|
||||||
|
时才允许规则回退。
|
||||||
|
|
||||||
|
当前生产使用的本地 hash embedding 只有 3 维,且连续中文可能被视为单个 token。即使正确
|
||||||
|
记忆已进入 pgvector,它也可能排在向量 Top-K 之外。因此 pgvector 读取必须合并有界的
|
||||||
|
向量候选与最近候选,再以中文字符 n-gram 查询覆盖率重排;无词法重合时保持向量排序。
|
||||||
|
|
||||||
## 必须保留的行为
|
## 必须保留的行为
|
||||||
|
|
||||||
1. `MEMORY_CANDIDATE_PERSISTENCE_ENABLED=1` 且 MySQL 可用时,Portal 必须先执行
|
1. `MEMORY_CANDIDATE_PERSISTENCE_ENABLED=1` 且 MySQL 可用时,Portal 必须先执行
|
||||||
@@ -28,11 +40,18 @@
|
|||||||
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 到现有可见会话,不能阻塞保存接口。
|
||||||
|
8. pgvector 召回必须同时覆盖有界向量候选和有界最近候选,去重后只返回请求的 limit;
|
||||||
|
中文词法重排用于弥补低维本地 hash 的排序缺陷,且候选池上限不得超过 100。
|
||||||
|
9. LLM 提取明确返回空数组时不得再走规则回退;提取提示必须明确排除没有提供答案的
|
||||||
|
问题、请求、指令和待办,避免把问句本身沉淀为长期记忆。
|
||||||
|
|
||||||
## 回归检查
|
## 回归检查
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
node --test memory-v2-personal-store.test.mjs memory-v2-lifecycle.test.mjs \
|
node --test conversation-memory.test.mjs 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
|
||||||
```
|
```
|
||||||
|
|||||||
+101
-7
@@ -24,6 +24,43 @@ function vectorLiteral(embedding) {
|
|||||||
return `[${embedding.join(',')}]`;
|
return `[${embedding.join(',')}]`;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function normalizeSearchText(value) {
|
||||||
|
return String(value ?? '')
|
||||||
|
.normalize('NFKC')
|
||||||
|
.toLowerCase()
|
||||||
|
.replace(/[^a-z0-9\u4e00-\u9fff]+/gu, '');
|
||||||
|
}
|
||||||
|
|
||||||
|
function buildCharacterNgrams(value, size = 2) {
|
||||||
|
const text = normalizeSearchText(value);
|
||||||
|
if (!text) return new Set();
|
||||||
|
if (text.length <= size) return new Set([text]);
|
||||||
|
const grams = new Set();
|
||||||
|
for (let index = 0; index <= text.length - size; index += 1) {
|
||||||
|
grams.add(text.slice(index, index + size));
|
||||||
|
}
|
||||||
|
return grams;
|
||||||
|
}
|
||||||
|
|
||||||
|
function lexicalQueryCoverage(query, text) {
|
||||||
|
const queryGrams = buildCharacterNgrams(query);
|
||||||
|
if (queryGrams.size === 0) return 0;
|
||||||
|
const textGrams = buildCharacterNgrams(text);
|
||||||
|
let overlap = 0;
|
||||||
|
for (const gram of queryGrams) {
|
||||||
|
if (textGrams.has(gram)) overlap += 1;
|
||||||
|
}
|
||||||
|
return overlap / queryGrams.size;
|
||||||
|
}
|
||||||
|
|
||||||
|
function timestampValue(value) {
|
||||||
|
if (value == null) return 0;
|
||||||
|
const numeric = Number(value);
|
||||||
|
if (Number.isFinite(numeric)) return numeric;
|
||||||
|
const parsed = Date.parse(String(value));
|
||||||
|
return Number.isFinite(parsed) ? parsed : 0;
|
||||||
|
}
|
||||||
|
|
||||||
function normalizeRow(row) {
|
function normalizeRow(row) {
|
||||||
const text = String(row?.content ?? row?.memory_text ?? row?.text ?? '').trim();
|
const text = String(row?.content ?? row?.memory_text ?? row?.text ?? '').trim();
|
||||||
if (!text) return null;
|
if (!text) return null;
|
||||||
@@ -33,9 +70,38 @@ function normalizeRow(row) {
|
|||||||
text,
|
text,
|
||||||
score: row?.score == null ? null : Number(row.score),
|
score: row?.score == null ? null : Number(row.score),
|
||||||
createdAt: row?.created_at ?? row?.createdAt ?? null,
|
createdAt: row?.created_at ?? row?.createdAt ?? null,
|
||||||
|
updatedAt: row?.updated_at ?? row?.updatedAt ?? row?.created_at ?? row?.createdAt ?? null,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function rankHybridCandidates(rows, query, limit) {
|
||||||
|
const byId = new Map();
|
||||||
|
for (const row of rows ?? []) {
|
||||||
|
const memory = normalizeRow(row);
|
||||||
|
if (!memory) continue;
|
||||||
|
const key = memory.id ?? `${memory.label}:${memory.text}`;
|
||||||
|
if (!byId.has(key)) byId.set(key, memory);
|
||||||
|
}
|
||||||
|
return [...byId.values()]
|
||||||
|
.map((memory) => ({
|
||||||
|
memory,
|
||||||
|
lexicalScore: lexicalQueryCoverage(query, memory.text),
|
||||||
|
vectorScore: Number.isFinite(memory.score) ? memory.score : -1,
|
||||||
|
updatedAt: timestampValue(memory.updatedAt),
|
||||||
|
}))
|
||||||
|
.sort((left, right) => {
|
||||||
|
if (left.lexicalScore !== right.lexicalScore) {
|
||||||
|
return right.lexicalScore - left.lexicalScore;
|
||||||
|
}
|
||||||
|
if (left.lexicalScore > 0 && left.updatedAt !== right.updatedAt) {
|
||||||
|
return right.updatedAt - left.updatedAt;
|
||||||
|
}
|
||||||
|
return right.vectorScore - left.vectorScore;
|
||||||
|
})
|
||||||
|
.slice(0, limit)
|
||||||
|
.map(({ memory }) => memory);
|
||||||
|
}
|
||||||
|
|
||||||
export function createPgvectorMemoryBackend({
|
export function createPgvectorMemoryBackend({
|
||||||
pool = null,
|
pool = null,
|
||||||
enabled = false,
|
enabled = false,
|
||||||
@@ -75,15 +141,38 @@ export function createPgvectorMemoryBackend({
|
|||||||
const embedding = await resolveEmbedding(input);
|
const embedding = await resolveEmbedding(input);
|
||||||
if (!embedding) return { memories: [], semanticMemories: [] };
|
if (!embedding) return { memories: [], semanticMemories: [] };
|
||||||
const limit = Math.max(1, Math.min(50, Number(input.limit ?? defaultLimit) || defaultLimit));
|
const limit = Math.max(1, Math.min(50, Number(input.limit ?? defaultLimit) || defaultLimit));
|
||||||
|
const candidateLimit = Math.max(
|
||||||
|
limit,
|
||||||
|
Math.min(100, Number(input.candidateLimit ?? 50) || 50),
|
||||||
|
);
|
||||||
const sql = `
|
const sql = `
|
||||||
SELECT id, content, type, created_at, 1 - (embedding <=> $2::vector) AS score
|
WITH vector_candidates AS (
|
||||||
FROM ${resolvedTableName}
|
SELECT id, content, type, created_at, updated_at,
|
||||||
WHERE user_id = $1
|
1 - (embedding <=> $2::vector) AS score,
|
||||||
ORDER BY embedding <=> $2::vector
|
0 AS source_priority
|
||||||
LIMIT $3
|
FROM ${resolvedTableName}
|
||||||
|
WHERE user_id = $1
|
||||||
|
ORDER BY embedding <=> $2::vector
|
||||||
|
LIMIT $3
|
||||||
|
), recent_candidates AS (
|
||||||
|
SELECT id, content, type, created_at, updated_at,
|
||||||
|
1 - (embedding <=> $2::vector) AS score,
|
||||||
|
1 AS source_priority
|
||||||
|
FROM ${resolvedTableName}
|
||||||
|
WHERE user_id = $1
|
||||||
|
ORDER BY updated_at DESC
|
||||||
|
LIMIT $3
|
||||||
|
)
|
||||||
|
SELECT DISTINCT ON (id) id, content, type, created_at, updated_at, score
|
||||||
|
FROM (
|
||||||
|
SELECT * FROM vector_candidates
|
||||||
|
UNION ALL
|
||||||
|
SELECT * FROM recent_candidates
|
||||||
|
) AS candidates
|
||||||
|
ORDER BY id, source_priority
|
||||||
`;
|
`;
|
||||||
const result = await pool.query(sql, [userId, vectorLiteral(embedding), limit]);
|
const result = await pool.query(sql, [userId, vectorLiteral(embedding), candidateLimit]);
|
||||||
const memories = (result?.rows ?? []).map((row) => normalizeRow(row)).filter(Boolean);
|
const memories = rankHybridCandidates(result?.rows ?? [], input.query, limit);
|
||||||
return {
|
return {
|
||||||
semanticMemories: memories.map((item) => item.text),
|
semanticMemories: memories.map((item) => item.text),
|
||||||
memories,
|
memories,
|
||||||
@@ -91,3 +180,8 @@ export function createPgvectorMemoryBackend({
|
|||||||
},
|
},
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export const pgvectorMemoryBackendInternals = {
|
||||||
|
lexicalQueryCoverage,
|
||||||
|
rankHybridCandidates,
|
||||||
|
};
|
||||||
|
|||||||
@@ -1,7 +1,10 @@
|
|||||||
import assert from 'node:assert/strict';
|
import assert from 'node:assert/strict';
|
||||||
import test from 'node:test';
|
import test from 'node:test';
|
||||||
import { createMemoryV2 } from './memory-v2.mjs';
|
import { createMemoryV2 } from './memory-v2.mjs';
|
||||||
import { createPgvectorMemoryBackend } from './memory-v2-pgvector.mjs';
|
import {
|
||||||
|
createPgvectorMemoryBackend,
|
||||||
|
pgvectorMemoryBackendInternals,
|
||||||
|
} from './memory-v2-pgvector.mjs';
|
||||||
|
|
||||||
test('pgvector backend is disabled by default and does not query storage', async () => {
|
test('pgvector backend is disabled by default and does not query storage', async () => {
|
||||||
let queried = false;
|
let queried = false;
|
||||||
@@ -85,7 +88,9 @@ test('pgvector backend performs parameterized vector lookup when explicitly enab
|
|||||||
|
|
||||||
assert.equal(queries.length, 1);
|
assert.equal(queries.length, 1);
|
||||||
assert.match(queries[0].sql, /FROM memory_embeddings/);
|
assert.match(queries[0].sql, /FROM memory_embeddings/);
|
||||||
assert.deepEqual(queries[0].params, ['user-1', '[0.25,0.5,0.75]', 5]);
|
assert.match(queries[0].sql, /WITH vector_candidates/);
|
||||||
|
assert.match(queries[0].sql, /recent_candidates/);
|
||||||
|
assert.deepEqual(queries[0].params, ['user-1', '[0.25,0.5,0.75]', 50]);
|
||||||
assert.deepEqual(result.semanticMemories, ['用户关注 Memory V2 的 facade 边界']);
|
assert.deepEqual(result.semanticMemories, ['用户关注 Memory V2 的 facade 边界']);
|
||||||
assert.deepEqual(result.memories, [
|
assert.deepEqual(result.memories, [
|
||||||
{
|
{
|
||||||
@@ -94,10 +99,68 @@ test('pgvector backend performs parameterized vector lookup when explicitly enab
|
|||||||
text: '用户关注 Memory V2 的 facade 边界',
|
text: '用户关注 Memory V2 的 facade 边界',
|
||||||
score: 0.87,
|
score: 0.87,
|
||||||
createdAt: '2026-07-02T00:00:00.000Z',
|
createdAt: '2026-07-02T00:00:00.000Z',
|
||||||
|
updatedAt: '2026-07-02T00:00:00.000Z',
|
||||||
},
|
},
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('pgvector hybrid ranking recovers a recent Chinese memory missed by vector top-k', async () => {
|
||||||
|
const marker = 'MEM-RECALL-NEW';
|
||||||
|
const backend = createPgvectorMemoryBackend({
|
||||||
|
enabled: true,
|
||||||
|
pool: {
|
||||||
|
async query() {
|
||||||
|
return {
|
||||||
|
rows: [
|
||||||
|
{
|
||||||
|
id: 1,
|
||||||
|
content: '用户以前关注贵州旅游攻略',
|
||||||
|
type: 'interest',
|
||||||
|
score: 0.91,
|
||||||
|
created_at: '2026-06-01T00:00:00.000Z',
|
||||||
|
updated_at: '2026-06-01T00:00:00.000Z',
|
||||||
|
},
|
||||||
|
{
|
||||||
|
id: 2,
|
||||||
|
content: `用户的记忆召回灰度测试代号是 ${marker}`,
|
||||||
|
type: 'fact',
|
||||||
|
score: -0.49,
|
||||||
|
created_at: '2026-07-22T00:00:00.000Z',
|
||||||
|
updated_at: '2026-07-22T00:00:00.000Z',
|
||||||
|
},
|
||||||
|
{
|
||||||
|
id: 3,
|
||||||
|
content: '用户偏好简洁回答',
|
||||||
|
type: 'preference',
|
||||||
|
score: 0.3,
|
||||||
|
created_at: '2026-07-21T00:00:00.000Z',
|
||||||
|
updated_at: '2026-07-21T00:00:00.000Z',
|
||||||
|
},
|
||||||
|
],
|
||||||
|
};
|
||||||
|
},
|
||||||
|
},
|
||||||
|
embedQuery: async () => [0.25, 0.5, 0.75],
|
||||||
|
});
|
||||||
|
|
||||||
|
const result = await backend.resolve({
|
||||||
|
userId: 'user-1',
|
||||||
|
query: '我之前让你记住的记忆召回灰度测试代号是什么?请只回答完整代号。',
|
||||||
|
limit: 3,
|
||||||
|
});
|
||||||
|
|
||||||
|
assert.match(result.memories[0].text, new RegExp(marker));
|
||||||
|
});
|
||||||
|
|
||||||
|
test('pgvector hybrid ranking keeps vector order when query has no lexical overlap', () => {
|
||||||
|
const ranked = pgvectorMemoryBackendInternals.rankHybridCandidates([
|
||||||
|
{ id: 1, content: 'alpha', score: 0.2 },
|
||||||
|
{ id: 2, content: 'beta', score: 0.8 },
|
||||||
|
], '完全无关的中文查询', 2);
|
||||||
|
assert.equal(ranked[0].id, '2');
|
||||||
|
assert.equal(ranked[1].id, '1');
|
||||||
|
});
|
||||||
|
|
||||||
test('pgvector backend validates table names before building SQL', () => {
|
test('pgvector backend validates table names before building SQL', () => {
|
||||||
assert.throws(
|
assert.throws(
|
||||||
() => createPgvectorMemoryBackend({ tableName: 'memory_embeddings;DROP TABLE users' }),
|
() => createPgvectorMemoryBackend({ tableName: 'memory_embeddings;DROP TABLE users' }),
|
||||||
|
|||||||
@@ -385,7 +385,8 @@ test('createMemoryV2Runtime selects pgvector only when pool and embedding are co
|
|||||||
assert.equal(queries.length, 1);
|
assert.equal(queries.length, 1);
|
||||||
assert.equal(queries[0].options.connectionString, 'postgresql://local/memory');
|
assert.equal(queries[0].options.connectionString, 'postgresql://local/memory');
|
||||||
assert.equal(queries[0].options.max, 2);
|
assert.equal(queries[0].options.max, 2);
|
||||||
assert.deepEqual(queries[0].params, ['u1', '[0.1,0.2,0.3]', 8]);
|
assert.match(queries[0].sql, /recent_candidates/);
|
||||||
|
assert.deepEqual(queries[0].params, ['u1', '[0.1,0.2,0.3]', 50]);
|
||||||
|
|
||||||
await memory.close();
|
await memory.close();
|
||||||
assert.equal(poolEnded, true);
|
assert.equal(poolEnded, true);
|
||||||
|
|||||||
+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) {
|
||||||
|
|||||||
+10
-5
@@ -1030,6 +1030,7 @@ export function isRecoverableWechatAgentSessionError(message) {
|
|||||||
const normalized = String(message ?? '').trim();
|
const normalized = String(message ?? '').trim();
|
||||||
if (!normalized) return false;
|
if (!normalized) return false;
|
||||||
if (/stale_session_poisoned_completion/i.test(normalized)) return true;
|
if (/stale_session_poisoned_completion/i.test(normalized)) return true;
|
||||||
|
if (/wechat_page_fresh_thumbnail_required:/i.test(normalized)) return true;
|
||||||
if (/403|404|not found|无权访问/i.test(normalized)) return true;
|
if (/403|404|not found|无权访问/i.test(normalized)) return true;
|
||||||
if (/tool_calls|tool_call_id|insufficient tool messages/i.test(normalized)) return true;
|
if (/tool_calls|tool_call_id|insufficient tool messages/i.test(normalized)) return true;
|
||||||
if (/session already has an active request|active request.*cancel/i.test(normalized)) return true;
|
if (/session already has an active request|active request.*cancel/i.test(normalized)) return true;
|
||||||
@@ -1840,6 +1841,7 @@ export function createWechatMpService({
|
|||||||
user,
|
user,
|
||||||
imagePolicy,
|
imagePolicy,
|
||||||
publishDir,
|
publishDir,
|
||||||
|
notifyFailure = true,
|
||||||
}) => {
|
}) => {
|
||||||
if (imagePolicy?.pageThumbnailMode !== WECHAT_PAGE_THUMBNAIL_MODE.REQUIRED_FRESH) return;
|
if (imagePolicy?.pageThumbnailMode !== WECHAT_PAGE_THUMBNAIL_MODE.REQUIRED_FRESH) return;
|
||||||
const images = collectWechatGeneratedImages(reply?.messages ?? []);
|
const images = collectWechatGeneratedImages(reply?.messages ?? []);
|
||||||
@@ -1860,14 +1862,16 @@ export function createWechatMpService({
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
const text = buildPagePublishFailureText({ missingFreshThumbnail: true });
|
const text = buildPagePublishFailureText({ missingFreshThumbnail: true });
|
||||||
try {
|
if (notifyFailure) {
|
||||||
await sendCustomerServiceText(openid, text, user);
|
try {
|
||||||
} catch (sendErr) {
|
await sendCustomerServiceText(openid, text, user);
|
||||||
logger.error?.('WeChat MP fresh thumbnail failure notice failed:', sendErr);
|
} catch (sendErr) {
|
||||||
|
logger.error?.('WeChat MP fresh thumbnail failure notice failed:', sendErr);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
const error = new Error(`wechat_page_fresh_thumbnail_required:${verification.reason}`);
|
const error = new Error(`wechat_page_fresh_thumbnail_required:${verification.reason}`);
|
||||||
error.code = 'WECHAT_PAGE_FRESH_THUMBNAIL_REQUIRED';
|
error.code = 'WECHAT_PAGE_FRESH_THUMBNAIL_REQUIRED';
|
||||||
throw markWechatUserNotified(error);
|
throw notifyFailure ? markWechatUserNotified(error) : error;
|
||||||
};
|
};
|
||||||
|
|
||||||
const ensureSessionProvider = async (sessionId) => {
|
const ensureSessionProvider = async (sessionId) => {
|
||||||
@@ -2351,6 +2355,7 @@ export function createWechatMpService({
|
|||||||
user,
|
user,
|
||||||
imagePolicy,
|
imagePolicy,
|
||||||
publishDir: workingDir,
|
publishDir: workingDir,
|
||||||
|
notifyFailure: false,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
const pageDataOutcome = await enforcePageDataCollectDelivery({
|
const pageDataOutcome = await enforcePageDataCollectDelivery({
|
||||||
|
|||||||
@@ -2723,6 +2723,12 @@ test('isRecoverableWechatAgentSessionError detects poisoned tool_calls history',
|
|||||||
true,
|
true,
|
||||||
);
|
);
|
||||||
assert.equal(isRecoverableWechatAgentSessionError('stale_session_poisoned_completion'), true);
|
assert.equal(isRecoverableWechatAgentSessionError('stale_session_poisoned_completion'), true);
|
||||||
|
assert.equal(
|
||||||
|
isRecoverableWechatAgentSessionError(
|
||||||
|
'wechat_page_fresh_thumbnail_required:fresh_image_not_generated',
|
||||||
|
),
|
||||||
|
true,
|
||||||
|
);
|
||||||
assert.equal(isRecoverableWechatAgentSessionError('无权访问该会话'), true);
|
assert.equal(isRecoverableWechatAgentSessionError('无权访问该会话'), true);
|
||||||
assert.equal(
|
assert.equal(
|
||||||
isRecoverableWechatAgentSessionError('Session already has an active request. Cancel it first.'),
|
isRecoverableWechatAgentSessionError('Session already has an active request. Cancel it first.'),
|
||||||
@@ -5119,6 +5125,210 @@ test('wechat mp page delivery requires and verifies a current-run fresh thumbnai
|
|||||||
assert.match(wechatPayloads[0].text.content, /fresh\.html/);
|
assert.match(wechatPayloads[0].text.content, /fresh\.html/);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('wechat mp retries a page in a new session before reporting a missing fresh thumbnail', async () => {
|
||||||
|
const token = 'token';
|
||||||
|
const timestamp = '1710000000';
|
||||||
|
const nonce = 'nonce';
|
||||||
|
const workspaceRoot = fs.mkdtempSync('/tmp/wechat-mp-fresh-thumbnail-retry-');
|
||||||
|
const htmlPath = path.join(workspaceRoot, 'public', 'retry.html');
|
||||||
|
const generated = {
|
||||||
|
ok: true,
|
||||||
|
jobId: 'job-fresh-page-retry',
|
||||||
|
source: { mimeType: 'image/webp', width: 1280, height: 720 },
|
||||||
|
asset: {
|
||||||
|
id: 'asset-fresh-page-retry',
|
||||||
|
htmlSrc: 'images/retry-cover.webp',
|
||||||
|
publicUrl: 'https://example.com/MindSpace/user-1/public/images/retry-cover.webp',
|
||||||
|
workspaceRelativePath: 'public/images/retry-cover.webp',
|
||||||
|
},
|
||||||
|
};
|
||||||
|
const firstHtml = previewReadyPageHtml({
|
||||||
|
title: 'Retry',
|
||||||
|
subtitle: '首次没有新缩略图',
|
||||||
|
cover: 'images/missing-cover.webp',
|
||||||
|
});
|
||||||
|
const retryHtml = previewReadyPageHtml({
|
||||||
|
title: 'Retry',
|
||||||
|
subtitle: '重试生成新缩略图',
|
||||||
|
cover: generated.asset.htmlSrc,
|
||||||
|
});
|
||||||
|
const pageImage = await sharp({
|
||||||
|
create: { width: 64, height: 64, channels: 3, background: '#cc6633' },
|
||||||
|
}).webp().toBuffer();
|
||||||
|
const eventFrame = (event) => `data: ${JSON.stringify(event)}\n\n`;
|
||||||
|
const wechatPayloads = [];
|
||||||
|
let routeCleared = false;
|
||||||
|
let startedSessions = 0;
|
||||||
|
|
||||||
|
const replyEvents = ({ requestId, html, includeImage = false }) => [
|
||||||
|
eventFrame({
|
||||||
|
type: 'Message',
|
||||||
|
request_id: requestId,
|
||||||
|
message: {
|
||||||
|
id: `${requestId}-tools`,
|
||||||
|
role: 'assistant',
|
||||||
|
metadata: { userVisible: true },
|
||||||
|
content: [
|
||||||
|
{
|
||||||
|
id: `${requestId}-page-skill`,
|
||||||
|
type: 'toolRequest',
|
||||||
|
toolCall: { value: { name: 'load_skill', arguments: { name: 'static-page-publish' } } },
|
||||||
|
},
|
||||||
|
...(includeImage
|
||||||
|
? [
|
||||||
|
{
|
||||||
|
id: `${requestId}-image`,
|
||||||
|
type: 'toolRequest',
|
||||||
|
toolCall: {
|
||||||
|
value: { name: 'sandbox-fs__generate_image', arguments: { purpose: 'hero' } },
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
id: `${requestId}-image`,
|
||||||
|
type: 'toolResponse',
|
||||||
|
toolResult: {
|
||||||
|
status: 'success',
|
||||||
|
value: { content: [{ type: 'text', text: JSON.stringify(generated) }] },
|
||||||
|
},
|
||||||
|
},
|
||||||
|
]
|
||||||
|
: []),
|
||||||
|
{
|
||||||
|
id: `${requestId}-write`,
|
||||||
|
type: 'toolRequest',
|
||||||
|
toolCall: {
|
||||||
|
value: {
|
||||||
|
name: 'sandbox-fs__write_file',
|
||||||
|
arguments: { path: 'public/retry.html', content: html },
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
],
|
||||||
|
},
|
||||||
|
}),
|
||||||
|
eventFrame({
|
||||||
|
type: 'Message',
|
||||||
|
request_id: requestId,
|
||||||
|
message: {
|
||||||
|
id: `${requestId}-final`,
|
||||||
|
role: 'assistant',
|
||||||
|
metadata: { userVisible: true },
|
||||||
|
content: [{
|
||||||
|
type: 'text',
|
||||||
|
text: '页面已完成:https://example.com/MindSpace/user-1/public/retry.html',
|
||||||
|
}],
|
||||||
|
},
|
||||||
|
}),
|
||||||
|
eventFrame({
|
||||||
|
type: 'Finish',
|
||||||
|
request_id: requestId,
|
||||||
|
token_state: { inputTokens: 1, outputTokens: 2 },
|
||||||
|
}),
|
||||||
|
].join('');
|
||||||
|
|
||||||
|
const service = createBoundWechatService({
|
||||||
|
token,
|
||||||
|
config: { requireFreshPageThumbnail: true },
|
||||||
|
startAgentSession: async () => {
|
||||||
|
startedSessions += 1;
|
||||||
|
return { id: 'session-2' };
|
||||||
|
},
|
||||||
|
userAuth: {
|
||||||
|
async getWechatAgentRoute() {
|
||||||
|
return routeCleared ? null : { agentSessionId: 'session-1' };
|
||||||
|
},
|
||||||
|
async clearWechatAgentRoute() {
|
||||||
|
routeCleared = true;
|
||||||
|
},
|
||||||
|
async upsertWechatAgentRoute({ agentSessionId }) {
|
||||||
|
assert.equal(agentSessionId, 'session-2');
|
||||||
|
},
|
||||||
|
async resolveWorkingDir() {
|
||||||
|
return workspaceRoot;
|
||||||
|
},
|
||||||
|
async getUserPublishLayout() {
|
||||||
|
return {
|
||||||
|
publishDir: workspaceRoot,
|
||||||
|
displayName: 'John',
|
||||||
|
username: 'john',
|
||||||
|
slug: 'john',
|
||||||
|
constraints: null,
|
||||||
|
};
|
||||||
|
},
|
||||||
|
},
|
||||||
|
sessionApiFetch: async (sessionId, pathname) => {
|
||||||
|
if (pathname === '/agent/harness_remember' || pathname === '/agent/harness_bootstrap') {
|
||||||
|
return new Response('{}', { status: 200, headers: { 'Content-Type': 'application/json' } });
|
||||||
|
}
|
||||||
|
if (pathname === `/sessions/${sessionId}/reply`) {
|
||||||
|
fs.mkdirSync(path.dirname(htmlPath), { recursive: true });
|
||||||
|
if (sessionId === 'session-1') {
|
||||||
|
fs.writeFileSync(htmlPath, firstHtml, 'utf8');
|
||||||
|
} else {
|
||||||
|
fs.writeFileSync(htmlPath, retryHtml, 'utf8');
|
||||||
|
fs.mkdirSync(path.join(workspaceRoot, 'public', 'images'), { recursive: true });
|
||||||
|
fs.writeFileSync(path.join(workspaceRoot, 'public', generated.asset.htmlSrc), pageImage);
|
||||||
|
}
|
||||||
|
return new Response('{}', { status: 200, headers: { 'Content-Type': 'application/json' } });
|
||||||
|
}
|
||||||
|
if (pathname === '/sessions/session-1/events') {
|
||||||
|
return new Response(
|
||||||
|
replyEvents({ requestId: 'req-thumbnail-first', html: firstHtml }),
|
||||||
|
{ status: 200, headers: { 'Content-Type': 'text/event-stream' } },
|
||||||
|
);
|
||||||
|
}
|
||||||
|
if (pathname === '/sessions/session-2/events') {
|
||||||
|
return new Response(
|
||||||
|
replyEvents({ requestId: 'req-thumbnail-retry', html: retryHtml, includeImage: true }),
|
||||||
|
{ status: 200, headers: { 'Content-Type': 'text/event-stream' } },
|
||||||
|
);
|
||||||
|
}
|
||||||
|
throw new Error(`unexpected session path: ${sessionId} ${pathname}`);
|
||||||
|
},
|
||||||
|
wechatFetch: async (url, init = {}) => {
|
||||||
|
if (String(url).includes('/cgi-bin/stable_token')) {
|
||||||
|
return new Response(JSON.stringify({ access_token: 'access-1', expires_in: 7200 }), {
|
||||||
|
status: 200,
|
||||||
|
headers: { 'Content-Type': 'application/json' },
|
||||||
|
});
|
||||||
|
}
|
||||||
|
if (String(url).includes('/cgi-bin/message/custom/send')) {
|
||||||
|
wechatPayloads.push(JSON.parse(init.body));
|
||||||
|
return new Response(JSON.stringify({ errcode: 0, errmsg: 'ok' }), {
|
||||||
|
status: 200,
|
||||||
|
headers: { 'Content-Type': 'application/json' },
|
||||||
|
});
|
||||||
|
}
|
||||||
|
throw new Error(`unexpected wechat url: ${url}`);
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
|
const originalRandomUuid = crypto.randomUUID;
|
||||||
|
crypto.randomUUID = (() => {
|
||||||
|
const ids = ['req-thumbnail-first', 'req-thumbnail-retry'];
|
||||||
|
return () => ids.shift() ?? 'req-thumbnail-retry';
|
||||||
|
})();
|
||||||
|
try {
|
||||||
|
const result = await service.handleInboundMessage(
|
||||||
|
inboundXml({ content: '生成一个活动页面,只要文字,不要正文图片' }),
|
||||||
|
{ timestamp, nonce, signature: signatureFor(token, timestamp, nonce) },
|
||||||
|
);
|
||||||
|
assert.equal(result.status, 200);
|
||||||
|
await result.task;
|
||||||
|
assert.equal(fs.existsSync(path.join(workspaceRoot, 'public', 'retry.thumbnail.svg')), true);
|
||||||
|
} finally {
|
||||||
|
crypto.randomUUID = originalRandomUuid;
|
||||||
|
fs.rmSync(workspaceRoot, { recursive: true, force: true });
|
||||||
|
}
|
||||||
|
|
||||||
|
assert.equal(routeCleared, true);
|
||||||
|
assert.equal(startedSessions, 1);
|
||||||
|
assert.equal(wechatPayloads.length, 1);
|
||||||
|
assert.equal(wechatPayloads[0].msgtype, 'text');
|
||||||
|
assert.match(wechatPayloads[0].text.content, /retry\.html/);
|
||||||
|
assert.doesNotMatch(wechatPayloads[0].text.content, /没有完成服务号要求的本轮新缩略图/);
|
||||||
|
});
|
||||||
|
|
||||||
test('wechat mp standalone image intent sends a native image message and text', async () => {
|
test('wechat mp standalone image intent sends a native image message and text', async () => {
|
||||||
const token = 'token';
|
const token = 'token';
|
||||||
const timestamp = '1710000000';
|
const timestamp = '1710000000';
|
||||||
|
|||||||
Reference in New Issue
Block a user