| 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 |