semillero-INVIMA/benchmarks/scripts/benchmark-lib.js

348 lines
12 KiB
JavaScript

const fs = require('fs');
const os = require('os');
const path = require('path');
const { performance } = require('perf_hooks');
const rootDir = path.resolve(__dirname, '..', '..');
const benchmarkDir = path.join(rootDir, 'benchmarks');
const resultsDir = path.join(benchmarkDir, 'results');
let PgClient = null;
try {
({ Client: PgClient } = require('pg'));
} catch {
({ Client: PgClient } = require(path.join(rootDir, 'backend', 'node_modules', 'pg')));
}
const readJson = (file) => JSON.parse(fs.readFileSync(file, 'utf8'));
const config = readJson(path.join(benchmarkDir, 'config.json'));
const ensureDir = (dir) => fs.mkdirSync(dir, { recursive: true });
const envInt = (name, fallback, min = 0, max = Number.MAX_SAFE_INTEGER) => {
const parsed = Number.parseInt(String(process.env[name] ?? ''), 10);
if (Number.isNaN(parsed)) return fallback;
return Math.max(min, Math.min(max, parsed));
};
const getDbConfig = (poolMax = 1) => ({
host: process.env.PGHOST || 'localhost',
user: process.env.PGUSER || 'postgres',
password: process.env.PGPASSWORD || '1234',
database: process.env.PGDATABASE || 'invima',
port: envInt('PGPORT', 5432, 1, 65535),
connectionTimeoutMillis: envInt('PGCONNECT_TIMEOUT_MS', 5000, 100, 60000),
statement_timeout: envInt('BENCHMARK_STATEMENT_TIMEOUT_MS', config.statementTimeoutMs, 1000, 600000),
application_name: 'semillero-benchmark',
max: poolMax
});
const createClient = async () => {
const client = new PgClient(getDbConfig());
await client.connect();
await client.query(`SET statement_timeout = ${envInt('BENCHMARK_STATEMENT_TIMEOUT_MS', config.statementTimeoutMs, 1000, 600000)}`);
return client;
};
const selectedQueryIds = () =>
String(process.env.BENCHMARK_QUERY_IDS || '')
.split(',')
.map((item) => item.trim())
.filter(Boolean);
const loadQueries = () => {
const selected = new Set(selectedQueryIds());
return config.queries
.filter((query) => !selected.size || selected.has(query.id))
.map((query) => ({
...query,
sql: fs.readFileSync(path.join(benchmarkDir, query.sqlFile), 'utf8').trim()
}));
};
const explainSql = (sql) => `EXPLAIN (ANALYZE, BUFFERS, VERBOSE, SETTINGS, FORMAT JSON) ${sql.replace(/;\s*$/, '')}`;
const flattenPlans = (plan) => {
const out = [];
const visit = (node) => {
if (!node) return;
out.push(node);
for (const child of node.Plans || []) visit(child);
};
visit(plan);
return out;
};
const planMetrics = (explainRows) => {
const doc = explainRows?.[0]?.['QUERY PLAN']?.[0] || {};
const plan = doc.Plan || {};
const nodes = flattenPlans(plan);
const sum = (field) => nodes.reduce((acc, item) => acc + Number(item[field] || 0), 0);
return {
planning_ms: Number(doc['Planning Time'] || 0),
execution_ms: Number(doc['Execution Time'] || 0),
estimated_rows: Number(plan['Plan Rows'] || 0),
actual_rows: Number(plan['Actual Rows'] || 0),
shared_hit_blocks: sum('Shared Hit Blocks'),
shared_read_blocks: sum('Shared Read Blocks'),
shared_written_blocks: sum('Shared Written Blocks'),
temp_read_blocks: sum('Temp Read Blocks'),
temp_written_blocks: sum('Temp Written Blocks'),
plan_node_type: plan['Node Type'] || '',
plan_summary: doc
};
};
const activeConnections = async (client) => {
try {
const { rows } = await client.query(
"SELECT COUNT(*)::int AS active FROM pg_stat_activity WHERE datname = current_database() AND state <> 'idle'"
);
return Number(rows[0]?.active || 0);
} catch {
return null;
}
};
const resourceStart = () => ({
at: performance.now(),
cpu: process.cpuUsage(),
usage: process.resourceUsage(),
mem: process.memoryUsage(),
freeMem: os.freemem()
});
const resourceEnd = (start, elapsedMs) => {
const cpu = process.cpuUsage(start.cpu);
const usage = process.resourceUsage();
const fsRead = Number(usage.fsRead || 0) - Number(start.usage.fsRead || 0);
const fsWrite = Number(usage.fsWrite || 0) - Number(start.usage.fsWrite || 0);
const cpuMs = (cpu.user + cpu.system) / 1000;
const mem = process.memoryUsage();
return {
cpu_ms: cpuMs,
cpu_percent_process: elapsedMs > 0 ? (cpuMs / elapsedMs) * 100 : null,
memory_rss_mb: mem.rss / 1024 / 1024,
memory_heap_used_mb: mem.heapUsed / 1024 / 1024,
system_memory_used_mb: (os.totalmem() - os.freemem()) / 1024 / 1024,
disk_fs_read_ops: fsRead,
disk_fs_write_ops: fsWrite
};
};
const executeMeasured = async (client, query, repetition, phase, concurrencyLevel = 1) => {
const startedAt = new Date().toISOString();
const start = resourceStart();
const startMs = performance.now();
let explain = {};
let rowsReturned = 0;
let success = false;
let errorMessage = '';
let httpStatus = 'SQL_ONLY';
let active = null;
try {
const result = await client.query(explainSql(query.sql));
const elapsedMs = performance.now() - startMs;
explain = planMetrics(result.rows);
rowsReturned = explain.actual_rows;
active = await activeConnections(client);
success = true;
return {
phase,
query_id: query.id,
query_name: query.name,
query_group: query.group,
repetition,
concurrency_level: concurrencyLevel,
timestamp: startedAt,
planning_ms: explain.planning_ms,
execution_ms: explain.execution_ms,
endpoint_total_ms: elapsedMs,
http_status: httpStatus,
rows_returned: rowsReturned,
estimated_rows: explain.estimated_rows,
actual_rows: explain.actual_rows,
shared_hit_blocks: explain.shared_hit_blocks,
shared_read_blocks: explain.shared_read_blocks,
shared_written_blocks: explain.shared_written_blocks,
temp_read_blocks: explain.temp_read_blocks,
temp_written_blocks: explain.temp_written_blocks,
active_pg_connections: active,
success,
error_message: errorMessage,
...resourceEnd(start, elapsedMs)
};
} catch (error) {
const elapsedMs = performance.now() - startMs;
active = await activeConnections(client);
errorMessage = String(error?.message || error);
return {
phase,
query_id: query.id,
query_name: query.name,
query_group: query.group,
repetition,
concurrency_level: concurrencyLevel,
timestamp: startedAt,
planning_ms: null,
execution_ms: null,
endpoint_total_ms: elapsedMs,
http_status: httpStatus,
rows_returned: 0,
estimated_rows: null,
actual_rows: null,
shared_hit_blocks: null,
shared_read_blocks: null,
shared_written_blocks: null,
temp_read_blocks: null,
temp_written_blocks: null,
active_pg_connections: active,
success: false,
error_message: errorMessage,
...resourceEnd(start, elapsedMs)
};
}
};
const csvEscape = (value) => {
if (value === null || value === undefined) return '';
const text = String(value);
if (/[",\n\r]/.test(text)) return `"${text.replace(/"/g, '""')}"`;
return text;
};
const writeCsv = (file, rows) => {
ensureDir(path.dirname(file));
if (!rows.length) {
fs.writeFileSync(file, '', 'utf8');
return;
}
const headers = Object.keys(rows[0]);
const body = [headers.join(',')]
.concat(rows.map((row) => headers.map((header) => csvEscape(row[header])).join(',')))
.join('\n');
fs.writeFileSync(file, `${body}\n`, 'utf8');
};
const percentile = (values, p) => {
if (!values.length) return null;
const sorted = [...values].sort((a, b) => a - b);
const index = Math.ceil((p / 100) * sorted.length) - 1;
return sorted[Math.max(0, Math.min(sorted.length - 1, index))];
};
const average = (values) => (values.length ? values.reduce((a, b) => a + b, 0) / values.length : null);
const stddev = (values) => {
if (values.length <= 1) return 0;
const mean = average(values);
const variance = values.reduce((acc, value) => acc + (value - mean) ** 2, 0) / (values.length - 1);
return Math.sqrt(variance);
};
const round = (value, digits = config.outputPrecisionDigits) =>
value === null || value === undefined || Number.isNaN(Number(value)) ? null : Number(Number(value).toFixed(digits));
const summarizeRows = (rows, groupKeys) => {
const groups = new Map();
for (const row of rows) {
const key = groupKeys.map((item) => row[item]).join('||');
if (!groups.has(key)) groups.set(key, []);
groups.get(key).push(row);
}
return Array.from(groups.entries()).map(([, items]) => {
const first = items[0];
const latencies = items.map((item) => Number(item.endpoint_total_ms)).filter(Number.isFinite);
const successes = items.filter((item) => item.success === true || item.success === 'true').length;
const failures = items.length - successes;
const deltaSeconds =
(Math.max(...items.map((item) => new Date(item.timestamp).getTime())) -
Math.min(...items.map((item) => new Date(item.timestamp).getTime()))) /
1000;
const totalObservedSeconds = latencies.reduce((acc, value) => acc + value, 0) / 1000;
const throughputWindow = deltaSeconds > 0 ? successes / deltaSeconds : null;
const throughputSequential = totalObservedSeconds > 0 ? successes / totalObservedSeconds : null;
const mean = average(latencies);
const sd = stddev(latencies);
const base = {};
for (const key of groupKeys) base[key] = first[key];
return {
...base,
n: items.length,
latency_min_ms: round(Math.min(...latencies)),
latency_max_ms: round(Math.max(...latencies)),
latency_avg_ms: round(mean),
latency_median_ms: round(percentile(latencies, 50)),
latency_p90_ms: round(percentile(latencies, 90)),
latency_p95_ms: round(percentile(latencies, 95)),
latency_p99_ms: round(percentile(latencies, 99)),
latency_stddev_ms: round(sd),
cv_percent: mean ? round((sd / mean) * 100) : null,
successes,
failures,
failure_rate_percent: round((failures / items.length) * 100),
throughput_per_second: round(throughputWindow || throughputSequential),
throughput_per_minute: round((throughputWindow || throughputSequential || 0) * 60),
cpu_avg_percent_process: round(average(items.map((item) => Number(item.cpu_percent_process)).filter(Number.isFinite))),
cpu_max_percent_process: round(Math.max(...items.map((item) => Number(item.cpu_percent_process)).filter(Number.isFinite))),
memory_avg_rss_mb: round(average(items.map((item) => Number(item.memory_rss_mb)).filter(Number.isFinite))),
memory_max_rss_mb: round(Math.max(...items.map((item) => Number(item.memory_rss_mb)).filter(Number.isFinite))),
disk_read_ops_total: round(items.reduce((acc, item) => acc + Number(item.disk_fs_read_ops || 0), 0)),
disk_write_ops_total: round(items.reduce((acc, item) => acc + Number(item.disk_fs_write_ops || 0), 0)),
max_active_pg_connections: Math.max(...items.map((item) => Number(item.active_pg_connections || 0))),
planning_avg_ms: round(average(items.map((item) => Number(item.planning_ms)).filter(Number.isFinite))),
execution_avg_ms: round(average(items.map((item) => Number(item.execution_ms)).filter(Number.isFinite))),
estimated_rows_avg: round(average(items.map((item) => Number(item.estimated_rows)).filter(Number.isFinite))),
actual_rows_avg: round(average(items.map((item) => Number(item.actual_rows)).filter(Number.isFinite)))
};
});
};
const writeJson = (file, data) => {
ensureDir(path.dirname(file));
fs.writeFileSync(file, `${JSON.stringify(data, null, 2)}\n`, 'utf8');
};
const readCsv = (file) => {
if (!fs.existsSync(file) || fs.statSync(file).size === 0) return [];
const lines = fs.readFileSync(file, 'utf8').trim().split(/\r?\n/);
const headers = lines.shift().split(',');
return lines.map((line) => {
const values = [];
let current = '';
let quoted = false;
for (let i = 0; i < line.length; i += 1) {
const char = line[i];
if (char === '"' && quoted && line[i + 1] === '"') {
current += '"';
i += 1;
} else if (char === '"') {
quoted = !quoted;
} else if (char === ',' && !quoted) {
values.push(current);
current = '';
} else {
current += char;
}
}
values.push(current);
return Object.fromEntries(headers.map((header, index) => [header, values[index] ?? '']));
});
};
module.exports = {
benchmarkDir,
resultsDir,
config,
createClient,
ensureDir,
executeMeasured,
getDbConfig,
loadQueries,
readCsv,
summarizeRows,
writeCsv,
writeJson
};