105 lines
2.6 KiB
JavaScript
105 lines
2.6 KiB
JavaScript
/**
|
|
* Per-session metric / anomaly push subscriptions.
|
|
*/
|
|
import { Pushes } from '../../shared/protocol.js'
|
|
import { peers } from '../core/peer-registry.js'
|
|
|
|
/**
|
|
* @param {import('../rpc/session.js').PeerSession} session
|
|
* @param {{ charts: string[], intervalMs: number }} opts
|
|
*/
|
|
export function subscribeMetrics(session, opts) {
|
|
const charts = opts.charts?.includes('*') ? ['*'] : opts.charts || ['*']
|
|
session.state.set('metricSub', {
|
|
charts,
|
|
intervalMs: opts.intervalMs || 1000,
|
|
})
|
|
return { success: true, charts, intervalMs: opts.intervalMs || 1000 }
|
|
}
|
|
|
|
/**
|
|
* @param {import('../rpc/session.js').PeerSession} session
|
|
*/
|
|
export function unsubscribeMetrics(session) {
|
|
session.state.delete('metricSub')
|
|
return { success: true }
|
|
}
|
|
|
|
/**
|
|
* @param {import('../rpc/session.js').PeerSession} session
|
|
*/
|
|
export function subscribeAnomalies(session) {
|
|
session.state.set('anomalySub', true)
|
|
return { success: true }
|
|
}
|
|
|
|
/**
|
|
* @param {import('../rpc/session.js').PeerSession} session
|
|
*/
|
|
export function unsubscribeAnomalies(session) {
|
|
session.state.delete('anomalySub')
|
|
return { success: true }
|
|
}
|
|
|
|
/**
|
|
* Fan-out a metric batch to subscribed peers (throttled per session preference).
|
|
* @param {Array<{ chart: string, context: string, ts: number, values: object }>} batch
|
|
*/
|
|
export function broadcastMetrics(batch) {
|
|
const list = peers.list()
|
|
if (!list.length) return
|
|
for (const session of list) {
|
|
const sub = session.state.get('metricSub')
|
|
if (!sub || session.closed) continue
|
|
const filtered =
|
|
sub.charts.includes('*')
|
|
? batch
|
|
: batch.filter((s) => sub.charts.includes(s.chart) || sub.charts.includes(s.context))
|
|
if (!filtered.length) continue
|
|
|
|
// simple throttle
|
|
const last = session.state.get('metricSubLast') || 0
|
|
const now = Date.now()
|
|
if (now - last < (sub.intervalMs || 1000) - 50) continue
|
|
session.state.set('metricSubLast', now)
|
|
|
|
try {
|
|
session.push(Pushes.metrics, { samples: filtered })
|
|
} catch {
|
|
// ignore dead peers
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* @param {object} anomaly
|
|
*/
|
|
export function broadcastAnomaly(anomaly) {
|
|
const list = peers.list()
|
|
if (!list.length) return
|
|
for (const session of list) {
|
|
if (!session.state.get('anomalySub') || session.closed) continue
|
|
try {
|
|
session.push(Pushes.anomaly, anomaly)
|
|
} catch {
|
|
// ignore
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* @param {object} payload
|
|
*/
|
|
export function broadcastHealth(payload) {
|
|
const list = peers.list()
|
|
if (!list.length) return
|
|
for (const session of list) {
|
|
if (session.closed) continue
|
|
try {
|
|
session.push(Pushes.health, payload)
|
|
} catch {
|
|
// ignore
|
|
}
|
|
}
|
|
}
|