Files
peardock/server/services/stats.js
T
Raven Scott bbf0607aaf
Release rolling / release (push) Successful in 9m13s
Fix container remove timeouts and make force-remove reliable.
Lifecycle RPCs now use a 120s operation timeout, remove cleans up
stats/terminal/log streams then SIGKILLs before force-remove, and the
UI awaits the request instead of a fragile 30s wait race.
2026-07-15 14:12:13 -04:00

436 lines
13 KiB
JavaScript

/**
* Container stats collection and broadcast.
*
* Docker stats streams are NDJSON and chunks may contain partial lines or
* multiple JSON objects — we buffer and parse line-by-line so CPU/memory
* actually update (JSON.parse on raw chunks often fails silently).
*/
import { docker, extractIpAddress } from './docker.js'
import { peers } from '../core/peer-registry.js'
import { Pushes } from '../../shared/protocol.js'
import { recordSample, pruneMissing } from './stats-history.js'
import logger from '../utils/logger.js'
const STATS_CACHE_TTL = 1000
const STATS_BROADCAST_INTERVAL = 2000
/** @type {Record<string, object>} */
const containerStats = {}
const statsCache = new Map()
const containerActivity = new Map()
let intervalHandle = null
/**
* Docker Engine CPU % (same idea as `docker stats`).
* First sample often has empty precpu_stats → 0 until the next tick.
*/
function calculateCPUPercent(stats) {
try {
const cpu = stats?.cpu_stats
const precpu = stats?.precpu_stats
if (!cpu?.cpu_usage || !precpu?.cpu_usage) return 0
const cpuDelta = (cpu.cpu_usage.total_usage || 0) - (precpu.cpu_usage.total_usage || 0)
const systemDelta = (cpu.system_cpu_usage || 0) - (precpu.system_cpu_usage || 0)
let cpuCount = cpu.online_cpus || 0
if (!cpuCount && Array.isArray(cpu.cpu_usage.percpu_usage)) {
cpuCount = cpu.cpu_usage.percpu_usage.length
}
if (!cpuCount) cpuCount = 1
if (systemDelta > 0 && cpuDelta >= 0) {
const pct = (cpuDelta / systemDelta) * cpuCount * 100.0
// Clamp absurd spikes from clock glitches
if (!Number.isFinite(pct) || pct < 0) return 0
return Math.min(pct, cpuCount * 100)
}
return 0
} catch {
return 0
}
}
/**
* Working-set style memory (closer to `docker stats` MEM USAGE).
*/
function calculateMemoryUsage(stats) {
try {
const mem = stats?.memory_stats
if (!mem) return { usage: 0, limit: 0 }
const usage = Number(mem.usage) || 0
const limit = Number(mem.limit) || 0
const s = mem.stats || {}
let working = usage
if (s.inactive_file != null) {
// cgroup v2
working = Math.max(0, usage - Number(s.inactive_file))
} else if (s.total_inactive_file != null) {
working = Math.max(0, usage - Number(s.total_inactive_file))
} else if (s.cache != null) {
// cgroup v1
working = Math.max(0, usage - Number(s.cache))
}
return { usage: working, limit }
} catch {
return { usage: 0, limit: 0 }
}
}
function sumNetworkBytes(stats) {
let rx = 0
let tx = 0
const nets = stats?.networks
if (!nets || typeof nets !== 'object') return { rx, tx }
for (const n of Object.values(nets)) {
rx += Number(n?.rx_bytes) || 0
tx += Number(n?.tx_bytes) || 0
}
return { rx, tx }
}
function sumBlkioBytes(stats) {
let read = 0
let write = 0
const arr =
stats?.blkio_stats?.io_service_bytes_recursive ||
stats?.blkio_stats?.io_service_bytes ||
[]
if (!Array.isArray(arr)) return { read, write }
for (const e of arr) {
const op = String(e?.op || '').toLowerCase()
const v = Number(e?.value) || 0
if (op === 'read') read += v
else if (op === 'write') write += v
}
return { read, write }
}
/**
* Bytes/sec rates from consecutive cumulative counters.
* @param {object} statsData
* @param {{ rx: number, tx: number }} net
* @param {{ read: number, write: number }} blk
* @param {number} now
*/
function updateIoRates(statsData, net, blk, now) {
const prevT = statsData._ioTs || 0
const dt = prevT ? Math.max(0.001, (now - prevT) / 1000) : 0
if (dt > 0 && statsData._netRx != null) {
statsData.netRxRate = Math.max(0, (net.rx - statsData._netRx) / dt)
statsData.netTxRate = Math.max(0, (net.tx - statsData._netTx) / dt)
statsData.blkReadRate = Math.max(0, (blk.read - statsData._blkRead) / dt)
statsData.blkWriteRate = Math.max(0, (blk.write - statsData._blkWrite) / dt)
} else {
statsData.netRxRate = statsData.netRxRate || 0
statsData.netTxRate = statsData.netTxRate || 0
statsData.blkReadRate = statsData.blkReadRate || 0
statsData.blkWriteRate = statsData.blkWriteRate || 0
}
statsData._netRx = net.rx
statsData._netTx = net.tx
statsData._blkRead = blk.read
statsData._blkWrite = blk.write
statsData._ioTs = now
statsData.netRxTotal = net.rx
statsData.netTxTotal = net.tx
statsData.blkReadTotal = blk.read
statsData.blkWriteTotal = blk.write
}
function isContainerActive(statsData) {
return statsData.cpu > 0.5 || statsData.memory > 512 * 1024
}
/**
* Attach a stats stream with NDJSON line buffering.
* @param {object} statsData
* @param {import('dockerode').Container} container
*/
function attachStatsStream(statsData, container) {
let buf = ''
const onChunk = (chunk) => {
try {
buf += chunk.toString('utf8')
// Docker may send one or more JSON objects per chunk, newline-delimited
let nl
while ((nl = buf.indexOf('\n')) >= 0) {
const line = buf.slice(0, nl).trim()
buf = buf.slice(nl + 1)
if (!line) continue
try {
applyDockerStatsSample(statsData, JSON.parse(line))
} catch {
// incomplete / bad line — drop
}
}
// Also try parse if buffer is a complete single JSON object without trailing NL yet
if (buf.length > 2 && buf.startsWith('{')) {
try {
const sample = JSON.parse(buf)
buf = ''
applyDockerStatsSample(statsData, sample)
} catch {
// wait for more data
if (buf.length > 2_000_000) buf = '' // safety
}
}
} catch (err) {
logger.debug('stats chunk parse failed', { id: statsData.id, error: err.message })
}
}
container
.stats({ stream: true })
.then((statsStream) => {
statsData.stream = statsStream
statsStream.on('data', onChunk)
statsStream.on('error', (err) => {
logger.error('Stats stream error', { id: statsData.id, error: err.message })
statsData.stream = null
})
statsStream.on('close', () => {
statsData.stream = null
})
statsStream.on('end', () => {
statsData.stream = null
})
})
.catch((err) => {
logger.error('Failed to start stats stream', { id: statsData.id, error: err.message })
statsData.stream = null
})
}
async function initializeContainerStats(containerInfo) {
const container = docker.getContainer(containerInfo.Id)
let ipAddress = 'No IP Assigned'
try {
const details = await container.inspect()
ipAddress = extractIpAddress(details)
} catch (err) {
logger.debug('inspect failed for stats', { id: containerInfo.Id, error: err.message })
}
const statsData = {
id: containerInfo.Id,
name: containerInfo.Names?.[0]?.replace(/^\//, '') || 'Unknown',
cpu: 0,
memory: 0,
memoryLimit: 0,
netRxRate: 0,
netTxRate: 0,
blkReadRate: 0,
blkWriteRate: 0,
netRxTotal: 0,
netTxTotal: 0,
blkReadTotal: 0,
blkWriteTotal: 0,
ip: ipAddress,
stream: null,
updatedAt: 0,
}
// Only running containers have meaningful live stats streams
const state = String(containerInfo.State || '').toLowerCase()
if (state === 'running') {
attachStatsStream(statsData, container)
}
return statsData
}
function applyDockerStatsSample(statsData, sample) {
const now = Date.now()
statsData.cpu = calculateCPUPercent(sample)
const mem = calculateMemoryUsage(sample)
statsData.memory = mem.usage
statsData.memoryLimit = mem.limit
updateIoRates(statsData, sumNetworkBytes(sample), sumBlkioBytes(sample), now)
statsData.updatedAt = now
}
function destroyStatsEntry(id) {
const statsData = containerStats[id]
if (!statsData) return
if (statsData.stream) {
try {
statsData.stream.destroy()
} catch {
// ignore
}
}
delete containerStats[id]
statsCache.delete(id)
}
/**
* Drop live stats stream for a container before stop/remove so Docker is not
* held open by NDJSON stats attachments (can delay force-remove).
* @param {string} id
*/
export function destroyStatsForContainer(id) {
if (!id) return
const needle = String(id)
destroyStatsEntry(needle)
for (const key of Object.keys(containerStats)) {
if (key === needle) continue
if (
(needle.length >= 12 && key.startsWith(needle)) ||
(key.length >= 12 && needle.startsWith(key.slice(0, Math.min(12, key.length))))
) {
destroyStatsEntry(key)
}
}
}
async function collectContainerStats() {
// Running only for streams; we still list all so stopped ids are cleaned up
const running = await docker.listContainers({ all: false })
const all = await docker.listContainers({ all: true })
const runningIds = new Set(running.map((c) => c.Id))
const allIds = new Set(all.map((c) => c.Id))
for (const containerInfo of running) {
const existing = containerStats[containerInfo.Id]
if (!existing) {
try {
containerStats[containerInfo.Id] = await initializeContainerStats(containerInfo)
} catch (err) {
logger.error('Failed to init stats', { id: containerInfo.Id, error: err.message })
}
} else if (!existing.stream) {
// Was stopped / stream died — reattach
try {
attachStatsStream(existing, docker.getContainer(containerInfo.Id))
} catch (err) {
logger.debug('reattach stats failed', { id: containerInfo.Id, error: err.message })
}
}
}
// Drop stats for containers that no longer exist, or stop streams for exited ones
for (const id of Object.keys(containerStats)) {
if (!allIds.has(id)) {
destroyStatsEntry(id)
continue
}
if (!runningIds.has(id) && containerStats[id]?.stream) {
try {
containerStats[id].stream.destroy()
} catch {
// ignore
}
containerStats[id].stream = null
containerStats[id].cpu = 0
// keep last memory sample or zero for stopped
containerStats[id].cpu = 0
containerStats[id].memory = 0
}
}
}
export function startStatsBroadcast() {
if (intervalHandle) return
let lastBroadcast = 0
let dockerDownLogged = false
let dockerBackoffUntil = 0
intervalHandle = setInterval(async () => {
try {
const now = Date.now()
if (now < dockerBackoffUntil) return
await collectContainerStats()
dockerDownLogged = false
if (now - lastBroadcast < STATS_BROADCAST_INTERVAL) return
if (peers.size === 0) return
const aggregatedStats = []
for (const [containerId, statsData] of Object.entries(containerStats)) {
const cached = statsCache.get(containerId)
if (
cached &&
now - cached.timestamp < STATS_CACHE_TTL &&
!isContainerActive(statsData)
) {
aggregatedStats.push(cached.data)
continue
}
if (isContainerActive(statsData)) {
containerActivity.set(containerId, now)
}
const statsObj = {
id: statsData.id,
name: statsData.name,
cpu: Number(statsData.cpu) || 0,
memory: Number(statsData.memory) || 0,
memoryLimit: Number(statsData.memoryLimit) || 0,
netRxRate: Number(statsData.netRxRate) || 0,
netTxRate: Number(statsData.netTxRate) || 0,
blkReadRate: Number(statsData.blkReadRate) || 0,
blkWriteRate: Number(statsData.blkWriteRate) || 0,
netRxTotal: Number(statsData.netRxTotal) || 0,
netTxTotal: Number(statsData.netTxTotal) || 0,
blkReadTotal: Number(statsData.blkReadTotal) || 0,
blkWriteTotal: Number(statsData.blkWriteTotal) || 0,
ip: statsData.ip,
}
statsCache.set(containerId, { data: statsObj, timestamp: now })
aggregatedStats.push(statsObj)
recordSample(containerId, {
cpu: statsObj.cpu,
memory: statsObj.memory,
memoryLimit: statsObj.memoryLimit,
netRxRate: statsObj.netRxRate,
netTxRate: statsObj.netTxRate,
blkReadRate: statsObj.blkReadRate,
blkWriteRate: statsObj.blkWriteRate,
})
}
pruneMissing(Object.keys(containerStats))
// Always broadcast when we have containers — even zeros so UI can show 0.00% not "—"
if (aggregatedStats.length > 0) {
peers.broadcast(Pushes.allStats, { type: 'allStats', data: aggregatedStats })
lastBroadcast = now
}
for (const [id, cached] of statsCache.entries()) {
if (now - cached.timestamp > STATS_CACHE_TTL * 10) statsCache.delete(id)
}
for (const [id, ts] of containerActivity.entries()) {
if (now - ts > 60000) containerActivity.delete(id)
}
} catch (err) {
const msg = err.message || ''
const dockerDown =
msg.includes('ENOENT') ||
msg.includes('ECONNREFUSED') ||
msg.includes('docker.sock')
if (dockerDown) {
dockerBackoffUntil = Date.now() + 15000
if (!dockerDownLogged) {
logger.error('Docker unavailable; stats paused', { error: msg })
dockerDownLogged = true
}
} else {
logger.error('Stats broadcast failed', { error: msg })
}
}
}, 1000)
}
export function stopStatsBroadcast() {
if (intervalHandle) {
clearInterval(intervalHandle)
intervalHandle = null
}
for (const id of Object.keys(containerStats)) {
destroyStatsEntry(id)
}
}