238 lines
12 KiB
TypeScript
238 lines
12 KiB
TypeScript
// src/lib/monitor-worker.ts — 后台健康检查循环
|
||
import { dbQuery, dbExec, escapeSql } from './db'
|
||
import { config } from './config'
|
||
import { HttpChecker, DockerChecker, HealthChecker } from '@shared/lib/alert/health-checker'
|
||
import { AlertManager } from '@shared/lib/alert/alert-manager'
|
||
import { WeChatPusher } from '@shared/lib/wechat/wechat-pusher'
|
||
import type { AlertChannelConfig, ServiceStatus, AlertLevel } from '@shared/lib/alert/types'
|
||
|
||
// SQLite 返回 snake_case,与 DbService camelCase 不同,使用此类型
|
||
interface DbService { id: number; name: string; category: string; alert_level: string; checks: string; check_interval: number; check_timeout: number; enabled: number; current_status: string; status_since: string; display_order: number }
|
||
interface WorkerStatus { ticking: boolean; lastTick: number; tickCount: number; startTime: number }
|
||
|
||
export class MonitorWorker {
|
||
ticking = false
|
||
lastTick = 0
|
||
tickCount = 0
|
||
startTime = Date.now()
|
||
private healthChecker: HealthChecker
|
||
private alertManager: AlertManager
|
||
private flappingCounter = new Map<number, number[]>()
|
||
private flappingAlertSent = new Set<number>()
|
||
private lastCleanupDate = ''
|
||
|
||
constructor() {
|
||
this.healthChecker = new HealthChecker()
|
||
this.healthChecker.registerChecker('http', new HttpChecker())
|
||
this.healthChecker.registerChecker('docker', new DockerChecker())
|
||
|
||
this.alertManager = new AlertManager({
|
||
getLastAlertTime: async (channelId, serviceId, level) => {
|
||
const rows = dbQuery<{ value: string }>(`SELECT value FROM alert_settings WHERE key = 'cooldown:${channelId}:${serviceId}:${level}'`)
|
||
return rows.length > 0 ? Number(rows[0].value) : null
|
||
},
|
||
setLastAlertTime: async (channelId, serviceId, level, timestamp) => {
|
||
dbExec(`INSERT OR REPLACE INTO alert_settings (key, value, description, updated_at) VALUES ('cooldown:${escapeSql(String(channelId))}:${escapeSql(String(serviceId))}:${escapeSql(level)}', '${escapeSql(String(timestamp))}', 'alert cooldown', datetime('now', '+8 hours'))`)
|
||
},
|
||
})
|
||
}
|
||
|
||
async start(): Promise<void> {
|
||
// 恢复逻辑:重启时重置所有状态为 unknown
|
||
dbExec("UPDATE services SET current_status = 'unknown', status_since = datetime('now', '+8 hours') WHERE current_status != 'unknown'")
|
||
dbExec(`UPDATE status_history SET ended_at = datetime('now', '+8 hours'), duration_seconds = (strftime('%s', 'now') - strftime('%s', started_at)), truncated_by_restart = 1 WHERE ended_at IS NULL`)
|
||
|
||
console.log('[worker] Started. State recovered.')
|
||
this.tick()
|
||
setInterval(() => this.tick(), 30_000)
|
||
}
|
||
|
||
private async tick(): Promise<void> {
|
||
if (this.ticking) { console.log('[worker] Skipping tick, previous not completed'); return }
|
||
this.ticking = true
|
||
try {
|
||
this.lastTick = Date.now()
|
||
this.tickCount++
|
||
|
||
const services = dbQuery<DbService>('SELECT *, checks AS checks FROM services WHERE enabled = 1')
|
||
for (const svc of services) {
|
||
try {
|
||
const checks = JSON.parse(String(svc.checks || '[]'))
|
||
const results = await this.healthChecker.check(checks, Number(svc.check_timeout) || 10)
|
||
|
||
const allSuccess = results.every(r => r.success)
|
||
const newStatus: ServiceStatus = allSuccess ? 'normal' : 'abnormal'
|
||
const oldStatus = svc.current_status as ServiceStatus
|
||
|
||
if (newStatus !== oldStatus) {
|
||
// 抖动检测
|
||
this.trackFlapping(Number(svc.id), newStatus)
|
||
const isFlapping = this.isFlapping(Number(svc.id))
|
||
|
||
await this.handleStatusChange(svc as unknown as DbService, oldStatus, newStatus, results, isFlapping)
|
||
} else if (newStatus === 'abnormal') {
|
||
// 持续异常超 30 分钟发提醒
|
||
this.sendProlongedReminder(svc as unknown as DbService)
|
||
}
|
||
} catch (err) {
|
||
console.error(`[worker] Error checking service ${svc.id}:`, err)
|
||
}
|
||
}
|
||
|
||
// WAL checkpoint(每小时)
|
||
if (this.tickCount % 120 === 0) {
|
||
dbExec('PRAGMA wal_checkpoint(PASSIVE)')
|
||
}
|
||
|
||
// 数据清理(每日,UTC+8 时区)
|
||
const d = new Date()
|
||
const today = `${d.getFullYear()}-${String(d.getMonth()+1).padStart(2,'0')}-${String(d.getDate()).padStart(2,'0')}`
|
||
if (today !== this.lastCleanupDate) {
|
||
this.cleanupOldData()
|
||
this.lastCleanupDate = today
|
||
}
|
||
} finally {
|
||
this.ticking = false
|
||
}
|
||
}
|
||
|
||
private async handleStatusChange(svc: DbService, oldStatus: ServiceStatus, newStatus: ServiceStatus, results: Array<{success: boolean; error?: string; checkType: string}>, isFlapping: boolean): Promise<void> {
|
||
if (oldStatus === 'unknown' && newStatus === 'normal') {
|
||
// 首次检查正常,记录但不告警
|
||
dbExec(`INSERT INTO status_history (service_id, from_status, to_status, started_at) VALUES (${svc.id}, 'unknown', 'normal', datetime('now', '+8 hours'))`)
|
||
dbExec(`UPDATE services SET current_status = 'normal', status_since = datetime('now', '+8 hours'), updated_at = datetime('now', '+8 hours') WHERE id = ${svc.id}`)
|
||
return
|
||
}
|
||
|
||
if (newStatus === 'abnormal') {
|
||
const failedChecks = results.filter(r => !r.success).map(r => ({ type: r.checkType, error: r.error }))
|
||
const errorMsg = failedChecks.map(c => c.error).join('; ')
|
||
|
||
dbExec(`INSERT INTO status_history (service_id, from_status, to_status, started_at, error_message, failed_checks) VALUES (${svc.id}, '${oldStatus}', 'abnormal', datetime('now', '+8 hours'), ${escapeSql(errorMsg)}, ${escapeSql(JSON.stringify(failedChecks))})`)
|
||
dbExec(`UPDATE services SET current_status = 'abnormal', status_since = datetime('now', '+8 hours'), updated_at = datetime('now', '+8 hours') WHERE id = ${svc.id}`)
|
||
|
||
if (isFlapping) {
|
||
// 仅在首次检测到抖动时发一条汇总告警
|
||
const alreadyFlapping = this.flappingAlertSent.has(svc.id)
|
||
if (!alreadyFlapping) {
|
||
await this.sendFlappingAlert(svc)
|
||
this.flappingAlertSent.add(svc.id)
|
||
}
|
||
// 抖动期间跳过个别告警
|
||
} else {
|
||
this.flappingAlertSent.delete(svc.id)
|
||
const level: AlertLevel = svc.alert_level as AlertLevel
|
||
await this.evaluateAndPush(svc, level)
|
||
}
|
||
} else {
|
||
// 恢复
|
||
dbExec(`UPDATE status_history SET ended_at = datetime('now', '+8 hours'), duration_seconds = (strftime('%s', 'now') - strftime('%s', started_at)) WHERE service_id = ${svc.id} AND ended_at IS NULL`)
|
||
dbExec(`UPDATE services SET current_status = 'normal', status_since = datetime('now', '+8 hours'), updated_at = datetime('now', '+8 hours') WHERE id = ${svc.id}`)
|
||
|
||
// 恢复通知(不受免打扰限制)
|
||
await this.sendRecovery(svc)
|
||
}
|
||
}
|
||
|
||
private async evaluateAndPush(svc: DbService, level: AlertLevel): Promise<void> {
|
||
const channels = dbQuery<AlertChannelConfig & { enabled_num: number; level_critical_num: number; level_warning_num: number; level_info_num: number }>(
|
||
'SELECT *, enabled AS enabled_num, level_critical AS level_critical_num, level_warning AS level_warning_num, level_info AS level_info_num FROM alert_channels WHERE enabled = 1'
|
||
)
|
||
|
||
const results = await Promise.allSettled(
|
||
channels.map(async (ch) => {
|
||
const raw = ch as unknown as Record<string, unknown>
|
||
const channel: AlertChannelConfig = {
|
||
id: ch.id as number,
|
||
name: ch.name as string,
|
||
webhookUrl: (raw.webhook_url as string) || '',
|
||
channelType: ((raw.channel_type as string) || 'wechat') as 'wechat' | 'email',
|
||
enabled: !!ch.enabled_num,
|
||
levelCritical: !!ch.level_critical_num,
|
||
levelWarning: !!ch.level_warning_num,
|
||
levelInfo: !!ch.level_info_num,
|
||
quietEnabled: !!raw.quiet_enabled,
|
||
quietStart: (raw.quiet_start as string) || '23:00',
|
||
quietEnd: (raw.quiet_end as string) || '07:00',
|
||
quietBypassCritical: !!raw.quiet_bypass_critical,
|
||
cooldownMinutes: Number(raw.cooldown_minutes) || 5,
|
||
}
|
||
const decision = await this.alertManager.evaluate(channel, svc.id, level)
|
||
if (decision.shouldPush) {
|
||
const pusher = new WeChatPusher(channel.webhookUrl)
|
||
const msg = `级别: ${level}\n服务: ${svc.name}\n状态: 异常\n时间: ${new Date().toLocaleString('zh-CN')}`
|
||
const result = await pusher.pushMarkdown('服务异常告警', msg)
|
||
await this.alertManager.recordAlert(channel.id, svc.id, level)
|
||
dbExec(`INSERT INTO alert_history (channel_id, service_id, level, title, content, sent_at, success, response_code) VALUES (${channel.id}, ${svc.id}, '${level}', '服务异常告警', ${escapeSql(msg)}, datetime('now', '+8 hours'), ${result.success ? 1 : 0}, ${result.responseCode ?? 0})`)
|
||
}
|
||
})
|
||
)
|
||
}
|
||
|
||
private async sendRecovery(svc: DbService): Promise<void> {
|
||
// 仅推送到启用了对应告警级别的 channel
|
||
const levelCol = svc.alert_level === 'critical' ? 'level_critical' : 'level_warning'
|
||
const channels = dbQuery<{ id: number; webhook_url: string }>(
|
||
`SELECT id, webhook_url FROM alert_channels WHERE enabled = 1 AND ${levelCol} = 1`
|
||
)
|
||
await Promise.allSettled(channels.map(async (ch) => {
|
||
const pusher = new WeChatPusher(String(ch.webhook_url))
|
||
const msg = `服务: ${svc.name}\n状态: 已恢复\n时间: ${new Date().toLocaleString('zh-CN')}`
|
||
const result = await pusher.pushMarkdown('服务已恢复', msg)
|
||
dbExec(`INSERT INTO alert_history (channel_id, service_id, status_event_id, alert_type, level, title, content, sent_at, success, response_code)
|
||
VALUES (${ch.id}, ${svc.id}, (SELECT id FROM status_history WHERE service_id = ${svc.id} AND ended_at IS NOT NULL ORDER BY id DESC LIMIT 1), 'recovery', 'info', '服务已恢复', ${escapeSql(msg)}, datetime('now', '+8 hours'), ${result.success ? 1 : 0}, ${result.responseCode ?? 0})`)
|
||
}))
|
||
}
|
||
|
||
private async sendFlappingAlert(svc: DbService): Promise<void> {
|
||
const channels = dbQuery<{ id: number; webhook_url: string }>('SELECT id, webhook_url FROM alert_channels WHERE enabled = 1 AND level_warning = 1')
|
||
await Promise.allSettled(channels.map(async (ch) => {
|
||
const pusher = new WeChatPusher(String(ch.webhook_url))
|
||
const msg = `服务: ${svc.name}\n近 10 分钟内频繁状态切换\n抖动期间暂停该服务的个别告警推送`
|
||
await pusher.pushMarkdown('服务状态不稳定', msg)
|
||
}))
|
||
}
|
||
|
||
// 抖动检测:追踪所有状态切换(abnormal + recovery)
|
||
private trackFlapping(serviceId: number, _toStatus: ServiceStatus): void {
|
||
const timestamps = this.flappingCounter.get(serviceId) || []
|
||
timestamps.push(Date.now())
|
||
const cutoff = Date.now() - 10 * 60 * 1000
|
||
this.flappingCounter.set(serviceId, timestamps.filter(t => t > cutoff))
|
||
}
|
||
|
||
private isFlapping(serviceId: number): boolean {
|
||
const timestamps = this.flappingCounter.get(serviceId) || []
|
||
return timestamps.length >= 5
|
||
}
|
||
|
||
// 持续异常提醒(>30分钟)
|
||
private sendProlongedReminder(svc: DbService): void {
|
||
const key = `reminder:${svc.id}`
|
||
const rows = dbQuery<{ value: string }>(`SELECT value FROM alert_settings WHERE key = '${key}'`)
|
||
const lastReminder = rows.length > 0 ? Number(rows[0].value) : 0
|
||
const now = Date.now()
|
||
if (now - lastReminder < 30 * 60 * 1000) return // 30分钟内已提醒过
|
||
|
||
// 发送提醒
|
||
const channels = dbQuery<{ id: number; webhook_url: string }>('SELECT id, webhook_url FROM alert_channels WHERE enabled = 1 AND level_warning = 1')
|
||
Promise.allSettled(channels.map(async (ch) => {
|
||
const pusher = new WeChatPusher(String(ch.webhook_url))
|
||
await pusher.pushMarkdown('服务持续异常', `服务: ${svc.name}\n已持续异常超过 30 分钟\n时间: ${new Date().toLocaleString('zh-CN')}`)
|
||
}))
|
||
|
||
dbExec(`INSERT OR REPLACE INTO alert_settings (key, value, description, updated_at) VALUES ('${key}', '${now}', 'prolonged reminder', datetime('now', '+8 hours'))`)
|
||
}
|
||
|
||
// 数据清理(180 天)
|
||
private cleanupOldData(): void {
|
||
dbExec(`DELETE FROM alert_history WHERE id IN (SELECT id FROM alert_history WHERE sent_at < datetime('now', '-180 days', '+8 hours') LIMIT 1000)`)
|
||
dbExec(`DELETE FROM status_history WHERE id IN (SELECT id FROM status_history WHERE created_at < datetime('now', '-180 days', '+8 hours') AND ended_at IS NOT NULL LIMIT 1000)`)
|
||
console.log('[worker] Data cleanup completed')
|
||
}
|
||
|
||
getStatus(): WorkerStatus {
|
||
return { ticking: this.ticking, lastTick: this.lastTick, tickCount: this.tickCount, startTime: this.startTime }
|
||
}
|
||
}
|