This commit is contained in:
+203
-95
@@ -9,10 +9,14 @@ 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 { CONFIG } from '../../config.js'
|
||||
import logger from '../utils/logger.js'
|
||||
|
||||
const STATS_CACHE_TTL = 1000
|
||||
const STATS_BROADCAST_INTERVAL = 2000
|
||||
const STATS_CACHE_TTL = CONFIG.STATS?.CACHE_TTL_MS ?? 1000
|
||||
/** How often we tick collection when peers are connected */
|
||||
const STATS_COLLECT_INTERVAL = CONFIG.STATS?.ACTIVE_INTERVAL_MS ?? 1000
|
||||
/** Minimum gap between allStats broadcasts */
|
||||
const STATS_BROADCAST_INTERVAL = CONFIG.STATS?.INTERVAL_MS ?? 2000
|
||||
|
||||
/** @type {Record<string, object>} */
|
||||
const containerStats = {}
|
||||
@@ -20,6 +24,15 @@ const statsCache = new Map()
|
||||
const containerActivity = new Map()
|
||||
|
||||
let intervalHandle = null
|
||||
/** Peer-registry unsubscribe; set while stats service is armed */
|
||||
let unsubPeers = null
|
||||
/** True after startStatsBroadcast(); false after stopStatsBroadcast() */
|
||||
let serviceArmed = false
|
||||
let lastBroadcast = 0
|
||||
let dockerDownLogged = false
|
||||
let dockerBackoffUntil = 0
|
||||
/** Serialize ticks so overlapping async collects cannot pile up */
|
||||
let tickInFlight = false
|
||||
|
||||
/**
|
||||
* Docker Engine CPU % (same idea as `docker stats`).
|
||||
@@ -332,104 +345,199 @@ async function collectContainerStats() {
|
||||
}
|
||||
}
|
||||
|
||||
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)
|
||||
/**
|
||||
* Tear down all live stats streams and in-memory collection state.
|
||||
* Used when the last peer disconnects so Docker is not polled idle.
|
||||
*/
|
||||
function destroyAllStatsEntries() {
|
||||
for (const id of Object.keys(containerStats)) {
|
||||
destroyStatsEntry(id)
|
||||
}
|
||||
statsCache.clear()
|
||||
containerActivity.clear()
|
||||
}
|
||||
|
||||
export function stopStatsBroadcast() {
|
||||
/**
|
||||
* Pause collection: stop interval and drop Docker stats streams.
|
||||
* Safe to call when already paused.
|
||||
*/
|
||||
export function pauseStatsCollection() {
|
||||
if (intervalHandle) {
|
||||
clearInterval(intervalHandle)
|
||||
intervalHandle = null
|
||||
}
|
||||
for (const id of Object.keys(containerStats)) {
|
||||
destroyStatsEntry(id)
|
||||
destroyAllStatsEntries()
|
||||
tickInFlight = false
|
||||
logger.debug('Stats collection paused (no peers)')
|
||||
}
|
||||
|
||||
async function statsTick() {
|
||||
if (tickInFlight) return
|
||||
if (peers.size === 0) {
|
||||
pauseStatsCollection()
|
||||
return
|
||||
}
|
||||
tickInFlight = true
|
||||
try {
|
||||
const now = Date.now()
|
||||
if (now < dockerBackoffUntil) return
|
||||
|
||||
await collectContainerStats()
|
||||
dockerDownLogged = false
|
||||
|
||||
// Peer may have left while we were collecting
|
||||
if (peers.size === 0) {
|
||||
pauseStatsCollection()
|
||||
return
|
||||
}
|
||||
|
||||
if (now - lastBroadcast < STATS_BROADCAST_INTERVAL) 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 })
|
||||
}
|
||||
} finally {
|
||||
tickInFlight = false
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resume collection when at least one peer is connected.
|
||||
* Idempotent; kicks an immediate tick so first client is not waiting a full interval.
|
||||
*/
|
||||
export function resumeStatsCollection() {
|
||||
if (!serviceArmed) return
|
||||
if (peers.size === 0) return
|
||||
if (!intervalHandle) {
|
||||
intervalHandle = setInterval(() => {
|
||||
statsTick().catch(() => {})
|
||||
}, STATS_COLLECT_INTERVAL)
|
||||
if (typeof intervalHandle.unref === 'function') intervalHandle.unref()
|
||||
logger.debug('Stats collection resumed', { peers: peers.size })
|
||||
}
|
||||
// Immediate kick so first connected client gets streams ASAP
|
||||
lastBroadcast = 0
|
||||
statsTick().catch(() => {})
|
||||
}
|
||||
|
||||
function onPeerCountChange(size, prevSize) {
|
||||
if (!serviceArmed) return
|
||||
if (size > 0 && prevSize === 0) {
|
||||
resumeStatsCollection()
|
||||
} else if (size === 0 && prevSize > 0) {
|
||||
pauseStatsCollection()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Arm the stats service. Collection only runs while peers.size > 0.
|
||||
* Call once at server boot (safe if already started).
|
||||
*/
|
||||
export function startStatsBroadcast() {
|
||||
if (serviceArmed) return
|
||||
serviceArmed = true
|
||||
lastBroadcast = 0
|
||||
dockerDownLogged = false
|
||||
dockerBackoffUntil = 0
|
||||
|
||||
if (!unsubPeers) {
|
||||
unsubPeers = peers.onChange(onPeerCountChange)
|
||||
}
|
||||
|
||||
if (peers.size > 0) {
|
||||
resumeStatsCollection()
|
||||
} else {
|
||||
// Explicit idle: no interval, no Docker stats streams
|
||||
pauseStatsCollection()
|
||||
logger.debug('Stats service armed (idle until first peer)')
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fully stop the stats service (process shutdown).
|
||||
*/
|
||||
export function stopStatsBroadcast() {
|
||||
serviceArmed = false
|
||||
if (unsubPeers) {
|
||||
try {
|
||||
unsubPeers()
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
unsubPeers = null
|
||||
}
|
||||
pauseStatsCollection()
|
||||
}
|
||||
|
||||
/** @returns {boolean} whether the collect interval is currently running */
|
||||
export function isStatsCollectionActive() {
|
||||
return Boolean(intervalHandle)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user