Files
redsen-lean-harness/src/lib/telemetry.mjs
T
mozempkandCopilot 383129f571 feat: scaffold redsen-lean-harness v0.1.0
Recovered from crashed session (Node OOM). Repo contains full P0-P6
scaffold: plugin.json/marketplace.json, AGENTS.md, ADRs 0001-0006,
lh CLI (init/index/graph/lane/run/memory/host/report/doctor), 10
.github/agents, 12 CLI skills, instructions, context7 mcp.json, and
unit/e2e test suite.

Fixed: run.mjs read --in-tokens/--out-tokens but tests and CLI docs
use --input-tokens/--output-tokens, so telemetry totals were always 0.
Now accepts both forms.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
2026-09-09 22:44:15 +02:00

319 lines
9.7 KiB
JavaScript

import crypto from 'node:crypto';
import fs from 'node:fs';
import path from 'node:path';
import { loadConfig, defaultConfig } from './config.mjs';
import { paths, ensureDir, readIfExists, writeFile } from './paths.mjs';
import { renderBoard } from './board.mjs';
function configFor(ctx = {}) {
return ctx.config ?? loadConfig(ctx.cwd ?? process.cwd())?.config ?? defaultConfig();
}
function pathsFor(ctx = {}) {
return paths(configFor(ctx), ctx.cwd ?? process.cwd());
}
function nowIso() {
return new Date().toISOString();
}
function spanId() {
return crypto.randomBytes(8).toString('hex');
}
function rand4() {
return crypto.randomBytes(3).toString('base64url').replace(/[^a-zA-Z0-9]/g, '').slice(0, 4).padEnd(4, '0');
}
function readHeader(ctx, runId) {
const file = path.join(pathsFor(ctx).run(runId), 'run.json');
const raw = readIfExists(file);
if (!raw) return { runId };
try {
return JSON.parse(raw);
} catch {
return { runId };
}
}
function writeHeader(ctx, runId, header) {
writeFile(path.join(pathsFor(ctx).run(runId), 'run.json'), `${JSON.stringify(header, null, 2)}\n`);
}
function refreshBoard(ctx, runId) {
const p = pathsFor(ctx);
const header = readHeader(ctx, runId);
const events = readEvents(ctx, runId);
writeFile(path.join(p.run(runId), 'board.md'), renderBoard(header, events));
}
export function newRunId() {
const d = new Date();
const stamp = d.toISOString().replace(/[-:]/g, '').replace(/\.\d{3}Z$/, '').replace('T', 'T');
return `${stamp}-${rand4()}`;
}
export function startRun(ctx = {}, { objective, spec, host, strategy, attrs = {} } = {}) {
const runId = attrs.runId ?? newRunId();
const p = pathsFor(ctx);
const dir = p.run(runId);
ensureDir(dir);
const header = {
runId,
objective: objective ?? '',
spec: spec ?? null,
host: host ?? null,
strategy: strategy ?? null,
status: 'running',
startedAt: nowIso(),
endedAt: null,
summary: null,
totals: null,
attrs,
};
writeHeader(ctx, runId, header);
appendEvent(ctx, runId, {
type: 'run.start',
name: 'run',
status: 'ok',
'gen_ai.operation.name': 'run.start',
attrs: { objective, spec, host, strategy, ...attrs },
});
return runId;
}
export function appendEvent(ctx = {}, runId, event = {}) {
if (!runId) throw new Error('runId is required');
const p = pathsFor(ctx);
const dir = p.run(runId);
ensureDir(dir);
const filled = {
ts: event.ts ?? nowIso(),
runId,
spanId: event.spanId ?? spanId(),
parentSpanId: event.parentSpanId ?? null,
type: event.type ?? 'note',
name: event.name ?? null,
status: event.status ?? null,
durationMs: event.durationMs === undefined ? null : Number(event.durationMs),
'gen_ai.request.model': event['gen_ai.request.model'] ?? null,
'gen_ai.operation.name': event['gen_ai.operation.name'] ?? event.type ?? null,
'gen_ai.usage.input_tokens':
event['gen_ai.usage.input_tokens'] === undefined ? 0 : Number(event['gen_ai.usage.input_tokens']),
'gen_ai.usage.output_tokens':
event['gen_ai.usage.output_tokens'] === undefined ? 0 : Number(event['gen_ai.usage.output_tokens']),
attrs: event.attrs ?? {},
};
fs.appendFileSync(path.join(dir, 'events.ndjson'), `${JSON.stringify(filled)}\n`, {
encoding: 'utf8',
flag: 'a',
});
try {
refreshBoard(ctx, runId);
} catch {
// Telemetry writes must survive board rendering bugs.
}
return filled;
}
export function endRun(ctx = {}, runId, { status = 'ok', summary = '' } = {}) {
const header = readHeader(ctx, runId);
const start = Date.parse(header.startedAt ?? new Date().toISOString());
const durationMs = Number.isFinite(start) ? Date.now() - start : null;
appendEvent(ctx, runId, {
type: 'run.end',
name: 'run',
status,
durationMs,
'gen_ai.operation.name': 'run.end',
attrs: { summary },
});
const events = readEvents(ctx, runId);
const totals = summarize(events).totals;
const updated = {
...header,
status,
summary,
endedAt: nowIso(),
totals,
};
writeHeader(ctx, runId, updated);
try {
refreshBoard(ctx, runId);
} catch {
// Best-effort board refresh only.
}
return updated;
}
export function readEvents(ctx = {}, runId) {
if (!runId) return [];
const file = path.join(pathsFor(ctx).run(runId), 'events.ndjson');
const raw = readIfExists(file);
if (!raw) return [];
const events = [];
for (const line of raw.split(/\r?\n/)) {
if (!line.trim()) continue;
try {
const event = JSON.parse(line);
if (event && typeof event === 'object') events.push(event);
} catch {
// Ignore malformed crash residue.
}
}
return events;
}
export function latestRunId(ctx = {}) {
const runs = pathsFor(ctx).runs;
try {
const dirs = fs
.readdirSync(runs, { withFileTypes: true })
.filter((entry) => entry.isDirectory())
.map((entry) => entry.name)
.sort();
return dirs.at(-1) ?? null;
} catch {
return null;
}
}
function emptyBucket() {
return {
events: 0,
inputTokens: 0,
outputTokens: 0,
totalTokens: 0,
durationMs: 0,
toolCalls: 0,
errors: 0,
status: null,
};
}
function add(bucket, event) {
const input = Number(event?.['gen_ai.usage.input_tokens'] ?? 0) || 0;
const output = Number(event?.['gen_ai.usage.output_tokens'] ?? 0) || 0;
bucket.events += 1;
bucket.inputTokens += input;
bucket.outputTokens += output;
bucket.totalTokens += input + output;
bucket.durationMs += Number(event?.durationMs ?? 0) || 0;
if (event.type === 'tool.call') bucket.toolCalls += 1;
if (event.type === 'error' || event.status === 'fail') bucket.errors += 1;
if (event.status) bucket.status = event.status;
}
function sortedEntries(obj) {
return Object.fromEntries(Object.entries(obj).sort(([a], [b]) => a.localeCompare(b)));
}
export function summarize(events = []) {
const totals = {
inputTokens: 0,
outputTokens: 0,
totalTokens: 0,
wallMs: 0,
toolCalls: 0,
errors: 0,
};
const byPhase = {};
const byAgent = {};
const byModel = {};
const byLane = {};
const verify = { passed: 0, failed: 0 };
const gates = {};
const ralph = { iterations: 0, retries: 0 };
const phaseStarts = {};
const laneStarts = {};
let activePhase = null;
const validDates = events.map((e) => Date.parse(e.ts)).filter(Number.isFinite);
if (validDates.length) totals.wallMs = Math.max(...validDates) - Math.min(...validDates);
for (const event of events) {
const input = Number(event?.['gen_ai.usage.input_tokens'] ?? 0) || 0;
const output = Number(event?.['gen_ai.usage.output_tokens'] ?? 0) || 0;
totals.inputTokens += input;
totals.outputTokens += output;
totals.totalTokens += input + output;
if (event.type === 'tool.call') totals.toolCalls += 1;
if (event.type === 'error' || event.status === 'fail') totals.errors += 1;
if (event.type === 'phase.start') {
activePhase = event.name ?? 'phase';
phaseStarts[activePhase] = Date.parse(event.ts);
}
const phaseName = event.attrs?.phase ?? (event.type?.startsWith('phase.') ? event.name : activePhase);
if (phaseName) {
byPhase[phaseName] ??= emptyBucket();
add(byPhase[phaseName], event);
if (event.type === 'phase.end') {
const start = phaseStarts[phaseName];
const end = Date.parse(event.ts);
if (Number.isFinite(start) && Number.isFinite(end)) byPhase[phaseName].wallMs = Math.max(0, end - start);
if (activePhase === phaseName) activePhase = null;
}
}
const agentName = event.attrs?.agent ?? (event.type?.startsWith('agent.') ? event.name : null);
if (agentName) {
byAgent[agentName] ??= emptyBucket();
add(byAgent[agentName], event);
}
const model = event['gen_ai.request.model'];
if (model) {
byModel[model] ??= emptyBucket();
add(byModel[model], event);
}
const laneId = event.attrs?.laneId ?? event.attrs?.lane ?? (event.type?.startsWith('lane.') ? event.name : null);
if (laneId) {
byLane[laneId] ??= { ...emptyBucket(), kind: event.attrs?.kind ?? null, iterations: 0, retries: 0, wallMs: 0 };
add(byLane[laneId], event);
byLane[laneId].kind = event.attrs?.kind ?? byLane[laneId].kind;
if (event.type === 'lane.start') laneStarts[laneId] = Date.parse(event.ts);
if (event.type === 'lane.end') {
const start = laneStarts[laneId];
const end = Date.parse(event.ts);
if (Number.isFinite(start) && Number.isFinite(end)) byLane[laneId].wallMs = Math.max(0, end - start);
}
if (event.type === 'ralph.iteration') {
byLane[laneId].iterations += 1;
if (event.status === 'fail' || event.attrs?.retry) byLane[laneId].retries += 1;
}
}
if (event.type === 'ralph.iteration') {
ralph.iterations += 1;
if (event.status === 'fail' || event.attrs?.retry) ralph.retries += 1;
}
if (event.type === 'verify') {
if (event.status === 'ok') verify.passed += 1;
if (event.status === 'fail') verify.failed += 1;
}
if (event.type === 'gate') {
const name = event.name ?? 'gate';
gates[name] ??= { passed: 0, failed: 0, skipped: 0, blocked: 0, lastStatus: null };
if (event.status === 'ok') gates[name].passed += 1;
else if (event.status === 'fail') gates[name].failed += 1;
else if (event.status === 'skipped') gates[name].skipped += 1;
else if (event.status === 'blocked') gates[name].blocked += 1;
gates[name].lastStatus = event.status ?? gates[name].lastStatus;
}
}
return {
totals,
byPhase: sortedEntries(byPhase),
byAgent: sortedEntries(byAgent),
byModel: sortedEntries(byModel),
byLane: sortedEntries(byLane),
ralph,
verify,
gates: sortedEntries(gates),
};
}