Files
memind/session-stream-store.test.mjs
T
john 14a00774d9 feat(h5): web 联网能力、实时查询路由与 session Finish 对齐
- 新增 web 能力并挂载 platform/web(web_search/fetch_url)
- 实时查询强制 web skill,router fallback 与 await session Finish
- Session Broker 覆盖率/指标、stream replay 与相关单测/E2E 脚本

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-06 16:06:26 +08:00

89 lines
2.9 KiB
JavaScript

import assert from 'node:assert/strict';
import test from 'node:test';
import { createSessionStreamStore } from './session-stream-store.mjs';
function createMockPool() {
let clock = 1;
const events = [];
return {
events,
async query(sql, params = []) {
if (sql.includes('INSERT INTO h5_session_stream_events')) {
const createdAt = clock++;
events.push({
id: params[0],
user_id: params[1],
agent_session_id: params[2],
event_type: params[3],
payload_json: params[4],
upstream_event_id: params[5],
created_at: createdAt,
});
return [{ insertId: 1 }];
}
if (sql.includes('WHERE id = ? AND agent_session_id = ?')) {
const row = events.find((item) => item.id === params[0]);
return [[row ? { created_at: row.created_at } : undefined].filter(Boolean)];
}
if (sql.includes('created_at > ?')) {
const after = Number(params[2]);
const rows = events
.filter(
(item) =>
item.agent_session_id === params[0] &&
item.user_id === params[1] &&
Number(item.created_at) > after,
)
.sort((a, b) => Number(a.created_at) - Number(b.created_at));
return [rows];
}
if (sql.includes('ORDER BY created_at ASC, id ASC')) {
const rows = events
.filter(
(item) => item.agent_session_id === params[0] && item.user_id === params[1],
)
.sort((a, b) => Number(a.created_at) - Number(b.created_at))
.slice(0, params[2]);
return [rows];
}
if (sql.includes('ORDER BY created_at DESC, id DESC')) {
const row = events
.filter(
(item) => item.agent_session_id === params[0] && item.user_id === params[1],
)
.sort((a, b) => Number(b.created_at) - Number(a.created_at))[0];
return [[row].filter(Boolean)];
}
throw new Error(`unexpected sql: ${sql}`);
},
};
}
test('session stream store append and replay with cursor', async () => {
const pool = createMockPool();
const store = createSessionStreamStore({ pool });
const first = await store.appendEvent({
userId: 'user-1',
sessionId: 'sess-1',
payload: { type: 'Message', message: { role: 'assistant', content: [] } },
id: 'evt-1',
});
await store.appendEvent({
userId: 'user-1',
sessionId: 'sess-1',
payload: { type: 'Finish' },
id: 'evt-2',
});
const head = await store.listEventsForUser('user-1', 'sess-1');
assert.equal(head.events.length, 2);
const tail = await store.listEventsForUser('user-1', 'sess-1', { afterEventId: 'evt-1' });
assert.equal(tail.events.length, 1);
assert.equal(tail.events[0].payload.type, 'Finish');
assert.equal(tail.cursorMiss, false);
const miss = await store.listEventsForUser('user-1', 'sess-1', { afterEventId: 'missing' });
assert.equal(miss.cursorMiss, true);
});