feat: add wechat publish recovery and conversation memory
This commit is contained in:
@@ -0,0 +1,152 @@
|
||||
import assert from 'node:assert/strict';
|
||||
import test from 'node:test';
|
||||
import { createConversationMemoryService, extractConversationMessageText } from './conversation-memory.mjs';
|
||||
|
||||
function createPool() {
|
||||
const state = {
|
||||
messages: [],
|
||||
memories: [],
|
||||
analyzed: new Set(),
|
||||
};
|
||||
return {
|
||||
state,
|
||||
async query(sql, params = []) {
|
||||
if (sql.includes('INSERT INTO h5_conversation_messages')) {
|
||||
for (const row of params[0]) {
|
||||
const [
|
||||
id,
|
||||
userId,
|
||||
sessionId,
|
||||
messageKey,
|
||||
sequenceNo,
|
||||
role,
|
||||
text,
|
||||
rawJson,
|
||||
createdAt,
|
||||
updatedAt,
|
||||
] = row;
|
||||
const existing = state.messages.find(
|
||||
(item) => item.agent_session_id === sessionId && item.message_key === messageKey,
|
||||
);
|
||||
const next = {
|
||||
id,
|
||||
user_id: userId,
|
||||
agent_session_id: sessionId,
|
||||
message_key: messageKey,
|
||||
sequence_no: sequenceNo,
|
||||
role,
|
||||
text,
|
||||
raw_json: rawJson,
|
||||
created_at: createdAt,
|
||||
updated_at: updatedAt,
|
||||
analyzed_at: existing?.analyzed_at ?? null,
|
||||
};
|
||||
if (existing) Object.assign(existing, next);
|
||||
else state.messages.push(next);
|
||||
}
|
||||
return [{ affectedRows: params[0].length }];
|
||||
}
|
||||
if (sql.includes('FROM h5_conversation_messages')) {
|
||||
const [userId, limit] = params;
|
||||
return [
|
||||
state.messages
|
||||
.filter((item) => item.user_id === userId && item.role === 'user' && item.analyzed_at == null)
|
||||
.slice(0, limit),
|
||||
];
|
||||
}
|
||||
if (sql.includes('INSERT INTO h5_user_memory_items')) {
|
||||
for (const row of params[0]) {
|
||||
const [
|
||||
id,
|
||||
userId,
|
||||
label,
|
||||
memoryHash,
|
||||
memoryText,
|
||||
evidenceMessageId,
|
||||
sourceSessionId,
|
||||
confidence,
|
||||
rawJson,
|
||||
createdAt,
|
||||
updatedAt,
|
||||
] = row;
|
||||
state.memories.push({
|
||||
id,
|
||||
user_id: userId,
|
||||
label,
|
||||
memory_hash: memoryHash,
|
||||
memory_text: memoryText,
|
||||
evidence_message_id: evidenceMessageId,
|
||||
source_session_id: sourceSessionId,
|
||||
confidence,
|
||||
raw_json: rawJson,
|
||||
created_at: createdAt,
|
||||
updated_at: updatedAt,
|
||||
});
|
||||
}
|
||||
return [{ affectedRows: params[0].length }];
|
||||
}
|
||||
if (sql.includes('UPDATE h5_conversation_messages SET analyzed_at')) {
|
||||
const [analyzedAt, ids] = params;
|
||||
for (const item of state.messages) {
|
||||
if (ids.includes(item.id)) item.analyzed_at = analyzedAt;
|
||||
}
|
||||
return [{ affectedRows: ids.length }];
|
||||
}
|
||||
if (sql.includes('FROM h5_user_memory_items')) {
|
||||
const [userId, limit] = params;
|
||||
return [
|
||||
state.memories
|
||||
.filter((item) => item.user_id === userId)
|
||||
.slice(0, limit)
|
||||
.map((item) => ({ ...item, status: 'active' })),
|
||||
];
|
||||
}
|
||||
if (sql.includes('FROM h5_llm_provider_keys')) return [[]];
|
||||
throw new Error(`Unexpected SQL: ${sql}`);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
test('extractConversationMessageText reads text content and display text', () => {
|
||||
assert.equal(
|
||||
extractConversationMessageText({ content: [{ type: 'text', text: ' hello ' }] }),
|
||||
'hello',
|
||||
);
|
||||
assert.equal(
|
||||
extractConversationMessageText({ content: [], metadata: { displayText: 'fallback' } }),
|
||||
'fallback',
|
||||
);
|
||||
});
|
||||
|
||||
test('saveAndAnalyze stores messages and fallback memories', async () => {
|
||||
const previous = process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||
process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = '0';
|
||||
const pool = createPool();
|
||||
const service = createConversationMemoryService(pool, { now: () => 1000 });
|
||||
|
||||
const result = await service.saveAndAnalyze('session-1', 'user-1', [
|
||||
{
|
||||
id: 'm1',
|
||||
role: 'user',
|
||||
content: [{ type: 'text', text: '我喜欢简洁直接的回答,也关注 AI 产品设计。' }],
|
||||
metadata: { userVisible: true },
|
||||
},
|
||||
{
|
||||
id: 'm2',
|
||||
role: 'assistant',
|
||||
content: [{ type: 'text', text: '好的。' }],
|
||||
metadata: { userVisible: true },
|
||||
},
|
||||
]);
|
||||
|
||||
assert.equal(result.saved, 2);
|
||||
assert.equal(result.analyzed, 1);
|
||||
assert.equal(result.memories, 1);
|
||||
assert.equal(pool.state.messages.length, 2);
|
||||
assert.equal(pool.state.memories[0].label, 'preference');
|
||||
assert.match(pool.state.memories[0].memory_text, /简洁直接/);
|
||||
assert.equal(pool.state.messages.find((item) => item.message_key === 'm1')?.analyzed_at, 1000);
|
||||
|
||||
if (previous == null) delete process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED;
|
||||
else process.env.USER_CONVERSATION_MEMORY_LLM_ENABLED = previous;
|
||||
});
|
||||
Reference in New Issue
Block a user