Files
Toolbox/src/server/cloud_balance_watch.js

975 lines
38 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
const fs = require('fs')
const path = require('path')
const crypto = require('crypto')
const axios = require('axios')
const { ImapFlow } = require('imapflow')
const { simpleParser } = require('mailparser')
const { db } = require('./db')
const { logJSON } = require('./logger')
const cfgPath = path.join(process.cwd(), 'config', 'cloud_balance_watch.json')
const defaultConfig = () => {
const currentPort = Number(process.env.PORT || '8976') || 8976
return {
_comment: '阿里云欠费邮件监听与分发配置',
enable_listener_comment: '是否启用欠费邮件监听服务(true/false)',
enable_listener: true,
automation_port_comment: '自动服务监听端口,仅当当前 PORT 与此值一致时才启动后台监听',
automation_port: 8977,
provider_comment: '邮箱提供方标识,仅用于说明',
provider: 'aliyun_qiye_mail',
mailbox_comment: '要监听的邮箱文件夹名称',
mailbox: 'INBOX',
reconnect_delay_ms_comment: 'IMAP 断开后的重连等待毫秒数',
reconnect_delay_ms: 15000,
poll_interval_ms_comment: '兜底轮询间隔毫秒数,用于弥补 IDLE 掉线或事件漏触发',
poll_interval_ms: 120000,
health_stale_ms_comment: '超过该毫秒数未成功扫描则视为监听异常',
health_stale_ms: 300000,
fetch_batch_comment: '每次补扫最多处理多少封新邮件,避免长时间阻塞',
fetch_batch: 50,
imap_comment: 'IMAP 收信配置(敏感信息请妥善保管)',
imap: {
host_comment: 'IMAP 服务器地址',
host: 'imap.qiye.aliyun.com',
port_comment: 'IMAP 端口,SSL 通常为 993',
port: 993,
secure_comment: '是否使用 SSL/TLS',
secure: true,
user_comment: '邮箱账号',
user: 'yangxiangyuan@umer.com.cn',
pass_comment: '邮箱授权码或应用密码(敏感,请勿泄露),从环境变量读取',
pass: process.env.YANGXIANGYUAN_EMAIL_PASS || 'y7FKf48oodBj3Bcb'
},
match_rules_comment: '邮件匹配规则,只抓指定发件人的可用额度/欠费类通知',
match_rules: {
sender_exact_comment: '发件人邮箱精确匹配',
sender_exact: 'system@notice.aliyun.com',
subject_keywords_comment: '主题模糊匹配关键词,命中任意一个即可参与判定',
subject_keywords: [
'可用额度',
'预警值',
'余额',
'欠费',
'提醒'
],
body_keywords_comment: '正文二次确认关键词,降低主题改名时的漏报风险',
body_keywords: [
'阿里云',
'可用额度',
'预警值',
'余额',
'欠费'
]
},
calendar_sync_enabled_comment: '是否启用同步到日历提醒',
calendar_sync_enabled: true,
calendar_api_base_comment: '日历提醒写入地址,沿用现有 /api/calendar/events 口径',
calendar_api_base: `http://127.0.0.1:${currentPort}/api/calendar/events`,
weekly_sync_enabled_comment: '是否启用同步到周报系统(按要求默认预留,不发送)',
weekly_sync_enabled: false,
weekly_api_base_comment: '周报系统基础地址(不含 /ext/event)',
weekly_api_base: `http://127.0.0.1:${currentPort}/api/weekly`,
weekly_key_comment: '推送到周报 /api/weekly/ext/event 使用的密钥',
weekly_key: process.env.CLOUD_BALANCE_WATCH_WEEKLY_KEY || process.env.SYSTEM_AUTH_PASS_93220 || '',
}
}
if (!fs.existsSync(path.dirname(cfgPath))) fs.mkdirSync(path.dirname(cfgPath), { recursive: true })
if (!fs.existsSync(cfgPath)) fs.writeFileSync(cfgPath, JSON.stringify(defaultConfig(), null, 2))
const readConfig = () => {
try {
return JSON.parse(fs.readFileSync(cfgPath, 'utf-8'))
} catch {
return defaultConfig()
}
}
db.exec(`CREATE TABLE IF NOT EXISTS cloud_balance_alerts (
id INTEGER PRIMARY KEY AUTOINCREMENT,
alert_key TEXT NOT NULL UNIQUE,
message_id TEXT,
message_uid INTEGER,
mailbox TEXT,
sender TEXT NOT NULL,
subject TEXT NOT NULL,
subject_norm TEXT,
received_at TEXT,
text_body TEXT,
html_body TEXT,
raw_headers TEXT,
match_reason TEXT,
risk_level TEXT DEFAULT 'medium',
seen_marked INTEGER DEFAULT 0,
seen_marked_at TEXT,
calendar_event_id TEXT,
calendar_sync_status TEXT DEFAULT 'pending',
calendar_sync_error TEXT,
calendar_sync_updated_at TEXT,
weekly_external_id TEXT,
weekly_sync_status TEXT DEFAULT 'disabled',
weekly_sync_error TEXT,
weekly_sync_updated_at TEXT,
created_at TEXT DEFAULT (datetime('now')),
updated_at TEXT DEFAULT (datetime('now'))
);`)
db.exec(`CREATE TABLE IF NOT EXISTS cloud_balance_state (
key TEXT PRIMARY KEY,
value TEXT,
updated_at TEXT DEFAULT (datetime('now'))
);`)
db.exec(`CREATE TABLE IF NOT EXISTS cloud_balance_dispatch_logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
alert_id INTEGER,
channel TEXT NOT NULL,
status TEXT NOT NULL,
detail TEXT,
created_at TEXT DEFAULT (datetime('now'))
);`)
db.exec(`CREATE INDEX IF NOT EXISTS idx_cloud_balance_message_uid ON cloud_balance_alerts(message_uid);`)
db.exec(`CREATE INDEX IF NOT EXISTS idx_cloud_balance_mailbox_uid ON cloud_balance_alerts(mailbox, message_uid);`)
db.exec(`CREATE INDEX IF NOT EXISTS idx_cloud_balance_message_id ON cloud_balance_alerts(message_id);`)
db.exec(`CREATE INDEX IF NOT EXISTS idx_cloud_balance_received_at ON cloud_balance_alerts(received_at DESC);`)
const getStateStmt = db.prepare(`SELECT value FROM cloud_balance_state WHERE key = ?`)
const setStateStmt = db.prepare(`INSERT INTO cloud_balance_state (key, value, updated_at)
VALUES (@key, @value, datetime('now'))
ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = datetime('now')`)
const insertDispatchLogStmt = db.prepare(`INSERT INTO cloud_balance_dispatch_logs (alert_id, channel, status, detail) VALUES (@alert_id, @channel, @status, @detail)`)
const countAlertsStmt = db.prepare(`SELECT COUNT(1) AS total FROM cloud_balance_alerts WHERE (? = '' OR subject LIKE ? OR text_body LIKE ?)` )
const listAlertsStmt = db.prepare(`SELECT id, sender, subject, received_at, text_body, html_body, risk_level, match_reason, seen_marked, seen_marked_at, calendar_event_id, calendar_sync_status, calendar_sync_error, weekly_external_id, weekly_sync_status, weekly_sync_error, created_at, updated_at
FROM cloud_balance_alerts
WHERE (? = '' OR subject LIKE ? OR text_body LIKE ?)
ORDER BY datetime(received_at) DESC, id DESC
LIMIT ? OFFSET ?`)
const getAlertByIdStmt = db.prepare(`SELECT * FROM cloud_balance_alerts WHERE id = ?`)
const getAlertByKeyStmt = db.prepare(`SELECT * FROM cloud_balance_alerts WHERE alert_key = ?`)
const getAlertByMailboxUidStmt = db.prepare(`SELECT * FROM cloud_balance_alerts WHERE mailbox = ? AND message_uid = ? ORDER BY id DESC LIMIT 1`)
const getAlertByMessageIdStmt = db.prepare(`SELECT * FROM cloud_balance_alerts WHERE message_id = ? ORDER BY id DESC LIMIT 1`)
const getAlertByCalendarIdStmt = db.prepare(`SELECT * FROM cloud_balance_alerts WHERE calendar_event_id = ?`)
const insertAlertStmt = db.prepare(`INSERT INTO cloud_balance_alerts (
alert_key, message_id, message_uid, mailbox, sender, subject, subject_norm, received_at, text_body, html_body, raw_headers, match_reason, risk_level, seen_marked, seen_marked_at, calendar_event_id, weekly_external_id, calendar_sync_status, weekly_sync_status
) VALUES (
@alert_key, @message_id, @message_uid, @mailbox, @sender, @subject, @subject_norm, @received_at, @text_body, @html_body, @raw_headers, @match_reason, @risk_level, @seen_marked, @seen_marked_at, @calendar_event_id, @weekly_external_id, @calendar_sync_status, @weekly_sync_status
)`)
const updateAlertStmt = db.prepare(`UPDATE cloud_balance_alerts SET
message_id = @message_id,
message_uid = @message_uid,
mailbox = @mailbox,
sender = @sender,
subject = @subject,
subject_norm = @subject_norm,
received_at = @received_at,
text_body = @text_body,
html_body = @html_body,
raw_headers = @raw_headers,
match_reason = @match_reason,
risk_level = @risk_level,
updated_at = datetime('now')
WHERE id = @id`)
const updateSeenStmt = db.prepare(`UPDATE cloud_balance_alerts SET seen_marked = 1, seen_marked_at = datetime('now'), updated_at = datetime('now') WHERE id = ?`)
const updateCalendarSyncStmt = db.prepare(`UPDATE cloud_balance_alerts SET
calendar_event_id = @calendar_event_id,
calendar_sync_status = @calendar_sync_status,
calendar_sync_error = @calendar_sync_error,
calendar_sync_updated_at = datetime('now'),
updated_at = datetime('now')
WHERE id = @id`)
const updateWeeklySyncStmt = db.prepare(`UPDATE cloud_balance_alerts SET
weekly_external_id = @weekly_external_id,
weekly_sync_status = @weekly_sync_status,
weekly_sync_error = @weekly_sync_error,
weekly_sync_updated_at = datetime('now'),
updated_at = datetime('now')
WHERE id = @id`)
const updateAlertTextBodyStmt = db.prepare(`UPDATE cloud_balance_alerts SET
text_body = @text_body,
updated_at = datetime('now')
WHERE id = @id`)
const listAlertsMissingTextStmt = db.prepare(`SELECT id, html_body FROM cloud_balance_alerts WHERE (text_body IS NULL OR trim(text_body) = '') AND html_body IS NOT NULL AND trim(html_body) <> ''`)
const recentAlertsSummaryStmt = db.prepare(`SELECT
COUNT(1) AS total,
SUM(CASE WHEN risk_level = 'high' THEN 1 ELSE 0 END) AS high_risk_count,
SUM(CASE WHEN calendar_sync_status = 'ok' THEN 1 ELSE 0 END) AS calendar_ok_count,
SUM(CASE WHEN calendar_sync_status = 'error' THEN 1 ELSE 0 END) AS calendar_error_count,
SUM(CASE WHEN weekly_sync_status = 'error' THEN 1 ELSE 0 END) AS weekly_error_count
FROM cloud_balance_alerts`)
const listenerState = {
enabled: false,
running: false,
connected: false,
mailbox: '',
current_port: Number(process.env.PORT || '8976') || 8976,
last_error: '',
last_connect_at: '',
last_scan_started_at: '',
last_scan_finished_at: '',
last_scan_reason: '',
reconnect_scheduled: false
}
let client = null
let pollTimer = null
let reconnectTimer = null
let scanPromise = null
let startInvoked = false
const getScanRuntimeState = () => ({
scan_in_progress: !!scanPromise,
last_scan_started_at: listenerState.last_scan_started_at || '',
last_scan_finished_at: listenerState.last_scan_finished_at || '',
last_scan_reason: listenerState.last_scan_reason || '',
last_error: listenerState.last_error || ''
})
const isoNow = () => new Date().toISOString()
const safeString = (value) => String(value == null ? '' : value).trim()
const normalizeSpace = (value) => safeString(value).replace(/\s+/g, ' ')
const sha1 = value => crypto.createHash('sha1').update(String(value || '')).digest('hex')
// #region debug-point cloud-balance-debug-helper
const reportCloudBalanceDebug = (hypothesisId, location, msg, data = {}) => {
try {
const envPath = path.join(process.cwd(), '.dbg', 'cloud-balance-stuck.env')
let debugUrl = 'http://127.0.0.1:7777/event'
let sessionId = 'cloud-balance-stuck'
try {
const envText = fs.readFileSync(envPath, 'utf8')
debugUrl = (envText.match(/DEBUG_SERVER_URL=(.+)/) || [])[1] || debugUrl
sessionId = (envText.match(/DEBUG_SESSION_ID=(.+)/) || [])[1] || sessionId
} catch {}
fetch(debugUrl, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
sessionId,
runId: 'pre-fix',
hypothesisId,
location,
msg,
data,
ts: Date.now()
})
}).catch(() => {})
} catch {}
}
// #endregion
const finalizeScanState = ({ highestCommittedUid, fetched, matched, inserted, updated, finishedAt, summaryFinishedAt }) => {
const finalFinishedAt = safeString(finishedAt) || isoNow()
const finalSummary = {
fetched: Number(fetched || 0) || 0,
matched: Number(matched || 0) || 0,
inserted: Number(inserted || 0) || 0,
updated: Number(updated || 0) || 0,
last_uid: Number(highestCommittedUid || 0) || 0,
finished_at: safeString(summaryFinishedAt) || finalFinishedAt
}
setState('last_uid', String(Number(highestCommittedUid || 0) || 0))
listenerState.last_scan_finished_at = finalFinishedAt
setState('last_scan_finished_at', finalFinishedAt)
setState('last_scan_summary', finalSummary)
}
const setState = (key, value) => setStateStmt.run({ key, value: typeof value === 'string' ? value : JSON.stringify(value) })
const getState = (key, fallback = '') => {
const row = getStateStmt.get(key)
return row && row.value != null ? row.value : fallback
}
const writeDispatchLog = (alertId, channel, status, detail) => {
try {
insertDispatchLogStmt.run({
alert_id: alertId || null,
channel: safeString(channel) || 'unknown',
status: safeString(status) || 'unknown',
detail: typeof detail === 'string' ? detail : JSON.stringify(detail || {})
})
} catch {}
}
const shouldStartListener = (cfg = readConfig()) => {
if (cfg.enable_listener === false) return false
const currentPort = Number(process.env.PORT || '8976') || 8976
const automationPort = Number(cfg.automation_port || 8977) || 8977
return currentPort === automationPort
}
const htmlToText = html => safeString(html)
.replace(/<style[\s\S]*?<\/style>/gi, ' ')
.replace(/<script[\s\S]*?<\/script>/gi, ' ')
.replace(/<[^>]+>/g, ' ')
.replace(/&nbsp;/gi, ' ')
.replace(/&amp;/gi, '&')
.replace(/&lt;/gi, '<')
.replace(/&gt;/gi, '>')
.replace(/&quot;/gi, '"')
.replace(/&#39;/gi, '\'')
.replace(/\s+/g, ' ')
.trim()
const normalizeMailBody = (textBody, htmlBody) => {
const text = normalizeSpace(textBody)
if (text) return text
return normalizeSpace(htmlToText(htmlBody))
}
const backfillAlertTextBodies = () => {
try {
const rows = listAlertsMissingTextStmt.all()
rows.forEach(row => {
const fallback = normalizeMailBody('', row.html_body || '')
if (fallback) updateAlertTextBodyStmt.run({ id: row.id, text_body: fallback })
})
} catch {}
}
backfillAlertTextBodies()
const extractSender = parsed => {
try {
const address = parsed && parsed.from && parsed.from.value && parsed.from.value[0] && parsed.from.value[0].address
if (address) return safeString(address).toLowerCase()
} catch {}
return ''
}
const matchesAlert = (parsed, cfg) => {
const rules = cfg.match_rules || {}
const sender = extractSender(parsed)
const senderExact = safeString(rules.sender_exact).toLowerCase()
const subject = normalizeSpace(parsed && parsed.subject)
const textBody = normalizeMailBody(parsed && parsed.text, parsed && parsed.html)
const htmlBody = safeString(parsed && parsed.html)
const fullBody = normalizeSpace(`${textBody}\n${htmlToText(htmlBody)}`)
const subjectLower = subject.toLowerCase()
const bodyLower = fullBody.toLowerCase()
const subjectKeywords = Array.isArray(rules.subject_keywords) ? rules.subject_keywords.map(item => safeString(item).toLowerCase()).filter(Boolean) : []
const bodyKeywords = Array.isArray(rules.body_keywords) ? rules.body_keywords.map(item => safeString(item).toLowerCase()).filter(Boolean) : []
const subjectHits = subjectKeywords.filter(keyword => subjectLower.includes(keyword))
const bodyHits = bodyKeywords.filter(keyword => bodyLower.includes(keyword))
const bodyHitsStrict = bodyHits.filter(keyword => keyword !== '阿里云')
const senderMatched = !!senderExact && sender === senderExact
const subjectMatched = subjectHits.length > 0
const bodyMatched = bodyHitsStrict.length > 0
const matched = senderMatched && (subjectMatched || bodyMatched)
const riskLevel = subjectMatched && bodyMatched ? 'high' : (matched ? 'medium' : 'low')
const reasons = []
if (senderMatched) reasons.push(`发件人=${sender}`)
if (subjectHits.length) reasons.push(`标题命中:${subjectHits.join('|')}`)
if (bodyHitsStrict.length) reasons.push(`正文命中:${bodyHitsStrict.join('|')}`)
return {
matched,
sender,
subject,
textBody,
htmlBody,
riskLevel,
matchReason: reasons.join(';')
}
}
const buildHeadersSnapshot = parsed => {
const snap = {}
try {
snap.messageId = safeString(parsed && parsed.messageId)
snap.date = parsed && parsed.date ? new Date(parsed.date).toISOString() : ''
snap.from = extractSender(parsed)
snap.to = parsed && parsed.to && parsed.to.text ? safeString(parsed.to.text) : ''
} catch {}
return JSON.stringify(snap)
}
const buildAlertKey = payload => {
const mailbox = safeString(payload.mailbox || 'INBOX')
const messageUid = Number(payload.message_uid || 0) || 0
const messageId = safeString(payload.message_id)
let raw = ''
if (mailbox && messageUid > 0) raw = `mailbox-uid|${mailbox}|${messageUid}`
else if (messageId) raw = `message-id|${messageId}`
else raw = `${payload.sender}|${payload.subject_norm}|${payload.received_at}|${payload.message_uid || ''}`
return sha1(raw)
}
const toSqliteDateTime = value => {
if (!value) return isoNow().replace('T', ' ').slice(0, 19)
try {
return new Date(value).toISOString().replace('T', ' ').slice(0, 19)
} catch {
return isoNow().replace('T', ' ').slice(0, 19)
}
}
const upsertAlert = payload => {
const alertKey = buildAlertKey(payload)
const next = Object.assign({}, payload, {
alert_key: alertKey,
subject_norm: normalizeSpace(payload.subject).toLowerCase(),
received_at: toSqliteDateTime(payload.received_at),
calendar_event_id: payload.calendar_event_id || `cloud-arrear-${alertKey.slice(0, 16)}`,
weekly_external_id: payload.weekly_external_id || `cloud-arrear-${alertKey.slice(0, 16)}`,
calendar_sync_status: payload.calendar_sync_status || 'pending',
weekly_sync_status: payload.weekly_sync_status || 'disabled',
seen_marked: payload.seen_marked ? 1 : 0,
seen_marked_at: payload.seen_marked ? toSqliteDateTime(payload.seen_marked_at || isoNow()) : null
})
const existing = (() => {
const mailbox = safeString(next.mailbox || 'INBOX')
const messageUid = Number(next.message_uid || 0) || 0
const messageId = safeString(next.message_id)
if (mailbox && messageUid > 0) {
const row = getAlertByMailboxUidStmt.get(mailbox, messageUid)
if (row) return row
}
if (messageId) {
const row = getAlertByMessageIdStmt.get(messageId)
if (row) return row
}
return getAlertByKeyStmt.get(alertKey)
})()
if (existing) {
updateAlertStmt.run(Object.assign({}, next, { id: existing.id }))
return Object.assign({}, getAlertByIdStmt.get(existing.id), { created: false })
}
const info = insertAlertStmt.run(next)
return Object.assign({}, getAlertByIdStmt.get(info.lastInsertRowid), { created: true })
}
const buildCalendarPayload = alert => {
const start = (alert.received_at || '').replace(' ', 'T')
const detailLines = [
`来源邮箱: ${alert.sender}`,
`接收时间: ${alert.received_at || ''}`,
`本地告警ID: ${alert.id}`,
`关联主键: ${alert.calendar_event_id}`,
`匹配原因: ${alert.match_reason || ''}`,
'',
'邮件正文:',
safeString(alert.text_body || htmlToText(alert.html_body || '')) || '(空)'
]
return {
id: alert.calendar_event_id,
title: `[欠费通知] ${safeString(alert.subject) || '阿里云额度提醒'}`,
content: detailLines.join('\n'),
start,
end: start,
remindAt: start,
allDay: 0,
deleted: false,
createdAt: isoNow(),
updatedAt: isoNow()
}
}
const dispatchCalendar = async (alert, cfg = readConfig()) => {
const enabled = cfg.calendar_sync_enabled !== false
const calendarEventId = alert.calendar_event_id || `cloud-arrear-${alert.alert_key.slice(0, 16)}`
if (!enabled) {
updateCalendarSyncStmt.run({ id: alert.id, calendar_event_id: calendarEventId, calendar_sync_status: 'disabled', calendar_sync_error: '' })
return { ok: false, skipped: true, reason: 'calendar_disabled' }
}
const url = safeString(cfg.calendar_api_base)
if (!url) {
updateCalendarSyncStmt.run({ id: alert.id, calendar_event_id: calendarEventId, calendar_sync_status: 'error', calendar_sync_error: 'calendar_api_base_missing' })
return { ok: false, error: 'calendar_api_base_missing' }
}
try {
await axios.post(url, buildCalendarPayload(Object.assign({}, alert, { calendar_event_id: calendarEventId })), {
timeout: 15000,
headers: { 'Content-Type': 'application/json' }
})
updateCalendarSyncStmt.run({ id: alert.id, calendar_event_id: calendarEventId, calendar_sync_status: 'ok', calendar_sync_error: '' })
writeDispatchLog(alert.id, 'calendar', 'ok', { url, event_id: calendarEventId })
return { ok: true, event_id: calendarEventId }
} catch (e) {
const error = String(e && e.message ? e.message : e)
updateCalendarSyncStmt.run({ id: alert.id, calendar_event_id: calendarEventId, calendar_sync_status: 'error', calendar_sync_error: error })
writeDispatchLog(alert.id, 'calendar', 'error', { url, error })
return { ok: false, error, event_id: calendarEventId }
}
}
const dispatchWeekly = async (alert, cfg = readConfig()) => {
const externalId = alert.weekly_external_id || `cloud-arrear-${alert.alert_key.slice(0, 16)}`
if (cfg.weekly_sync_enabled !== true) {
updateWeeklySyncStmt.run({ id: alert.id, weekly_external_id: externalId, weekly_sync_status: 'disabled', weekly_sync_error: '' })
return { ok: false, skipped: true, reason: 'weekly_disabled' }
}
const base = safeString(cfg.weekly_api_base)
if (!base) {
updateWeeklySyncStmt.run({ id: alert.id, weekly_external_id: externalId, weekly_sync_status: 'error', weekly_sync_error: 'weekly_api_base_missing' })
return { ok: false, error: 'weekly_api_base_missing' }
}
try {
await axios.post(`${base.replace(/\/$/, '')}/ext/event`, {
key: safeString(cfg.weekly_key) || process.env.SYSTEM_AUTH_PASS_93220 || '',
source: 'cloud_arrear',
external_id: externalId,
date: String(alert.received_at || '').slice(0, 10),
title: `[欠费通知] ${safeString(alert.subject) || '阿里云额度提醒'}`,
content: safeString(alert.text_body || htmlToText(alert.html_body || '')),
status_code: 'new',
status_label: '⚠️欠费通知',
deadline: (alert.received_at || '').replace(' ', 'T')
}, {
timeout: 15000,
headers: { 'Content-Type': 'application/json' }
})
updateWeeklySyncStmt.run({ id: alert.id, weekly_external_id: externalId, weekly_sync_status: 'ok', weekly_sync_error: '' })
writeDispatchLog(alert.id, 'weekly', 'ok', { external_id: externalId })
return { ok: true, external_id: externalId }
} catch (e) {
const error = String(e && e.message ? e.message : e)
updateWeeklySyncStmt.run({ id: alert.id, weekly_external_id: externalId, weekly_sync_status: 'error', weekly_sync_error: error })
writeDispatchLog(alert.id, 'weekly', 'error', { external_id: externalId, error })
return { ok: false, error, external_id: externalId }
}
}
const markSeen = async (uid, alertId) => {
// #region debug-point cloud-balance-mark-seen-start
reportCloudBalanceDebug('D', 'cloud_balance_watch.js:markSeen:start', '[DEBUG] mark seen start', { uid, alertId })
// #endregion
if (!client) return
await client.messageFlagsAdd(uid, ['\\Seen'], { uid: true })
updateSeenStmt.run(alertId)
// #region debug-point cloud-balance-mark-seen-done
reportCloudBalanceDebug('D', 'cloud_balance_watch.js:markSeen:done', '[DEBUG] mark seen done', { uid, alertId })
// #endregion
}
const processParsedMessage = async (parsed, meta, cfg) => {
const match = matchesAlert(parsed, cfg)
// #region debug-point cloud-balance-process-message
reportCloudBalanceDebug('B', 'cloud_balance_watch.js:processParsedMessage', '[DEBUG] parsed message result', {
uid: meta && meta.uid,
mailbox: meta && meta.mailbox,
sender: match && match.sender,
subject: match && match.subject,
matched: !!(match && match.matched),
matchReason: match && match.matchReason
})
// #endregion
if (!match.matched) return { matched: false }
const row = upsertAlert({
message_id: safeString(parsed && parsed.messageId),
message_uid: meta.uid,
mailbox: meta.mailbox,
sender: match.sender,
subject: match.subject,
received_at: parsed && parsed.date ? parsed.date.toISOString() : isoNow(),
text_body: match.textBody,
html_body: match.htmlBody,
raw_headers: buildHeadersSnapshot(parsed),
match_reason: match.matchReason,
risk_level: match.riskLevel
})
await markSeen(meta.uid, row.id)
const fresh = getAlertByIdStmt.get(row.id)
await dispatchCalendar(fresh, cfg)
await dispatchWeekly(getAlertByIdStmt.get(row.id), cfg)
return { matched: true, created: !!row.created, id: row.id }
}
const closeClient = async () => {
const current = client
client = null
listenerState.connected = false
if (!current) return
try { current.removeAllListeners() } catch {}
try { await current.logout() } catch {}
}
const scheduleReconnect = () => {
if (reconnectTimer || !shouldStartListener(readConfig())) return
listenerState.reconnect_scheduled = true
const delay = Math.max(5000, Number(readConfig().reconnect_delay_ms || 15000) || 15000)
reconnectTimer = setTimeout(async () => {
reconnectTimer = null
listenerState.reconnect_scheduled = false
try {
await connectClient('reconnect')
} catch {}
}, delay)
}
const connectClient = async (reason = 'manual') => {
const cfg = readConfig()
if (!shouldStartListener(cfg)) {
listenerState.enabled = false
listenerState.running = false
return false
}
await closeClient()
const imapCfg = cfg.imap || {}
const nextClient = new ImapFlow({
host: safeString(imapCfg.host),
port: Number(imapCfg.port || 993) || 993,
secure: imapCfg.secure !== false,
auth: {
user: safeString(imapCfg.user),
pass: safeString(imapCfg.pass)
},
logger: false,
disableAutoIdle: false,
clientInfo: { name: 'TRAE-Toolbox', version: '1.0.0' }
})
nextClient.on('error', err => {
listenerState.last_error = String(err && err.message ? err.message : err)
logJSON('cloud.balance.listener.error', { error: listenerState.last_error }, 'cloud')
// #region debug-point cloud-balance-listener-error
reportCloudBalanceDebug('A', 'cloud_balance_watch.js:connectClient:error', '[DEBUG] listener error', { reason, error: listenerState.last_error })
// #endregion
})
nextClient.on('close', () => {
listenerState.connected = false
scheduleReconnect()
// #region debug-point cloud-balance-listener-close
reportCloudBalanceDebug('A', 'cloud_balance_watch.js:connectClient:close', '[DEBUG] listener close', { reason, reconnect_scheduled: true })
// #endregion
})
nextClient.on('exists', () => {
triggerScan('exists').catch(() => {})
})
await nextClient.connect()
await nextClient.mailboxOpen(safeString(cfg.mailbox) || 'INBOX', { readOnly: false })
// #region debug-point cloud-balance-connect-done
reportCloudBalanceDebug('C', 'cloud_balance_watch.js:connectClient:done', '[DEBUG] mailbox connected', {
reason,
mailbox: safeString(cfg.mailbox) || 'INBOX',
host: safeString(imapCfg.host),
port: Number(imapCfg.port || 993) || 993
})
// #endregion
client = nextClient
listenerState.enabled = true
listenerState.running = true
listenerState.connected = true
listenerState.mailbox = safeString(cfg.mailbox) || 'INBOX'
listenerState.last_connect_at = isoNow()
listenerState.last_error = ''
setState('last_connect_reason', reason)
return true
}
const scanMailbox = async reason => {
const cfg = readConfig()
if (!shouldStartListener(cfg)) return { ok: false, skipped: true, reason: 'not_automation_port' }
if (!client) await connectClient(reason)
listenerState.last_scan_reason = reason
listenerState.last_scan_started_at = isoNow()
setState('last_scan_started_at', listenerState.last_scan_started_at)
setState('last_scan_reason', reason)
let lastUid = Number(getState('last_uid', '0') || '0') || 0
let lastSenderHistoryUid = Number(getState('last_sender_history_uid', '0') || '0') || 0
let highestCommittedUid = lastUid
const batchLimit = Math.max(1, Number(cfg.fetch_batch || 50) || 50)
let fetched = 0
let matched = 0
let inserted = 0
let updated = 0
const processedUids = new Set()
try {
const range = `${Math.max(1, lastUid + 1)}:*`
// #region debug-point cloud-balance-scan-start
reportCloudBalanceDebug('A', 'cloud_balance_watch.js:scanMailbox:start', '[DEBUG] scan start', {
reason,
lastUid,
range,
batchLimit,
mailbox: listenerState.mailbox || safeString(cfg.mailbox) || 'INBOX'
})
// #endregion
const processMessage = async (message, sourceTag) => {
if (!message || !message.uid || processedUids.has(message.uid)) return
processedUids.add(message.uid)
fetched += 1
if (message.uid > highestCommittedUid) highestCommittedUid = message.uid
try {
// #region debug-point cloud-balance-message-fetch
reportCloudBalanceDebug('A', 'cloud_balance_watch.js:scanMailbox:message', '[DEBUG] processing message', {
uid: message.uid,
sourceTag,
fetched,
highestCommittedUid
})
// #endregion
const parsed = await simpleParser(message.source)
const result = await processParsedMessage(parsed, {
uid: message.uid,
mailbox: listenerState.mailbox || safeString(cfg.mailbox) || 'INBOX'
}, cfg)
if (result.matched) {
matched += 1
if (result.created) inserted += 1
else updated += 1
}
} catch (e) {
listenerState.last_error = String(e && e.message ? e.message : e)
logJSON('cloud.balance.scan.message.error', { uid: message.uid, error: listenerState.last_error }, 'cloud')
// #region debug-point cloud-balance-message-error
reportCloudBalanceDebug('D', 'cloud_balance_watch.js:scanMailbox:message:error', '[DEBUG] message error', {
uid: message && message.uid,
sourceTag,
error: listenerState.last_error
})
// #endregion
throw e
}
}
for await (const message of client.fetch(range, {
uid: true,
envelope: true,
flags: true,
internalDate: true,
source: true
}, { uid: true })) {
if (!message || !message.uid) continue
if (message.uid <= lastUid) {
highestCommittedUid = Math.max(highestCommittedUid, message.uid)
continue
}
await processMessage(message, 'uid-forward')
if (fetched >= batchLimit) break
}
let legacyCandidateUids = []
if (client && listenerState.connected) {
try {
const senderExact = safeString(cfg && cfg.match_rules && cfg.match_rules.sender_exact).toLowerCase()
const searchQuery = senderExact ? { from: senderExact } : { seen: false }
const legacyUidsRaw = await client.search(searchQuery, { uid: true })
// #region debug-point cloud-balance-search-unseen
reportCloudBalanceDebug('C', 'cloud_balance_watch.js:scanMailbox:search-unseen', '[DEBUG] legacy search result', {
reason,
lastUid,
lastSenderHistoryUid,
mode: senderExact ? 'sender-history' : 'unread-fallback',
sender: senderExact || '',
count: Array.isArray(legacyUidsRaw) ? legacyUidsRaw.length : -1
})
// #endregion
legacyCandidateUids = (Array.isArray(legacyUidsRaw) ? legacyUidsRaw : [])
.map(item => Number(item) || 0)
.filter(uid => uid > 0 && uid <= lastUid && !processedUids.has(uid) && (senderExact ? uid > lastSenderHistoryUid : true))
.sort((a, b) => a - b)
} catch (e) {
const searchError = String(e && e.message ? e.message : e)
listenerState.last_error = searchError
reportCloudBalanceDebug('C', 'cloud_balance_watch.js:scanMailbox:search-unseen:error', '[DEBUG] unseen search error', {
reason,
lastUid,
error: searchError
})
}
}
if (legacyCandidateUids.length > 0 && client && listenerState.connected) {
for await (const message of client.fetch(legacyCandidateUids.join(','), {
uid: true,
envelope: true,
flags: true,
internalDate: true,
source: true
}, { uid: true })) {
await processMessage(message, 'legacy-history')
}
const lastHistoryProcessedUid = legacyCandidateUids[legacyCandidateUids.length - 1] || 0
if (lastHistoryProcessedUid > lastSenderHistoryUid) {
lastSenderHistoryUid = lastHistoryProcessedUid
setState('last_sender_history_uid', String(lastSenderHistoryUid))
}
}
finalizeScanState({
highestCommittedUid,
fetched,
matched,
inserted,
updated
})
listenerState.last_error = ''
// #region debug-point cloud-balance-scan-finish
reportCloudBalanceDebug('A', 'cloud_balance_watch.js:scanMailbox:finish', '[DEBUG] scan finish', {
reason,
fetched,
matched,
inserted,
updated,
last_uid: highestCommittedUid
})
// #endregion
return { ok: true, fetched, matched, inserted, updated, last_uid: highestCommittedUid }
} catch (e) {
listenerState.last_error = String(e && e.message ? e.message : e)
logJSON('cloud.balance.scan.error', { reason, error: listenerState.last_error }, 'cloud')
scheduleReconnect()
finalizeScanState({
highestCommittedUid,
fetched,
matched,
inserted,
updated
})
// #region debug-point cloud-balance-scan-error
reportCloudBalanceDebug('A', 'cloud_balance_watch.js:scanMailbox:error', '[DEBUG] scan error', {
reason,
error: listenerState.last_error,
fetched,
matched,
inserted,
updated,
last_uid: highestCommittedUid
})
// #endregion
return { ok: false, error: listenerState.last_error, fetched, matched, inserted, updated, last_uid: highestCommittedUid }
}
}
const triggerScan = async reason => {
if (scanPromise) return scanPromise
scanPromise = scanMailbox(reason).finally(() => {
scanPromise = null
})
return scanPromise
}
const requestScan = reason => {
if (scanPromise) {
// #region debug-point cloud-balance-request-already-running
reportCloudBalanceDebug('E', 'cloud_balance_watch.js:requestScan:already-running', '[DEBUG] request scan already running', {
reason,
runtime: getScanRuntimeState()
})
// #endregion
return {
ok: true,
accepted: true,
already_running: true,
message: 'scan_in_progress',
runtime: getScanRuntimeState()
}
}
scanPromise = scanMailbox(reason).finally(() => {
scanPromise = null
})
// #region debug-point cloud-balance-request-started
reportCloudBalanceDebug('E', 'cloud_balance_watch.js:requestScan:started', '[DEBUG] request scan started', {
reason,
runtime: getScanRuntimeState()
})
// #endregion
return {
ok: true,
accepted: true,
already_running: false,
message: 'scan_started',
runtime: getScanRuntimeState()
}
}
const start = async () => {
if (startInvoked) return
startInvoked = true
const cfg = readConfig()
listenerState.enabled = cfg.enable_listener !== false
if (!shouldStartListener(cfg)) {
listenerState.running = false
listenerState.connected = false
return
}
const interval = Math.max(30000, Number(cfg.poll_interval_ms || 600000) || 600000)
if (!pollTimer) {
pollTimer = setInterval(() => {
triggerScan('poll').catch(() => {})
}, interval)
}
try {
await connectClient('startup')
await triggerScan('startup')
} catch (e) {
listenerState.last_error = String(e && e.message ? e.message : e)
logJSON('cloud.balance.start.error', { error: listenerState.last_error }, 'cloud')
scheduleReconnect()
}
}
const listAlerts = ({ page = 1, pageSize = 20, keyword = '' } = {}) => {
backfillAlertTextBodies()
const safePage = Math.max(1, parseInt(page, 10) || 1)
const safePageSize = Math.min(100, Math.max(1, parseInt(pageSize, 10) || 20))
const q = safeString(keyword)
const like = `%${q}%`
const total = (countAlertsStmt.get(q, like, like) || {}).total || 0
const rows = listAlertsStmt.all(q, like, like, safePageSize, (safePage - 1) * safePageSize).map(row => Object.assign({}, row, {
text_body: normalizeMailBody(row.text_body || '', row.html_body || '')
}))
return {
page: safePage,
page_size: safePageSize,
total,
rows
}
}
const getStatus = () => {
const cfg = readConfig()
const summary = recentAlertsSummaryStmt.get() || {}
const lastScanFinishedAt = listenerState.last_scan_finished_at || getState('last_scan_finished_at', '')
const staleMs = Math.max(60000, Number(cfg.health_stale_ms || 300000) || 300000)
let healthy = false
if (lastScanFinishedAt) {
const diff = Date.now() - new Date(lastScanFinishedAt).getTime()
healthy = Number.isFinite(diff) && diff <= staleMs
}
return {
config: {
enabled: cfg.enable_listener !== false,
automation_port: Number(cfg.automation_port || 8977) || 8977,
current_port: Number(process.env.PORT || '8976') || 8976,
mailbox: safeString(cfg.mailbox) || 'INBOX',
sender_exact: safeString(cfg.match_rules && cfg.match_rules.sender_exact),
calendar_sync_enabled: cfg.calendar_sync_enabled !== false,
weekly_sync_enabled: cfg.weekly_sync_enabled === true
},
listener: Object.assign({}, listenerState, {
healthy,
scan_in_progress: !!scanPromise
}),
state: {
last_uid: Number(getState('last_uid', '0') || '0') || 0,
last_scan_summary: (() => {
try { return JSON.parse(getState('last_scan_summary', '{}')) } catch { return {} }
})()
},
summary: {
total: summary.total || 0,
high_risk_count: summary.high_risk_count || 0,
calendar_ok_count: summary.calendar_ok_count || 0,
calendar_error_count: summary.calendar_error_count || 0,
weekly_error_count: summary.weekly_error_count || 0
}
}
}
const retryCalendarById = async id => {
const alert = getAlertByIdStmt.get(id)
if (!alert) throw new Error('alert_not_found')
return dispatchCalendar(alert, readConfig())
}
const retryWeeklyById = async id => {
const alert = getAlertByIdStmt.get(id)
if (!alert) throw new Error('alert_not_found')
return dispatchWeekly(alert, readConfig())
}
module.exports = {
start,
readConfig,
triggerScan,
requestScan,
listAlerts,
getAlert: id => getAlertByIdStmt.get(id),
getStatus,
retryCalendarById,
retryWeeklyById
}