| 1 | import test from 'node:test'; |
| 2 | import assert from 'node:assert/strict'; |
| 3 | import { createServer } from 'node:http'; |
| 4 | import { spawn } from 'node:child_process'; |
| 5 | import { mkdtemp, readFile } from 'node:fs/promises'; |
| 6 | import { tmpdir } from 'node:os'; |
| 7 | import { join } from 'node:path'; |
| 8 | import { once } from 'node:events'; |
| 9 | import { setTimeout as delay } from 'node:timers/promises'; |
| 10 | import { followRuntime } from '../scripts/lib/pet-runtime.mjs'; |
| 11 | import { decodePetJSONL } from '../dist/core/pet-telemetry.js'; |
| 12 | import { spawnRecorder } from './helpers/recorder-process.mjs'; |
| 13 | |
| 14 | test('Runtime pet input refuses remote hosts, credentials, paths and missing thread selection before connecting', async () => { |
| 15 | for (const baseUrl of ['https://127.0.0.1:1', 'http://example.com', 'http://localhost:1', 'http://user:secret@127.0.0.1:1', 'http://127.0.0.1:1/private', 'http://127.0.0.1:1/?token=secret']) |
| 16 | await assert.rejects(followRuntime({ baseUrl, threadId: 't' }), /loopback IP origin/); |
| 17 | await assert.rejects(followRuntime({ baseUrl: 'http://127.0.0.1:1', threadId: '' }), /thread/); |
| 18 | }); |
| 19 | |
| 20 | test('Runtime shutdown closes an idle SSE body after garbage collection', { timeout: 10_000 }, async t => { |
| 21 | let response, closed = false; |
| 22 | const server = createServer((req, res) => { |
| 23 | response = res; |
| 24 | res.once('close', () => { closed = true; }); |
| 25 | res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); |
| 26 | res.write(`data: ${JSON.stringify({ seq: 1, previous_seq: 0, event: 'thread.updated', |
| 27 | thread_id: 'fixture', timestamp: new Date().toISOString(), payload: {} })}\n\n`); |
| 28 | res.write(`data: ${JSON.stringify({ event: 'stream.progress', state: 'live', thread_id: 'fixture', seq: 1 })}\n\n`); |
| 29 | // Stay open without new chunks: cancellation must wake the idle reader. |
| 30 | }); |
| 31 | await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); |
| 32 | const script = ` |
| 33 | import assert from 'node:assert/strict'; |
| 34 | import { setTimeout as delay } from 'node:timers/promises'; |
| 35 | import { followRuntime } from './scripts/lib/pet-runtime.mjs'; |
| 36 | const input = await followRuntime({ baseUrl: 'http://127.0.0.1:${server.address().port}', threadId: 'fixture' }); |
| 37 | for (let i = 0; !input.connected && i < 200; i++) await delay(10); |
| 38 | assert.equal(input.cursor, 1); |
| 39 | globalThis.gc(); await delay(20); globalThis.gc(); |
| 40 | const deadline = setTimeout(() => { console.error('Idle Runtime reader did not stop'); process.exit(2); }, 2_000); |
| 41 | await input.close(); clearTimeout(deadline); |
| 42 | assert.equal(input.connected, false); |
| 43 | `; |
| 44 | const child = spawn(process.execPath, ['--expose-gc', '--input-type=module', '-e', script], |
| 45 | { cwd: new URL('../', import.meta.url), stdio: ['ignore', 'pipe', 'pipe'] }); |
| 46 | const exited = once(child, 'exit'); let log = ''; |
| 47 | child.stdout.on('data', b => log += b); child.stderr.on('data', b => log += b); |
| 48 | t.after(async () => { |
| 49 | if (child.exitCode === null) child.kill('SIGKILL'); |
| 50 | response?.destroy(); server.closeAllConnections(); |
| 51 | await new Promise(resolve => server.close(resolve)); |
| 52 | }); |
| 53 | const [code] = await exited; |
| 54 | assert.equal(code, 0, log); assert.ok(closed, 'The server must see the reader disconnect'); |
| 55 | }); |
| 56 | |
| 57 | test('the CLI follows real Runtime SSE envelopes through disconnect and cursor recovery, recording no prompt content', { timeout: 20_000 }, async t => { |
| 58 | let sequence = 0, connections = 0, stream, pulse; |
| 59 | const requests = [], responses = new Set(); |
| 60 | const emit = (res, tool) => { |
| 61 | const now = Date.now(), previous = sequence; sequence += 7; |
| 62 | const event = { seq: sequence, previous_seq: previous, event: 'item.completed', thread_id: 'fixture-thread', item_id: `i${sequence}`, |
| 63 | timestamp: new Date(now).toISOString(), payload: { item: { id: `i${sequence}`, kind: 'tool_call', status: 'completed', |
| 64 | started_at: new Date(now - 180).toISOString(), ended_at: new Date(now).toISOString(), summary: `${tool}: fixture-private-text` }, tool } }; |
| 65 | res.write(`data: ${JSON.stringify(event)}\n\n`); |
| 66 | }; |
| 67 | const server = createServer((req, res) => { |
| 68 | requests.push({ method: req.method, url: req.url, authorization: req.headers.authorization }); |
| 69 | res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); res.flushHeaders(); responses.add(res); |
| 70 | res.on('close', () => responses.delete(res)); |
| 71 | const connection = ++connections; |
| 72 | res.write(`data: ${JSON.stringify({ event: 'stream.progress', state: 'live', thread_id: 'fixture-thread', seq: Number(new URL(req.url, 'http://local').searchParams.get('since_seq')) })}\n\n`); |
| 73 | if (connection === 1) { stream = res; emit(res, 'bash'); pulse = setInterval(() => emit(res, 'bash'), 120); } |
| 74 | else if (connection === 2) { |
| 75 | // Reject a hole; the next reconnect must request the same last cursor. |
| 76 | res.end(`data: ${JSON.stringify({ seq: sequence + 20, previous_seq: sequence + 1, event: 'thread.updated', thread_id: 'fixture-thread', timestamp: new Date().toISOString(), payload: {} })}\n\n`); |
| 77 | } else { |
| 78 | const later = setTimeout(() => { emit(res, 'browser'); pulse = setInterval(() => emit(res, 'browser'), 120); }, 700); |
| 79 | res.once('close', () => clearTimeout(later)); |
| 80 | } |
| 81 | }); |
| 82 | await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); |
| 83 | const dir = await mkdtemp(join(tmpdir(), 'pet-runtime-')), output = join(dir, 'pet.jsonl'); |
| 84 | const child = spawnRecorder([`--runtime=http://127.0.0.1:${server.address().port}`, '--thread=fixture-thread', `--output=${output}`], |
| 85 | { ...process.env, CODEWHALE_RUNTIME_TOKEN: 'fixture-token' }); |
| 86 | const exited = once(child, 'exit'); let log = ''; |
| 87 | child.stdout.on('data', b => log += b); child.stderr.on('data', b => log += b); |
| 88 | t.after(async () => { clearInterval(pulse); if (child.exitCode === null) child.kill('SIGTERM'); for (const res of responses) res.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); }); |
| 89 | for (let i = 0; !stream && i < 60; i++) await delay(25); |
| 90 | assert.ok(stream, log); |
| 91 | const waitForRecordedChannel = async channel => { |
| 92 | const until = Date.now() + 6_000; |
| 93 | while (Date.now() < until) { |
| 94 | try { |
| 95 | const tape = decodePetJSONL(await readFile(output, 'utf8')); |
| 96 | if (tape.some(b => b.channel === channel && b.observed === 1)) return; |
| 97 | } catch { /* The first file or an in-flight final line is not ready. */ } |
| 98 | await delay(40); |
| 99 | } |
| 100 | assert.fail(`The recorder did not persist observed ${channel} work. ${log}`); |
| 101 | }; |
| 102 | // Assert actual recorder output before moving the fixture to its next phase. |
| 103 | // A fixed sleep can expire before reconnect + a complete bin on a busy runner. |
| 104 | await waitForRecordedChannel('code'); clearInterval(pulse); stream.destroy(); |
| 105 | await waitForRecordedChannel('browser'); child.stopRecorder(); const [code] = await exited; assert.equal(code, 0, log); clearInterval(pulse); |
| 106 | const text = await readFile(output, 'utf8'), tape = decodePetJSONL(text); |
| 107 | assert.ok(tape.some(b => b.channel === 'code' && b.observed === 1)); |
| 108 | assert.ok(tape.some(b => b.sequence > 1 && b.observed === 0)); |
| 109 | assert.ok(tape.some(b => b.channel === 'browser' && b.observed === 1)); |
| 110 | assert.equal(requests.length, 3); assert.equal(new URL(requests[1].url, 'http://127.0.0.1').search, new URL(requests[2].url, 'http://127.0.0.1').search); |
| 111 | assert.ok(requests.every(r => r.method === 'GET' && r.url.startsWith('/v1/threads/fixture-thread/events?') && r.authorization === 'Bearer fixture-token')); |
| 112 | assert.doesNotMatch(text + log, /fixture-private-text|fixture-token/); |
| 113 | }); |
| 114 | |
| 115 | test('live Runtime recording retains human waits and a brief late failure exactly once between timer ticks', { timeout: 15_000 }, async t => { |
| 116 | let sequence = 0, response, answer; |
| 117 | const timers = [], requests = []; |
| 118 | const server = createServer((req, res) => { |
| 119 | requests.push({ method: req.method, url: req.url }); response = res; |
| 120 | res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); res.flushHeaders(); |
| 121 | res.write(`data: ${JSON.stringify({ event: 'stream.progress', state: 'live', thread_id: 'fixture-thread', seq: 0 })}\n\n`); |
| 122 | const oldStart = new Date(Date.now() - 20_000).toISOString(), oldEnd = new Date(Date.now() - 16_000).toISOString(); |
| 123 | const emit = (event, payload) => { |
| 124 | const previous_seq = sequence; sequence++; |
| 125 | res.write(`data: ${JSON.stringify({ seq: sequence, previous_seq, event, thread_id: 'fixture-thread', turn_id: 'turn-a', timestamp: new Date().toISOString(), payload })}\n\n`); |
| 126 | }; |
| 127 | emit('item.started', { item: { id: 'old-work', kind: 'tool_call', status: 'running', started_at: oldStart }, tool: 'bash' }); |
| 128 | timers.push(setTimeout(() => emit('user_input.required', { id: 'question', request: { questions: ['fixture-private-question'] } }), 110)); |
| 129 | timers.push(setTimeout(() => emit('item.completed', { item: { id: 'old-work', kind: 'tool_call', status: 'failed', started_at: oldStart, ended_at: oldEnd, detail: 'fixture-private-result' }, tool: 'bash' }), 650)); |
| 130 | answer = () => emit('user_input.answered', { input_id: 'question', answers: ['fixture-private-answer'] }); |
| 131 | }); |
| 132 | await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); |
| 133 | const dir = await mkdtemp(join(tmpdir(), 'pet-runtime-lifecycle-')), output = join(dir, 'pet.jsonl'); |
| 134 | const child = spawnRecorder([`--runtime=http://127.0.0.1:${server.address().port}`, '--thread=fixture-thread', `--output=${output}`]); |
| 135 | const exited = once(child, 'exit'); let log = ''; |
| 136 | child.stdout.on('data', b => log += b); child.stderr.on('data', b => log += b); |
| 137 | t.after(async () => { timers.forEach(clearTimeout); if (child.exitCode === null) child.kill('SIGTERM'); response?.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); }); |
| 138 | for (let i = 0; !response && i < 100; i++) await delay(20); |
| 139 | assert.ok(response, log); |
| 140 | const waitForTape = async (condition, message) => { |
| 141 | for (let i = 0; i < 125; i++) { |
| 142 | let text = ''; |
| 143 | try { text = await readFile(output, 'utf8'); } catch (error) { if (error.code !== 'ENOENT') throw error; } |
| 144 | const rows = decodePetJSONL(text.slice(0, text.lastIndexOf('\n') + 1)); |
| 145 | if (condition(rows)) return; |
| 146 | await delay(40); |
| 147 | } |
| 148 | assert.fail(message + '\n' + log); |
| 149 | }; |
| 150 | // Drive the answer after actual recorded coverage. Wall-clock sleeps alone |
| 151 | // can stop the child before it seals the final unknown bins on a busy runner. |
| 152 | await waitForTape(rows => rows.filter(b => b.waiting).length >= 3 && rows.some(b => b.errors), 'Waiting/error receipts were not recorded'); |
| 153 | answer(); |
| 154 | await waitForTape(rows => rows.length >= 2 && rows.slice(-2).every(b => !b.waiting && !b.observed), 'Answered input did not expire to unknown'); |
| 155 | child.stopRecorder(); const [code] = await exited; assert.equal(code, 0, log); |
| 156 | const text = await readFile(output, 'utf8'), tape = decodePetJSONL(text); |
| 157 | assert.equal(tape.reduce((sum, b) => sum + b.errors, 0), 1); |
| 158 | assert.ok(tape.some(b => b.channel === 'error' && b.observed === 1)); |
| 159 | assert.ok(tape.filter(b => b.waiting).length >= 3); |
| 160 | assert.ok(tape.slice(-2).every(b => !b.waiting && !b.observed)); |
| 161 | assert.doesNotMatch(text + log, /fixture-private-question|fixture-private-result|fixture-private-answer/); |
| 162 | assert.ok(requests.every(r => r.method === 'GET' && r.url.startsWith('/v1/threads/fixture-thread/events?'))); |
| 163 | }); |
| 164 | |
| 165 | |
| 166 | test('replayed requests remain unknown until catch-up, including reentry to replay after a live connection', { timeout: 10_000 }, async t => { |
| 167 | let response, input, sequence = 0; |
| 168 | const requests = [], reports = []; |
| 169 | const server = createServer((req, res) => { |
| 170 | requests.push(req.url); response = res; |
| 171 | res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); res.flushHeaders(); |
| 172 | }); |
| 173 | await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); |
| 174 | t.after(async () => { await input?.close(); response?.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); }); |
| 175 | input = await followRuntime({ baseUrl: `http://127.0.0.1:${server.address().port}`, threadId: 'fixture', report: text => reports.push(text) }); |
| 176 | const wait = async predicate => { |
| 177 | for (let n = 0; n < 250 && !predicate(); n++) await delay(10); |
| 178 | assert.ok(predicate(), reports.join('\n')); |
| 179 | }; |
| 180 | await wait(() => response); |
| 181 | const write = packet => response.write('data: ' + JSON.stringify(packet) + '\n\n'); |
| 182 | const progress = state => write({ event: 'stream.progress', thread_id: 'fixture', seq: sequence, state }); |
| 183 | const event = (event, payload, age = 0) => write({ seq: ++sequence, previous_seq: sequence - 1, |
| 184 | event, thread_id: 'fixture', turn_id: 'turn-a', timestamp: new Date(Date.now() - age).toISOString(), payload }); |
| 185 | progress('replaying'); event('user_input.required', { id: 'settled-old-request' }, 60_000); |
| 186 | await wait(() => input.cursor === 1); |
| 187 | assert.equal(input.connected, false); assert.equal(input.snapshot(), undefined); |
| 188 | // Simulate a slow backlog while the 400 ms recorder clock could run. |
| 189 | await delay(450); assert.equal(input.snapshot(), undefined); |
| 190 | event('user_input.answered', { id: 'settled-old-request' }, 50_000); progress('live'); |
| 191 | await wait(() => input.connected); |
| 192 | assert.equal(input.snapshot(), undefined, 'A historical answer must arrive before old pending input can become current'); |
| 193 | event('user_input.required', { id: 'fresh-request' }); await wait(() => input.cursor === 3); |
| 194 | const { compilePetTelemetry } = await import('../dist/core/pet-telemetry.js'); |
| 195 | let snapshot = input.snapshot(Date.now() + 400); |
| 196 | assert.ok(compilePetTelemetry(snapshot.events, snapshot.duration).some(b => b.waiting && b.observed === 1)); |
| 197 | progress('replaying'); await wait(() => !input.connected); assert.equal(input.snapshot(), undefined); |
| 198 | event('user_input.answered', { id: 'fresh-request' }); await wait(() => input.cursor === 4); |
| 199 | assert.equal(input.snapshot(), undefined, 'Journal data alone cannot establish readiness'); |
| 200 | progress('live'); await wait(() => input.connected); |
| 201 | snapshot = input.snapshot(Date.now() + 800); |
| 202 | assert.ok(!snapshot || !compilePetTelemetry(snapshot.events, snapshot.duration).at(-1).waiting); |
| 203 | assert.equal(requests.length, 1); assert.equal(new URL(requests[0], 'http://local').searchParams.get('progress'), 'true'); |
| 204 | assert.deepEqual(reports, []); |
| 205 | }); |
| 206 | |
| 207 | test('a Runtime without replay progress stops explicitly before any old request becomes current', { timeout: 5000 }, async t => { |
| 208 | let response, input, closed = false; |
| 209 | const reports = []; |
| 210 | const server = createServer((_req, res) => { |
| 211 | response = res; res.once('close', () => { closed = true; }); |
| 212 | res.writeHead(200, { 'content-type': 'text/event-stream' }); |
| 213 | res.write('data: ' + JSON.stringify({ seq: 1, event: 'user_input.required', thread_id: 'fixture', timestamp: new Date().toISOString(), payload: { id: 'old' } }) + '\n\n'); |
| 214 | }); |
| 215 | await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); |
| 216 | t.after(async () => { await input?.close(); response?.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); }); |
| 217 | input = await followRuntime({ baseUrl: `http://127.0.0.1:${server.address().port}`, threadId: 'fixture', report: text => reports.push(text) }); |
| 218 | for (let n = 0; n < 200 && !reports.length; n++) await delay(10); |
| 219 | assert.match(reports.join('\n'), /stopped.*replay-progress support/); |
| 220 | assert.equal(input.connected, false); assert.equal(input.cursor, 0); assert.equal(input.snapshot(), undefined); |
| 221 | for (let n = 0; n < 100 && !closed; n++) await delay(10); |
| 222 | assert.ok(closed); |
| 223 | }); |
| 224 |