fix(chat): translate session replay cursors
Memind CI / Test, build, and release guards (pull_request) Successful in 2m55s
Memind CI / Test, build, and release guards (pull_request) Successful in 2m55s
This commit is contained in:
@@ -533,6 +533,178 @@ test('proxySessionEvents attaches session taxonomy when flag enabled', async ()
|
||||
}
|
||||
});
|
||||
|
||||
test('proxySessionEvents omits an unmapped Portal cursor so Goose can replay Finish', async () => {
|
||||
const previousReplay = process.env.MEMIND_SESSION_STREAM_REPLAY;
|
||||
process.env.MEMIND_SESSION_STREAM_REPLAY = '1';
|
||||
let upstream;
|
||||
let receivedLastEventId = 'not-requested';
|
||||
|
||||
try {
|
||||
upstream = createServer((req, res) => {
|
||||
if (req.method === 'GET' && req.url === '/sessions/session-1/events') {
|
||||
receivedLastEventId = req.headers['last-event-id'] ?? null;
|
||||
res.writeHead(200, { 'Content-Type': 'text/event-stream' });
|
||||
res.end(
|
||||
'id: upstream-finish\n' +
|
||||
'data: {"type":"Finish","reason":"stop","request_id":"req-1"}\n\n',
|
||||
);
|
||||
return;
|
||||
}
|
||||
res.writeHead(404, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ message: `unexpected ${req.method} ${req.url}` }));
|
||||
});
|
||||
const upstreamPort = await listen(upstream);
|
||||
const persisted = [];
|
||||
const proxy = createTkmindProxy({
|
||||
apiTarget: `http://127.0.0.1:${upstreamPort}`,
|
||||
apiSecret: 'test-secret',
|
||||
userAuth: createMemoryTestUserAuth(process.cwd()),
|
||||
sessionStreamStore: {
|
||||
async listEventsForUser(_userId, _sessionId, { afterEventId } = {}) {
|
||||
assert.equal(afterEventId, 'portal-active-request-uuid');
|
||||
return {
|
||||
cursorMiss: false,
|
||||
cursorEvent: {
|
||||
id: 'portal-active-request-uuid',
|
||||
payload: { type: 'ActiveRequests', request_ids: ['req-1'] },
|
||||
upstreamEventId: null,
|
||||
createdAt: 1,
|
||||
},
|
||||
events: [],
|
||||
};
|
||||
},
|
||||
async appendEvent(event) {
|
||||
persisted.push(event);
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
const req = new EventEmitter();
|
||||
req.currentUser = { id: 'user-1', username: 'john' };
|
||||
req.get = (name) =>
|
||||
name.toLowerCase() === 'last-event-id' ? 'portal-active-request-uuid' : '';
|
||||
req.once = req.once.bind(req);
|
||||
req.off = req.off.bind(req);
|
||||
|
||||
const chunks = [];
|
||||
const res = new EventEmitter();
|
||||
res.headersSent = false;
|
||||
res.writableEnded = false;
|
||||
res.statusCode = 200;
|
||||
res.status = (code) => {
|
||||
res.statusCode = code;
|
||||
return res;
|
||||
};
|
||||
res.setHeader = () => {};
|
||||
res.flushHeaders = () => {
|
||||
res.headersSent = true;
|
||||
};
|
||||
res.write = (chunk) => {
|
||||
res.headersSent = true;
|
||||
chunks.push(String(chunk));
|
||||
return true;
|
||||
};
|
||||
res.end = () => {
|
||||
res.writableEnded = true;
|
||||
};
|
||||
res.json = () => {
|
||||
throw new Error('json must not be called after SSE started');
|
||||
};
|
||||
|
||||
await proxy.proxySessionEvents(req, res, 'session-1');
|
||||
|
||||
assert.equal(receivedLastEventId, null);
|
||||
assert.match(chunks.join(''), /"type":"Finish"/);
|
||||
await new Promise((resolve) => setImmediate(resolve));
|
||||
assert.equal(persisted.at(-1)?.upstreamEventId, 'upstream-finish');
|
||||
} finally {
|
||||
if (previousReplay == null) delete process.env.MEMIND_SESSION_STREAM_REPLAY;
|
||||
else process.env.MEMIND_SESSION_STREAM_REPLAY = previousReplay;
|
||||
await closeServer(upstream);
|
||||
}
|
||||
});
|
||||
|
||||
test('proxySessionEvents translates a Portal replay cursor to its Goose cursor', async () => {
|
||||
const previousReplay = process.env.MEMIND_SESSION_STREAM_REPLAY;
|
||||
process.env.MEMIND_SESSION_STREAM_REPLAY = '1';
|
||||
let upstream;
|
||||
let receivedLastEventId = null;
|
||||
|
||||
try {
|
||||
upstream = createServer((req, res) => {
|
||||
if (req.method === 'GET' && req.url === '/sessions/session-1/events') {
|
||||
receivedLastEventId = req.headers['last-event-id'] ?? null;
|
||||
res.writeHead(200, { 'Content-Type': 'text/event-stream' });
|
||||
res.end('id: upstream-finish\ndata: {"type":"Finish","reason":"stop"}\n\n');
|
||||
return;
|
||||
}
|
||||
res.writeHead(404, { 'Content-Type': 'application/json' });
|
||||
res.end(JSON.stringify({ message: `unexpected ${req.method} ${req.url}` }));
|
||||
});
|
||||
const upstreamPort = await listen(upstream);
|
||||
const proxy = createTkmindProxy({
|
||||
apiTarget: `http://127.0.0.1:${upstreamPort}`,
|
||||
apiSecret: 'test-secret',
|
||||
userAuth: createMemoryTestUserAuth(process.cwd()),
|
||||
sessionStreamStore: {
|
||||
async listEventsForUser() {
|
||||
return {
|
||||
cursorMiss: false,
|
||||
cursorEvent: {
|
||||
id: 'portal-message-uuid',
|
||||
payload: { type: 'Message' },
|
||||
upstreamEventId: 'upstream-message-17',
|
||||
createdAt: 1,
|
||||
},
|
||||
events: [],
|
||||
};
|
||||
},
|
||||
async appendEvent() {},
|
||||
},
|
||||
});
|
||||
|
||||
const req = new EventEmitter();
|
||||
req.currentUser = { id: 'user-1', username: 'john' };
|
||||
req.get = (name) => name.toLowerCase() === 'last-event-id' ? 'portal-message-uuid' : '';
|
||||
req.once = req.once.bind(req);
|
||||
req.off = req.off.bind(req);
|
||||
|
||||
const chunks = [];
|
||||
const res = new EventEmitter();
|
||||
res.headersSent = false;
|
||||
res.writableEnded = false;
|
||||
res.statusCode = 200;
|
||||
res.status = (code) => {
|
||||
res.statusCode = code;
|
||||
return res;
|
||||
};
|
||||
res.setHeader = () => {};
|
||||
res.flushHeaders = () => {
|
||||
res.headersSent = true;
|
||||
};
|
||||
res.write = (chunk) => {
|
||||
res.headersSent = true;
|
||||
chunks.push(String(chunk));
|
||||
return true;
|
||||
};
|
||||
res.end = () => {
|
||||
res.writableEnded = true;
|
||||
};
|
||||
res.json = () => {
|
||||
throw new Error('json must not be called after SSE started');
|
||||
};
|
||||
|
||||
await proxy.proxySessionEvents(req, res, 'session-1');
|
||||
|
||||
assert.equal(receivedLastEventId, 'upstream-message-17');
|
||||
assert.match(chunks.join(''), /"type":"Finish"/);
|
||||
} finally {
|
||||
if (previousReplay == null) delete process.env.MEMIND_SESSION_STREAM_REPLAY;
|
||||
else process.env.MEMIND_SESSION_STREAM_REPLAY = previousReplay;
|
||||
await closeServer(upstream);
|
||||
}
|
||||
});
|
||||
|
||||
test('startSessionForUser resolves memories through Memory V2 facade', async () => {
|
||||
let resolveInput = null;
|
||||
await withFakeGoosedSession(async ({ apiTarget, workingDir, harnessEntries }) => {
|
||||
|
||||
Reference in New Issue
Block a user