返回 CodeWhale
pet.mjs
根目录 / pet / scripts / pet.mjs
1 #!/usr/bin/env node
2 /** Read-only adapter: existing Whalesong ingestion is the only event parser. */
3 import { readFile, stat, open } from 'node:fs/promises';
4 import { watch } from 'node:fs';
5 import { basename, dirname } from 'node:path';
6 import { importTrace } from '../dist/core/ingest.js';
7 import { compilePetTelemetry, encodePetJSONL, encodePetTSV } from '../dist/core/pet-telemetry.js';
8 import { petDemoEvents } from '../dist/core/pet-demo.js';
9 import { followRuntime } from './lib/pet-runtime.mjs';
10 import { createPetRecorder } from './lib/pet-recorder.mjs';
11
12 const args = process.argv.slice(2);
13 const option = name => args.find(a => a.startsWith(`--${name}=`))?.slice(name.length + 3);
14 if (args.includes('--help')) {
15 console.log('node scripts/pet.mjs --input=trace.jsonl --output=pet.jsonl [--trace=ID] [--format=jsonl|tsv] [--watch]\nnode scripts/pet.mjs --runtime=http://127.0.0.1:7878 --thread=ID --output=pet.jsonl [--segment-buckets=216000] [--resume]\nUse --demo instead of --input for synthetic telemetry. Output must not already exist unless --resume is used for live recording. A resumed recorder preserves the previous segment and starts unknown at the same path.\nLive recording rotates at 216000 buckets or 64 MiB into OUTPUT.segment-NNNNNN.jsonl and continues at the same live path. All archives are retained.\nRuntime reads only the existing local event journal. Optional authentication comes from CODEWHALE_RUNTIME_TOKEN; never put a token in the URL. No agent or provider is started.');
16 process.exit(0);
17 }
18 let output, recorder, monitor, timer, runtime;
19 try {
20 for (const a of args) if (!['--demo', '--watch', '--resume'].includes(a) && !/^--(input|output|trace|format|runtime|thread|segment-buckets)=.+/.test(a)) throw new Error('Unknown or empty option. Use --help.');
21 const input = option('input'), runtimeURL = option('runtime'), path = option('output'), format = option('format') ?? 'jsonl', live = args.includes('--watch') || !!runtimeURL;
22 if (!path || [!!input, args.includes('--demo'), !!runtimeURL].filter(Boolean).length !== 1
23 || !['jsonl', 'tsv'].includes(format) || live && format !== 'jsonl' || args.includes('--watch') && !input
24 || !!runtimeURL !== !!option('thread') || option('trace') && !input || option('segment-buckets') && !live || args.includes('--resume') && !live)
25 throw new Error('Choose one input source, an unused --output path (or --resume), and JSONL for live recording. Runtime requires --thread.');
26 const load = async () => {
27 if (!input) return { events: petDemoEvents(), duration: 80_000 };
28 if ((await stat(input)).size > 64 * 1024 * 1024) throw new Error('Input exceeds 64 MiB.');
29 const traces = importTrace(await readFile(input, 'utf8'), input, { privacy: 'metadata' });
30 const trace = option('trace') ? traces.find(t => t.id === option('trace')) : traces.length === 1 ? traces[0] : undefined;
31 if (!trace) throw new Error('Select an existing --trace ID when input contains multiple traces.');
32 return trace;
33 };
34 let trace = runtimeURL ? undefined : await load(), buckets = compilePetTelemetry(trace?.events ?? [], trace?.duration ?? 0);
35 const segmentBuckets = option('segment-buckets') === undefined ? 216_000 : Number(option('segment-buckets'));
36 if (live) recorder = await createPetRecorder(path, { resume: args.includes('--resume'), maxBuckets: segmentBuckets, report: text => console.error(text) });
37 else output = await open(path, 'wx', 0o600);
38 if (runtimeURL) runtime = await followRuntime({ baseUrl: runtimeURL, threadId: option('thread'),
39 token: process.env.CODEWHALE_RUNTIME_TOKEN, report: text => console.error(text) });
40 if (!live) {
41 await output.writeFile(format === 'tsv' ? encodePetTSV(buckets) : encodePetJSONL(buckets));
42 await output.close(); output = undefined;
43 console.log(`Wrote ${buckets.length} pet buckets (${args.includes('--demo') ? 'demo' : 'trace replay'}).`);
44 } else {
45 // The driver owns wall time. The core only sees recorded relative timestamps.
46 const started = performance.now(), startedWall = Date.now();
47 const origin = trace && 'originTime' in trace && trace.originTime ? Date.parse(trace.originTime) : NaN;
48 const offset = Number.isFinite(origin) ? Math.max(0, Date.now() - origin) : trace?.duration ?? 0;
49 let dirty = false, running = false, sequence = 0, failed = false, stopping = false, lastBin = -1, outage = false;
50 const empty = compilePetTelemetry([])[0];
51 if (input) {
52 monitor = watch(dirname(input), (_event, filename) => { if (!filename || String(filename) === basename(input)) dirty = true; });
53 monitor.on('error', () => { failed = true; dirty = true; });
54 }
55 const tick = async () => {
56 if (running || stopping) return;
57 running = true;
58 try {
59 const elapsed = performance.now() - started, target = Math.floor(elapsed / 400);
60 if (sequence > target) return;
61 if (runtime) {
62 failed = !runtime.connected;
63 if (!failed) {
64 try {
65 trace = runtime.snapshot(startedWall + elapsed);
66 } catch { failed = true; console.error('Runtime snapshot is invalid; recording an unobserved gap.'); }
67 }
68 }
69 if (dirty) {
70 dirty = false;
71 try { trace = await load(); buckets = compilePetTelemetry(trace.events, trace.duration); failed = false; }
72 catch { failed = true; console.error('Source unavailable or invalid; recording an unobserved gap.'); }
73 }
74 // A stalled host records skipped intervals as unknown instead of silently
75 // compressing time. Never repeat onsets when timer jitter hits a source bin twice.
76 // At most one segment of them: a longer suspension would otherwise replay
77 // every missed 400 ms append (hours of I/O and archives of pure unknown)
78 // before observing again. It starts a new segment whose first bucket is
79 // unknown, as --resume does after an outage; the gap length is reported,
80 // not recorded.
81 if (target - sequence > segmentBuckets) {
82 console.error(`Recorder was suspended or stalled for ${Math.round((target - sequence) * 0.4)} s; starting a new segment instead of recording that many unknown buckets.`);
83 recorder.restart(); sequence = target; outage = true;
84 }
85 while (sequence < target && !stopping) {
86 await recorder.append({ ...empty, sequence, simTimeMs: sequence * 400 }); sequence++;
87 }
88 if (stopping) return;
89 let state = empty;
90 if (outage) outage = false;
91 else if (runtime) {
92 // Seal the preceding observation interval before recording its state.
93 // A fixed recorder origin survives imports discovering older starts.
94 // Accepting this state one bucket later matches the foreground host.
95 if (!failed && trace && sequence > 0) {
96 try { state = compilePetTelemetry(trace.events, trace.duration,
97 sequence - 1, startedWall - Date.parse(trace.originTime))[0] ?? empty; }
98 catch { console.error('Runtime snapshot is invalid; recording an unobserved gap.'); }
99 }
100 } else {
101 const bin = Math.floor((offset + elapsed) / 400);
102 state = failed || sequence === 0 ? empty : buckets[bin] ?? empty;
103 if (bin === lastBin) state = { ...state, onsets: Array(13).fill(0), errors: 0 };
104 lastBin = bin;
105 }
106 await recorder.append({ ...state, sequence, simTimeMs: sequence * 400 });
107 sequence++;
108 } finally { running = false; }
109 };
110 await tick();
111 timer = setInterval(() => { tick().catch(async error => { console.error(error.message); process.exitCode = 1; await stop(); }); }, 400);
112 const stop = async () => {
113 stopping = true; clearInterval(timer); monitor?.close(); await runtime?.close();
114 while (running) await new Promise(resolve => setTimeout(resolve, 5));
115 await recorder.close();
116 };
117 process.once('SIGINT', stop); process.once('SIGTERM', stop);
118 console.log('Recording local live pet states. Ctrl+C to stop.');
119 }
120 } catch (error) {
121 console.error(error.message); process.exitCode = 1;
122 clearInterval(timer); monitor?.close(); await runtime?.close(); await recorder?.close(); if (output) await output.close();
123 }
124
124 lines Plain Text