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 };