chore: 初始提交 - 移除所有硬编码凭据,统一从环境变量读取

This commit is contained in:
yangxiangyuan
2026-06-11 17:10:58 +08:00
commit c223c79b36
1699 changed files with 313940 additions and 0 deletions
+974
View File
@@ -0,0 +1,974 @@
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
}