返回 CodeWhale
pet-runtime.test.mjs
根目录 / pet / tests / pet-runtime.test.mjs
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('a Runtime stream.end resumes from its cursor and reports what Runtime said', { timeout: 10_000 }, async t => {
58 const cursors = [], reports = [], responses = new Set();
59 const frame = value => `data: ${JSON.stringify(value)}\n\n`;
60 const server = createServer((req, res) => {
61 cursors.push(new URL(req.url, 'http://local').searchParams.get('since_seq')); responses.add(res);
62 res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1', 'x-codewhale-stream-end': '1' });
63 if (cursors.length === 1) {
64 res.write(frame({ event: 'stream.progress', state: 'replaying', thread_id: 'fixture', seq: 0 }));
65 res.write(frame({ seq: 4, previous_seq: 0, event: 'thread.updated', thread_id: 'fixture', timestamp: new Date().toISOString(), payload: {} }));
66 res.end(frame({ schema_version: 1, event: 'stream.end', kind: 'stream.end', thread_id: 'fixture', reason: 'runtime_shutdown', last_seq: 4, retryable: true }));
67 } else {
68 res.write(frame({ event: 'stream.progress', state: 'live', thread_id: 'fixture', seq: 4 }));
69 }
70 });
71 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
72 const input = await followRuntime({ baseUrl: `http://127.0.0.1:${server.address().port}`, threadId: 'fixture', report: m => reports.push(m) });
73 t.after(async () => { await input.close(); for (const res of responses) res.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); });
74 for (let i = 0; !input.connected && i < 300; i++) await delay(10);
75 assert.equal(input.connected, true);
76 assert.equal(input.cursor, 4);
77 assert.deepEqual(cursors, ['0', '4'], 'resume from stream.end last_seq, never from zero');
78 assert.deepEqual(reports, ['Runtime input paused: Runtime is shutting down; reconnecting from the last cursor.']);
79 });
80
81 test('the CLI follows real Runtime SSE envelopes through disconnect and cursor recovery, recording no prompt content', { timeout: 20_000 }, async t => {
82 let sequence = 0, connections = 0, stream, pulse;
83 const requests = [], responses = new Set();
84 const emit = (res, tool) => {
85 const now = Date.now(), previous = sequence; sequence += 7;
86 const event = { seq: sequence, previous_seq: previous, event: 'item.completed', thread_id: 'fixture-thread', item_id: `i${sequence}`,
87 timestamp: new Date(now).toISOString(), payload: { item: { id: `i${sequence}`, kind: 'tool_call', status: 'completed',
88 started_at: new Date(now - 180).toISOString(), ended_at: new Date(now).toISOString(), summary: `${tool}: fixture-private-text` }, tool } };
89 res.write(`data: ${JSON.stringify(event)}\n\n`);
90 };
91 const server = createServer((req, res) => {
92 requests.push({ method: req.method, url: req.url, authorization: req.headers.authorization });
93 res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); res.flushHeaders(); responses.add(res);
94 res.on('close', () => responses.delete(res));
95 const connection = ++connections;
96 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`);
97 if (connection === 1) { stream = res; emit(res, 'bash'); pulse = setInterval(() => emit(res, 'bash'), 120); }
98 else if (connection === 2) {
99 // Reject a hole; the next reconnect must request the same last cursor.
100 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`);
101 } else {
102 const later = setTimeout(() => { emit(res, 'browser'); pulse = setInterval(() => emit(res, 'browser'), 120); }, 700);
103 res.once('close', () => clearTimeout(later));
104 }
105 });
106 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
107 const dir = await mkdtemp(join(tmpdir(), 'pet-runtime-')), output = join(dir, 'pet.jsonl');
108 const child = spawnRecorder([`--runtime=http://127.0.0.1:${server.address().port}`, '--thread=fixture-thread', `--output=${output}`],
109 { ...process.env, CODEWHALE_RUNTIME_TOKEN: 'fixture-token' });
110 const exited = once(child, 'exit'); let log = '';
111 child.stdout.on('data', b => log += b); child.stderr.on('data', b => log += b);
112 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)); });
113 for (let i = 0; !stream && i < 60; i++) await delay(25);
114 assert.ok(stream, log);
115 const waitForRecordedChannel = async channel => {
116 const until = Date.now() + 6_000;
117 while (Date.now() < until) {
118 try {
119 const tape = decodePetJSONL(await readFile(output, 'utf8'));
120 if (tape.some(b => b.channel === channel && b.observed === 1)) return;
121 } catch { /* The first file or an in-flight final line is not ready. */ }
122 await delay(40);
123 }
124 assert.fail(`The recorder did not persist observed ${channel} work. ${log}`);
125 };
126 // Assert actual recorder output before moving the fixture to its next phase.
127 // A fixed sleep can expire before reconnect + a complete bin on a busy runner.
128 await waitForRecordedChannel('code'); clearInterval(pulse); stream.destroy();
129 await waitForRecordedChannel('browser'); child.stopRecorder(); const [code] = await exited; assert.equal(code, 0, log); clearInterval(pulse);
130 const text = await readFile(output, 'utf8'), tape = decodePetJSONL(text);
131 assert.ok(tape.some(b => b.channel === 'code' && b.observed === 1));
132 assert.ok(tape.some(b => b.sequence > 1 && b.observed === 0));
133 assert.ok(tape.some(b => b.channel === 'browser' && b.observed === 1));
134 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);
135 assert.ok(requests.every(r => r.method === 'GET' && r.url.startsWith('/v1/threads/fixture-thread/events?') && r.authorization === 'Bearer fixture-token'));
136 assert.doesNotMatch(text + log, /fixture-private-text|fixture-token/);
137 });
138
139 test('live Runtime recording retains human waits and a brief late failure exactly once between timer ticks', { timeout: 15_000 }, async t => {
140 let sequence = 0, response, answer;
141 const timers = [], requests = [];
142 const server = createServer((req, res) => {
143 requests.push({ method: req.method, url: req.url }); response = res;
144 res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); res.flushHeaders();
145 res.write(`data: ${JSON.stringify({ event: 'stream.progress', state: 'live', thread_id: 'fixture-thread', seq: 0 })}\n\n`);
146 const oldStart = new Date(Date.now() - 20_000).toISOString(), oldEnd = new Date(Date.now() - 16_000).toISOString();
147 const emit = (event, payload) => {
148 const previous_seq = sequence; sequence++;
149 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`);
150 };
151 emit('item.started', { item: { id: 'old-work', kind: 'tool_call', status: 'running', started_at: oldStart }, tool: 'bash' });
152 timers.push(setTimeout(() => emit('user_input.required', { id: 'question', request: { questions: ['fixture-private-question'] } }), 110));
153 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));
154 answer = () => emit('user_input.answered', { input_id: 'question', answers: ['fixture-private-answer'] });
155 });
156 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
157 const dir = await mkdtemp(join(tmpdir(), 'pet-runtime-lifecycle-')), output = join(dir, 'pet.jsonl');
158 const child = spawnRecorder([`--runtime=http://127.0.0.1:${server.address().port}`, '--thread=fixture-thread', `--output=${output}`]);
159 const exited = once(child, 'exit'); let log = '';
160 child.stdout.on('data', b => log += b); child.stderr.on('data', b => log += b);
161 t.after(async () => { timers.forEach(clearTimeout); if (child.exitCode === null) child.kill('SIGTERM'); response?.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); });
162 for (let i = 0; !response && i < 100; i++) await delay(20);
163 assert.ok(response, log);
164 const waitForTape = async (condition, message) => {
165 for (let i = 0; i < 125; i++) {
166 let text = '';
167 try { text = await readFile(output, 'utf8'); } catch (error) { if (error.code !== 'ENOENT') throw error; }
168 const rows = decodePetJSONL(text.slice(0, text.lastIndexOf('\n') + 1));
169 if (condition(rows)) return;
170 await delay(40);
171 }
172 assert.fail(message + '\n' + log);
173 };
174 // Drive the answer after actual recorded coverage. Wall-clock sleeps alone
175 // can stop the child before it seals the final unknown bins on a busy runner.
176 await waitForTape(rows => rows.filter(b => b.waiting).length >= 3 && rows.some(b => b.errors), 'Waiting/error receipts were not recorded');
177 answer();
178 await waitForTape(rows => rows.length >= 2 && rows.slice(-2).every(b => !b.waiting && !b.observed), 'Answered input did not expire to unknown');
179 child.stopRecorder(); const [code] = await exited; assert.equal(code, 0, log);
180 const text = await readFile(output, 'utf8'), tape = decodePetJSONL(text);
181 assert.equal(tape.reduce((sum, b) => sum + b.errors, 0), 1);
182 assert.ok(tape.some(b => b.channel === 'error' && b.observed === 1));
183 assert.ok(tape.filter(b => b.waiting).length >= 3);
184 assert.ok(tape.slice(-2).every(b => !b.waiting && !b.observed));
185 assert.doesNotMatch(text + log, /fixture-private-question|fixture-private-result|fixture-private-answer/);
186 assert.ok(requests.every(r => r.method === 'GET' && r.url.startsWith('/v1/threads/fixture-thread/events?')));
187 });
188
189
190 test('replayed requests remain unknown until catch-up, including reentry to replay after a live connection', { timeout: 10_000 }, async t => {
191 let response, input, sequence = 0;
192 const requests = [], reports = [];
193 const server = createServer((req, res) => {
194 requests.push(req.url); response = res;
195 res.writeHead(200, { 'content-type': 'text/event-stream', 'x-codewhale-event-progress': '1' }); res.flushHeaders();
196 });
197 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
198 t.after(async () => { await input?.close(); response?.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); });
199 input = await followRuntime({ baseUrl: `http://127.0.0.1:${server.address().port}`, threadId: 'fixture', report: text => reports.push(text) });
200 const wait = async predicate => {
201 for (let n = 0; n < 250 && !predicate(); n++) await delay(10);
202 assert.ok(predicate(), reports.join('\n'));
203 };
204 await wait(() => response);
205 const write = packet => response.write('data: ' + JSON.stringify(packet) + '\n\n');
206 const progress = state => write({ event: 'stream.progress', thread_id: 'fixture', seq: sequence, state });
207 const event = (event, payload, age = 0) => write({ seq: ++sequence, previous_seq: sequence - 1,
208 event, thread_id: 'fixture', turn_id: 'turn-a', timestamp: new Date(Date.now() - age).toISOString(), payload });
209 progress('replaying'); event('user_input.required', { id: 'settled-old-request' }, 60_000);
210 await wait(() => input.cursor === 1);
211 assert.equal(input.connected, false); assert.equal(input.snapshot(), undefined);
212 // Simulate a slow backlog while the 400 ms recorder clock could run.
213 await delay(450); assert.equal(input.snapshot(), undefined);
214 event('user_input.answered', { id: 'settled-old-request' }, 50_000); progress('live');
215 await wait(() => input.connected);
216 assert.equal(input.snapshot(), undefined, 'A historical answer must arrive before old pending input can become current');
217 event('user_input.required', { id: 'fresh-request' }); await wait(() => input.cursor === 3);
218 const { compilePetTelemetry } = await import('../dist/core/pet-telemetry.js');
219 let snapshot = input.snapshot(Date.now() + 400);
220 assert.ok(compilePetTelemetry(snapshot.events, snapshot.duration).some(b => b.waiting && b.observed === 1));
221 progress('replaying'); await wait(() => !input.connected); assert.equal(input.snapshot(), undefined);
222 event('user_input.answered', { id: 'fresh-request' }); await wait(() => input.cursor === 4);
223 assert.equal(input.snapshot(), undefined, 'Journal data alone cannot establish readiness');
224 progress('live'); await wait(() => input.connected);
225 snapshot = input.snapshot(Date.now() + 800);
226 assert.ok(!snapshot || !compilePetTelemetry(snapshot.events, snapshot.duration).at(-1).waiting);
227 assert.equal(requests.length, 1); assert.equal(new URL(requests[0], 'http://local').searchParams.get('progress'), 'true');
228 assert.deepEqual(reports, []);
229 });
230
231 test('a Runtime without replay progress stops explicitly before any old request becomes current', { timeout: 5000 }, async t => {
232 let response, input, closed = false;
233 const reports = [];
234 const server = createServer((_req, res) => {
235 response = res; res.once('close', () => { closed = true; });
236 res.writeHead(200, { 'content-type': 'text/event-stream' });
237 res.write('data: ' + JSON.stringify({ seq: 1, event: 'user_input.required', thread_id: 'fixture', timestamp: new Date().toISOString(), payload: { id: 'old' } }) + '\n\n');
238 });
239 await new Promise(resolve => server.listen(0, '127.0.0.1', resolve));
240 t.after(async () => { await input?.close(); response?.destroy(); server.closeAllConnections(); await new Promise(resolve => server.close(resolve)); });
241 input = await followRuntime({ baseUrl: `http://127.0.0.1:${server.address().port}`, threadId: 'fixture', report: text => reports.push(text) });
242 for (let n = 0; n < 200 && !reports.length; n++) await delay(10);
243 assert.match(reports.join('\n'), /stopped.*replay-progress support/);
244 assert.equal(input.connected, false); assert.equal(input.cursor, 0); assert.equal(input.snapshot(), undefined);
245 for (let n = 0; n < 100 && !closed; n++) await delay(10);
246 assert.ok(closed);
247 });
248
248 lines Plain Text