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(//gi, ' ') .replace(//gi, ' ') .replace(/<[^>]+>/g, ' ') .replace(/ /gi, ' ') .replace(/&/gi, '&') .replace(/</gi, '<') .replace(/>/gi, '>') .replace(/"/gi, '"') .replace(/'/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 }