Setting the file. One moment.
Rp Telemetry · Rp Telemetry · wix/skills · Skills Docs
ContentsBack to the top of the page Bundled file Bi Sink
scripts/ rp-telemetry.js
JavaScript · 211 lines · 9 KB
15 // node scripts/rp-telemetry.js wait start [--halt <subtype> --skill <s> [--what '<text>']]
16 // node scripts/rp-telemetry.js wait end
17 // node scripts/rp-telemetry.js digest --transcript <path to CC session jsonl> [--offline]
18 // node scripts/rp-telemetry.js finalize '<rollup-json>'
19 // node scripts/rp-telemetry.js rebuild [--attempt <n>] [--push]
20 // node scripts/rp-telemetry.js status
21 //
22 // `digest` (spec 0039): call it just before `finalize`. It parses the Claude
23 // Code session transcript that produced this run — no model, no network, so a
24 // failure here must never abort the migration (`|| true` it). Pass `--offline`
25 // when re-running it later against a checkpointed or local transcript to
26 // produce the digest for a run that died before ever calling it; omitted, the
27 // digest is the in-run call and necessarily misses the final few turns it is
28 // itself part of (`source.tail_complete: false`). Also copies this run's
29 // sessions into telemetry/transcripts/<sessionId>.jsonl.gz + manifest.json.
30 //
31 // Output: one JSON object on stdout. {"ok":true,...} on success; on rejection
32 // {"ok":false,"errors":[...],"hint":...} with exit code 1 — fix and retry.
33
34 const fs = require ( 'node:fs' );
35 const path = require ( 'node:path' );
36 const zlib = require ( 'node:zlib' );
37 const recorder = require ( '../lib/telemetry-recorder.js' );
38 const transcriptDigestLib = require ( '../lib/transcript-digest.js' );
39
40 const VALUE_FLAGS = new Set ([ '--project' , '--outcome' , '--halt' , '--skill' , '--what' , '--attempt' , '--stage' , '--transcript' ]);
41 // Meter measurements are accepted as kebab-case flags so a shell script can report its own timing
42 // without composing JSON: --api-ms 1234 --api-calls 73.
43 const METER_FLAGS = new Map (
44 [ 'model_ms' , 'api_ms' , 'script_ms' , 'input_tokens' , 'output_tokens' , 'cache_read_tokens' , 'cache_write_tokens' , 'api_calls' , 'api_retries' ]
45 . map (( key ) => [ `--${ key . replace ( /_/ g , '-' ) }` , key]),
46 );
47
48 function parseArgs ( argv ) {
49 const positional = [];
50 const flags = {};
51 const meter = {};
52 for ( let i = 0 ; i < argv. length ; i += 1 ) {
53 const arg = argv[i];
54 if ( VALUE_FLAGS . has (arg)) {
55 flags[arg. slice ( 2 )] = argv[i + 1 ];
56 i += 1 ;
57 } else if ( METER_FLAGS . has (arg)) {
58 const raw = argv[i + 1 ];
59 const value = Number (raw);
60 if ( ! Number. isFinite (value)) fail ([ `${ arg } expects a number (got: ${ raw === undefined ? '(none)' : raw })` ]);
61 meter[ METER_FLAGS . get (arg)] = value;
62 i += 1 ;
63 } else if (arg === '--push' ) {
64 flags.push = true ;
65 } else if (arg === '--offline' ) {
66 flags.offline = true ;
67 } else if (arg. startsWith ( '--' )) {
68 fail ([ `unknown flag: ${ arg }` ]);
69 } else {
70 positional. push (arg);
71 }
72 }
73 return { positional, flags, meter };
74 }
75
76 function fail ( errors , hint ) {
77 process.stdout. write ( `${ JSON . stringify ({ ok: false , errors , ... (hint ? { hint } : {}) }) } \n ` );
78 process. exit ( 1 );
79 }
80
81 function parseJsonArg ( raw , label ) {
82 if (raw === undefined ) fail ([ `missing ${ label } JSON argument` ]);
83 try {
84 return JSON . parse (raw);
85 } catch (error) {
86 fail ([ `${ label } is not valid JSON: ${ error . message }` ]);
87 }
88 return undefined ;
89 }
90
91 // Locates a `stage` for a timestamp from this run's own journal, so a retry
92 // loop or hook block in the digest can be attributed to the stage it fell in
93 // without the pure parser having to know anything about telemetry stages.
94 function stageIntervalsFor ( state , nowMs ) {
95 const intervals = state.entries. map (( e ) => ({ stage: e.stage, start: Date. parse (e.start), end: Date. parse (e.end) }));
96 if (state.openStage) {
97 intervals. push ({ stage: state.openStage, start: Date. parse (state.openStageStart), end: nowMs });
98 }
99 return ( ms ) => {
100 if ( ! Number. isFinite (ms)) return null ;
101 const hit = intervals. find (( iv ) => ms >= iv.start && ms <= iv.end);
102 return hit ? hit.stage : null ;
103 };
104 }
105
106 async function runDigest ( projectDir , transcriptPath , { offline }) {
107 const journal = recorder. readJournal (projectDir);
108 if ( ! journal) fail ([ 'no active telemetry run in this project' ], "call `rp-telemetry.js start` first" );
109 const state = recorder. replay (journal.records);
110 if ( ! state.runStart) fail ([ 'telemetry journal is corrupt: missing run-start record' ]);
111
112 const windowStartMs = Date. parse (state.runStart.ts);
113 const nowMs = Date. now ();
114 const sessions = transcriptDigestLib. discoverRunSessions (transcriptPath, { windowStartMs, windowEndMs: nowMs });
115 const digest = transcriptDigestLib. computeDigest (sessions, {
116 tailComplete: offline === true ,
117 stageForTs: stageIntervalsFor (state, nowMs),
118 });
119
120 // Copy, never move (spec 0039 §4.1) — Claude Code owns the original.
121 const transcriptsDir = path. join (projectDir, 'telemetry' , 'transcripts' );
122 fs. mkdirSync (transcriptsDir, { recursive: true });
123 const manifestPath = path. join (transcriptsDir, 'manifest.json' );
124 const manifest = fs. existsSync (manifestPath) ? JSON . parse (fs. readFileSync (manifestPath, 'utf8' )) : {};
125 const runId = state.runStart.run_id;
126 const sessionIds = new Set (manifest[runId] || []);
127 for ( const session of sessions) {
128 const dest = path. join (transcriptsDir, `${ session . sessionId }.jsonl.gz` );
129 if ( ! fs. existsSync (dest)) fs. writeFileSync (dest, zlib. gzipSync (fs. readFileSync (session.path)));
130 sessionIds. add (session.sessionId);
131 }
132 manifest[runId] = [ ... sessionIds];
133 const manifestTemp = `${ manifestPath }.${ process . pid }.tmp` ;
134 fs. writeFileSync (manifestTemp, `${ JSON . stringify ( manifest , null , 2 ) } \n ` , 'utf8' );
135 fs. renameSync (manifestTemp, manifestPath);
136
137 const recorded = await recorder. transcriptDigest (projectDir, digest);
138 return { ... recorded, sessions_included: sessions. length , tail_complete: digest.source.tail_complete };
139 }
140
141 async function main () {
142 const { positional , flags , meter } = parseArgs (process.argv. slice ( 2 ));
143 const [ command , ... rest ] = positional;
144 const projectDir = flags.project || process. cwd ();
145 if ( ! fs. existsSync (projectDir) || ! fs. statSync (projectDir). isDirectory ()) {
146 fail ([ `project directory not found: ${ projectDir }` ], 'pass --project <migration project dir>' );
147 }
148
149 let result;
150 switch (command) {
151 case 'start' :
152 result = await recorder. start (projectDir, rest[ 0 ] === undefined ? {} : parseJsonArg (rest[ 0 ], 'dims' ));
153 break ;
154 case 'dims' :
155 result = await recorder. dims (projectDir, parseJsonArg (rest[ 0 ], 'dims' ));
156 break ;
157 case 'record' :
158 result = await recorder. record (projectDir, parseJsonArg (rest[ 0 ], 'event' ));
159 break ;
160 case 'stage' :
161 result = recorder. stage (projectDir, rest[ 0 ], rest[ 1 ], { outcome: flags.outcome });
162 break ;
163 case 'meter' : {
164 // Flags and a JSON payload both work; flags win on conflict so a wrapper script can override
165 // a template payload without rebuilding the JSON.
166 const fromJson = rest[ 0 ] === undefined ? {} : parseJsonArg (rest[ 0 ], 'meter' );
167 const payload = { ... fromJson, ... meter };
168 if (flags.stage) payload.stage = flags.stage;
169 result = recorder. meter (projectDir, payload);
170 break ;
171 }
172 case 'wait' :
173 result = await recorder. wait (projectDir, rest[ 0 ], {
174 haltSubtype: flags.halt,
175 skill: flags.skill,
176 what: flags.what,
177 });
178 break ;
179 case 'digest' :
180 if ( ! flags.transcript) fail ([ '--transcript <path to Claude Code session jsonl> is required' ]);
181 if ( ! fs. existsSync (flags.transcript)) fail ([ `transcript not found: ${ flags . transcript }` ]);
182 result = await runDigest (projectDir, flags.transcript, { offline: flags.offline === true });
183 break ;
184 case 'finalize' :
185 result = await recorder. finalize (projectDir, parseJsonArg (rest[ 0 ], 'rollup' ));
186 break ;
187 case 'rebuild' :
188 result = await recorder. rebuild (projectDir, {
189 attempt: flags.attempt === undefined ? undefined : Number (flags.attempt),
190 push: flags.push === true ,
191 });
192 break ;
193 case 'status' :
194 result = recorder. status (projectDir);
195 break ;
196 default :
197 fail (
198 [ `unknown command: ${ command || '(none)'}` ],
199 'commands: start, dims, record, stage start|end, meter, wait start|end, digest, finalize, rebuild, status' ,
200 );
201 return ;
202 }
203 process.stdout. write ( `${ JSON . stringify ({ ok: true , ... result }) } \n ` );
204 }
205
206 main (). catch (( error ) => {
207 if (error instanceof recorder . ValidationError ) {
208 fail (error.errors, error.hint || undefined );
209 }
210 fail ([error && error.message ? error.message : String (error)]);
211 });