返回 CodeWhale
pet-runtime-continuous.test.mjs
根目录 / pet / tests / pet-runtime-continuous.test.mjs
1 import test from 'node:test';
2 import assert from 'node:assert/strict';
3 import { createServer } from 'node:http';
4 import { once } from 'node:events';
5 import { setTimeout as delay } from 'node:timers/promises';
6 import { followRuntime } from '../scripts/lib/pet-runtime.mjs';
7 import { compilePetTelemetry } from '../dist/core/pet-telemetry.js';
8
9 // Real HTTP/SSE transport, virtual event timestamps: no provider calls or
10 // 29-hour wall-clock sleep. This exceeded the old raw-journal lifetime limit.
11 test('Runtime consumes more than 250000 records across a day while retaining current state and unfinished requests', { timeout: 60_000 }, async t => {
12 const total = 260_001, began = Date.now() - total * 400, reports = [];
13 let response, transport, serverError;
14 const server = createServer(async (_req, res) => {
15 response = res;
16 const closed = new AbortController(); res.once('close', () => closed.abort());
17 res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' });
18 res.write(`data: ${JSON.stringify({ event: 'stream.progress', state: 'replaying', thread_id: 'long-fixture', seq: 0 })}\n\n`);
19 try {
20 for (let first = 1; first <= total; first += 128) {
21 let block = '';
22 for (let seq = first; seq < first + 128 && seq <= total; seq++) {
23 block += `data: ${JSON.stringify({ seq, previous_seq: seq - 1, event: seq === 1 ? 'user_input.required' : 'thread.updated',
24 thread_id: 'long-fixture', timestamp: new Date(began + seq * 400).toISOString(),
25 payload: seq === 1 ? { id: 'still-waiting' } : { description: 'fixture-private-journal'.repeat(12) } })}\n\n`;
26 }
27 if (!res.write(block)) await once(res, 'drain', { signal: closed.signal });
28 }
29 res.write(`data: ${JSON.stringify({ event: 'stream.progress', state: 'live', thread_id: 'long-fixture', seq: total })}\n\n`);
30 } catch (error) { if (!closed.signal.aborted) serverError = error; }
31 // Keep the final cursor healthy, including the still-open human request.
32 });
33 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
34 t.after(async () => {
35 await transport?.close(); response?.destroy(); server.closeAllConnections();
36 await new Promise(resolve => server.close(resolve));
37 });
38 transport = await followRuntime({ baseUrl: `http://127.0.0.1:${server.address().port}`, threadId: 'long-fixture', report: text => reports.push(text) });
39 const deadline = Date.now() + 50_000;
40 while ((transport.cursor < total || !transport.connected) && Date.now() < deadline && !reports.some(text => text.includes('stopped'))) await delay(20);
41 assert.equal(serverError, undefined);
42 assert.equal(transport.cursor, total, reports.join('\n'));
43 assert.equal(transport.connected, true);
44 const now = Date.now(), snapshot = transport.snapshot(now);
45 assert.ok(snapshot); assert.doesNotMatch(JSON.stringify(snapshot), /fixture-private-journal/);
46 assert.ok(transport.retainedEvents < 512); assert.ok(transport.retainedBytes < 256 * 1024);
47 const live = compilePetTelemetry(snapshot.events, snapshot.duration, 0, now - 800 - Date.parse(snapshot.originTime));
48 assert.equal(live[0].waiting, true); assert.equal(live[0].activeMs[11], 400);
49 assert.deepEqual(reports, []);
50 });
51
51 lines Plain Text