mailer.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434
  1. """邮件模块:IMAP 检测模型下线通知 + 转发给群发列表 + 转发记录去重。
  2. 去重机制(本地不落盘,状态在已发送邮件的信头里):
  3. - 转发时在转发件上写入指向原邮件的 References 与 X-Forwarded-Msgid 信头
  4. - 每轮先扫"已发送"近 N 天这两个信头,收集所有"已被转发过的原邮件 Message-ID"
  5. - 命中即说明这封下线通知已转发过,跳过
  6. 这是邮件协议原生的"转发关系"表达,标题可以保持干净,也不怕标题被改动。
  7. """
  8. import imaplib
  9. import logging
  10. import re
  11. import smtplib
  12. import ssl
  13. import time
  14. from datetime import datetime, timedelta
  15. from email import message_from_bytes
  16. from email.header import Header, decode_header, make_header
  17. from email.mime.message import MIMEMessage
  18. from email.mime.multipart import MIMEMultipart
  19. from email.mime.text import MIMEText
  20. from email.utils import formataddr, formatdate, parseaddr
  21. log = logging.getLogger("mailer")
  22. # 标记转发关系的自定义信头(与标准 References 双保险,防止服务商丢弃其一)
  23. FORWARD_HEADER = "X-Forwarded-Msgid"
  24. _ANGLE_MSGID_RE = re.compile(r"<([^<>\s]+)>")
  25. _FWD_HEADER_RE = re.compile(
  26. r"^%s:\s*(.+)$" % re.escape(FORWARD_HEADER), re.IGNORECASE | re.MULTILINE
  27. )
  28. def normalize_message_id(message_id):
  29. """把 <abc@host> 规整成 abc@host。"""
  30. return (message_id or "").strip().strip("<>").strip()
  31. def build_forward_subject(origin_subject):
  32. """转发件标题:干净的转发标识 + 原邮件标题(不再往标题里塞 Message-ID)。"""
  33. origin = (origin_subject or "").strip()
  34. if len(origin) > 120:
  35. origin = origin[:120] + "…"
  36. return f"【模型下线通知转发】{origin}"
  37. def normalize_domain(value):
  38. """把用户可能写成 "@aliyun.com" / "noreply@aliyun.com" / "ALIYUN.COM"
  39. 的配置统一成小写裸域名 "aliyun.com"。
  40. """
  41. d = (value or "").strip().lower().rstrip(".")
  42. if "@" in d: # 误填了完整邮箱地址,取 @ 后面的部分
  43. d = d.rsplit("@", 1)[1]
  44. return d.lstrip("@").strip()
  45. def match_sender_domain(from_addr, domains):
  46. """发件人域名是否属于 domains 之一(本域或其子域)。
  47. 只接受"完全相等"或"以 .域名 结尾"两种情况,**不能用裸 endswith**:
  48. aliyun.com -> 命中 aliyun.com、mail.aliyun.com
  49. notaliyun.com -> 不命中(裸 endswith 会误放行,等于开后门)
  50. aliyun.com.evil.cn-> 不命中
  51. """
  52. addr = (from_addr or "").strip().lower()
  53. if "@" not in addr:
  54. return False
  55. host = addr.rsplit("@", 1)[1].strip().rstrip(".")
  56. if not host:
  57. return False
  58. for d in domains:
  59. d = normalize_domain(d)
  60. if d and (host == d or host.endswith("." + d)):
  61. return True
  62. return False
  63. def _imap_quote(value):
  64. """IMAP 字符串字面量转义(反斜杠和双引号)。"""
  65. return '"%s"' % str(value).replace("\\", "\\\\").replace('"', '\\"')
  66. def _imap_or_from(domains):
  67. """把多个域名拼成 IMAP 的 FROM 或条件。
  68. IMAP 的 OR 只接受两个 key,多个要嵌套:
  69. 1 个 -> FROM "a"
  70. 2 个 -> OR FROM "a" FROM "b"
  71. 3 个 -> OR FROM "a" OR FROM "b" FROM "c"
  72. """
  73. terms = ["FROM %s" % _imap_quote(normalize_domain(d)) for d in domains if normalize_domain(d)]
  74. if not terms:
  75. return ""
  76. expr = terms[-1]
  77. for term in reversed(terms[:-1]):
  78. expr = "OR %s %s" % (term, expr)
  79. return expr
  80. def _imap_and(*parts):
  81. """IMAP 搜索条件默认就是 AND,空条件跳过。"""
  82. return " ".join(p for p in parts if p)
  83. # ========== IMAP 检测 ==========
  84. class MailChecker:
  85. def __init__(self, mail_cfg, notice_cfg):
  86. self.mail = mail_cfg
  87. self.notice = notice_cfg
  88. def find_offline_notices(self):
  89. """返回需要转发的下线通知列表(已排除"已发送"中转发过的)。"""
  90. days = max(1, int(self.notice.get("lookback_days", 3)))
  91. since = (datetime.now() - timedelta(days=days)).strftime("%d-%b-%Y")
  92. conn = self._connect()
  93. try:
  94. forwarded = self._forwarded_message_ids(conn, since)
  95. log.info("[邮件] 已发送中近 %d 天的转发记录: %d 封原邮件", days, len(forwarded))
  96. notices = self._scan_inbox(conn, since)
  97. log.info("[邮件] 收件箱近 %d 天命中下线通知: %d 封", days, len(notices))
  98. fresh = [n for n in notices if not self._already_forwarded(n["message_id"], forwarded)]
  99. skipped = len(notices) - len(fresh)
  100. if skipped:
  101. log.info("[邮件] 其中 %d 封已转发过,跳过", skipped)
  102. return fresh
  103. finally:
  104. try:
  105. conn.logout()
  106. except Exception: # noqa: BLE001
  107. pass
  108. # ----- 内部实现 -----
  109. def _connect(self):
  110. host, port = self.mail["imap_host"], int(self.mail["imap_port"])
  111. if port == 993:
  112. conn = imaplib.IMAP4_SSL(host, port)
  113. else:
  114. conn = imaplib.IMAP4(host, port)
  115. conn.starttls()
  116. conn.login(self.mail["username"], self.mail["password"])
  117. return conn
  118. def _scan_inbox(self, conn, since):
  119. """扫收件箱:发件人域名属于 sender_domains 且标题/正文含 keyword。"""
  120. if conn.select(self.mail["inbox_folder"], readonly=True)[0] != "OK":
  121. raise RuntimeError(f"打开收件箱 {self.mail['inbox_folder']} 失败")
  122. domains = self.notice.get("sender_domains") or []
  123. criteria = f"SINCE {since}"
  124. if domains:
  125. # IMAP 的 FROM 是子串匹配,先用它把范围缩小,减少 fetch 量;
  126. # 精确的域名归属校验在下面逐封做。
  127. criteria = _imap_and(criteria, _imap_or_from(domains))
  128. else:
  129. log.warning("[邮件] offline_notice.sender_domains 未配置,本轮不做发件人过滤")
  130. typ, data = conn.search(None, criteria)
  131. if typ != "OK" or not data or not data[0]:
  132. return []
  133. keyword = self.notice.get("keyword") or "模型下线通知"
  134. notices = []
  135. for num in data[0].split():
  136. typ, msg_data = conn.fetch(num, "(RFC822)")
  137. if typ != "OK" or not msg_data or not msg_data[0]:
  138. continue
  139. raw = msg_data[0][1]
  140. msg = message_from_bytes(raw)
  141. from_addr = parseaddr(msg.get("From", ""))[1]
  142. # IMAP 的子串匹配不可信(evil@notaliyun.com 也会命中),必须精确校验域名归属
  143. if domains and not match_sender_domain(from_addr, domains):
  144. continue
  145. # 告警邮件本身也会进自己的收件箱,必须排除,否则会自我触发
  146. if from_addr.lower() == self.mail["username"].lower():
  147. continue
  148. subject = _decode_mime(msg.get("Subject"))
  149. body = _body_text(msg)
  150. if keyword not in f"{subject}\n{body}":
  151. continue
  152. message_id = normalize_message_id(msg.get("Message-ID"))
  153. if not message_id:
  154. # 极少数邮件没有 Message-ID,用发件人+日期+标题兜底出一个稳定标记
  155. message_id = "nomid-%s-%s" % (from_addr, abs(hash(subject + msg.get("Date", ""))))
  156. notices.append({
  157. "message_id": message_id,
  158. "subject": subject,
  159. "from": from_addr,
  160. "date": msg.get("Date", ""),
  161. "body": body,
  162. "raw": raw, # 原始邮件字节,用于转发时附上完整原文
  163. })
  164. return notices
  165. def _forwarded_message_ids(self, conn, since):
  166. """扫"已发送"近 N 天的转发关系信头,返回已转发过的原邮件 Message-ID 集合。"""
  167. folder = self._select_sent(conn)
  168. if not folder:
  169. return set()
  170. typ, data = conn.search(None, f"SINCE {since}")
  171. if typ != "OK" or not data or not data[0]:
  172. return set()
  173. nums = data[0].split()
  174. # 一次 fetch 取回转发关系信头,避免逐封往返
  175. typ, msg_data = conn.fetch(
  176. b",".join(nums),
  177. "(BODY.PEEK[HEADER.FIELDS (REFERENCES %s)])" % FORWARD_HEADER.upper(),
  178. )
  179. if typ != "OK" or not msg_data:
  180. return set()
  181. ids = set()
  182. for item in msg_data:
  183. if not isinstance(item, tuple) or len(item) < 2 or not item[1]:
  184. continue
  185. ids |= _extract_msgids(item[1].decode("utf-8", errors="ignore"))
  186. return ids
  187. def _select_sent(self, conn):
  188. """依次尝试候选"已发送"文件夹名,返回成功选中的名字。"""
  189. for folder in _sent_candidates(self.mail.get("sent_folder")):
  190. try:
  191. if conn.select(_quote_folder(folder), readonly=True)[0] == "OK":
  192. return folder
  193. except Exception: # noqa: BLE001 - 中文文件夹名编码问题等,换下一个候选
  194. pass
  195. log.info("[邮件] 已发送文件夹 %s 不可用,尝试下一个", folder)
  196. log.warning("[邮件] 未找到可用的已发送文件夹,本轮无法去重(可能重复转发)")
  197. return None
  198. @staticmethod
  199. def _already_forwarded(message_id, forwarded_ids):
  200. """该原邮件 Message-ID 出现在已发送的转发关系信头里 -> 已转发过。"""
  201. mid = normalize_message_id(message_id)
  202. return bool(mid) and mid in forwarded_ids
  203. def _sent_candidates(configured):
  204. names = [n for n in (configured, "Sent Messages", "Sent", "已发送", "INBOX.Sent") if n]
  205. seen, out = set(), []
  206. for n in names:
  207. if n not in seen:
  208. seen.add(n)
  209. out.append(n)
  210. return out
  211. def _quote_folder(name):
  212. name = (name or "").strip()
  213. if name.startswith('"') and name.endswith('"'):
  214. return name
  215. if any(ch.isspace() or ch in '(){%*"\\' for ch in name):
  216. return '"%s"' % name
  217. return name
  218. def _extract_msgids(raw_headers):
  219. """从 References / X-Forwarded-Msgid 信头原文里抽出所有 Message-ID(已规整)。"""
  220. ids = {m.group(1).strip() for m in _ANGLE_MSGID_RE.finditer(raw_headers)}
  221. # 自定义信头可能是不带尖括号的裸值
  222. for m in _FWD_HEADER_RE.finditer(raw_headers):
  223. v = normalize_message_id(m.group(1))
  224. if v:
  225. ids.add(v)
  226. return {i for i in ids if i}
  227. def _decode_mime(value):
  228. """解码 MIME 编码标题(=?UTF-8?B?...?=)。"""
  229. if not value:
  230. return ""
  231. try:
  232. return str(make_header(decode_header(value)))
  233. except Exception: # noqa: BLE001
  234. return value
  235. def _body_text(msg):
  236. """提取纯文本正文;没有 text/plain 时退化为 text/html 原文。"""
  237. plain, html = [], []
  238. for part in (msg.walk() if msg.is_multipart() else [msg]):
  239. ctype = part.get_content_type()
  240. if ctype not in ("text/plain", "text/html"):
  241. continue
  242. payload = part.get_payload(decode=True)
  243. if not payload:
  244. continue
  245. text = payload.decode(part.get_content_charset() or "utf-8", errors="ignore")
  246. (plain if ctype == "text/plain" else html).append(text)
  247. return "\n".join(plain or html)
  248. # ========== SMTP 群发 ==========
  249. def send_mail(mail_cfg, subject, body):
  250. """向 notify_emails 群发一封纯文本邮件(用于模型可用性告警,无需去重)。"""
  251. _send(mail_cfg, MIMEText(body, "plain", "utf-8"), subject, action="群发告警")
  252. def forward_notice(mail_cfg, notice):
  253. """把检测到的下线通知原邮件转发给 notify_emails。
  254. 转发件结构:
  255. text/plain 转发说明 + 原邮件正文(可直接阅读)
  256. message/rfc822 原始邮件完整存档(.eml 附件,保留全部头信息)
  257. 去重关系写在信头里:References + X-Forwarded-Msgid 指向原邮件 Message-ID,
  258. 下轮扫"已发送"这两个信头即可判定该原邮件是否已转发过。
  259. """
  260. mid = normalize_message_id(notice["message_id"])
  261. subject = build_forward_subject(notice["subject"])
  262. intro = "\n".join([
  263. "检测到模型下线通知,以下为原邮件转发。",
  264. "",
  265. f"原发件人 : {notice['from']}",
  266. f"原发送时间: {notice['date']}",
  267. f"原邮件标题: {notice['subject']}",
  268. f"Message-ID: {mid}",
  269. "",
  270. "-" * 56,
  271. "以下为原邮件正文:",
  272. "",
  273. notice.get("body") or "(原邮件无纯文本正文,请查看附件 original.eml)",
  274. "",
  275. "-" * 56,
  276. "本邮件由 model-watchdog 自动转发。",
  277. ])
  278. msg = MIMEMultipart()
  279. msg.attach(MIMEText(intro, "plain", "utf-8"))
  280. if notice.get("raw"):
  281. try:
  282. att = MIMEMessage(message_from_bytes(notice["raw"]))
  283. att.add_header("Content-Disposition", "attachment", filename="original.eml")
  284. msg.attach(att)
  285. except Exception as e: # noqa: BLE001 - 附件失败不该阻断转发
  286. log.warning("[邮件] 原邮件附件构造失败,仅转发正文: %s", e)
  287. if notice.get("from"):
  288. msg["Reply-To"] = notice["from"]
  289. # 转发关系(去重依据)
  290. msg["References"] = f"<{mid}>"
  291. msg[FORWARD_HEADER] = mid
  292. _send(mail_cfg, msg, subject, sent_marker=mid, action="转发下线通知")
  293. def _send(mail_cfg, msg, subject, sent_marker=None, action="发送"):
  294. """填充信头 -> SMTP 投递 -> 必要时确认已发送归档。"""
  295. recipients = mail_cfg["notify_emails"]
  296. if not recipients:
  297. raise ValueError("notify_emails 为空,无法群发")
  298. msg["Subject"] = Header(subject, "utf-8")
  299. msg["From"] = formataddr(("model-watchdog", mail_cfg["username"]))
  300. msg["To"] = ", ".join(recipients)
  301. msg["Date"] = formatdate(localtime=True)
  302. if mail_cfg.get("smtp_ssl"):
  303. server = smtplib.SMTP_SSL(
  304. mail_cfg["smtp_host"], int(mail_cfg["smtp_port"]),
  305. context=ssl.create_default_context(), timeout=30,
  306. )
  307. else:
  308. server = smtplib.SMTP(mail_cfg["smtp_host"], int(mail_cfg["smtp_port"]), timeout=30)
  309. server.starttls(context=ssl.create_default_context())
  310. try:
  311. server.login(mail_cfg["username"], mail_cfg["password"])
  312. server.sendmail(mail_cfg["username"], recipients, msg.as_string())
  313. log.info("[邮件] 已%s给 %d 人: %s", action, len(recipients), subject)
  314. finally:
  315. try:
  316. server.quit()
  317. except Exception: # noqa: BLE001
  318. pass
  319. if sent_marker and mail_cfg.get("append_to_sent", True):
  320. _ensure_in_sent(mail_cfg, msg, sent_marker)
  321. if sent_marker and mail_cfg.get("append_to_sent", True):
  322. _ensure_in_sent(mail_cfg, msg, sent_marker)
  323. def _ensure_in_sent(mail_cfg, msg, marker, attempts=4, delay=4):
  324. """确认转发关系信头已出现在已发送里,否则 APPEND 补档。失败只告警不抛异常。
  325. 邮箱服务商自动归档存在延迟(腾讯企业邮实测 >2s),所以轮询重试若干次,
  326. 避免"其实已归档成功"却误报为可能重复转发。
  327. """
  328. checker = MailChecker(mail_cfg, {})
  329. since = (datetime.now() - timedelta(days=1)).strftime("%d-%b-%Y")
  330. try:
  331. conn = checker._connect()
  332. except Exception as e: # noqa: BLE001
  333. log.warning("[邮件] 无法确认已发送归档(IMAP 登录失败): %s", e)
  334. return
  335. try:
  336. for _ in range(max(1, attempts)):
  337. time.sleep(delay)
  338. if checker._already_forwarded(marker, checker._forwarded_message_ids(conn, since)):
  339. log.info("[邮件] 转发记录已确认在已发送中(服务商自动归档),去重生效")
  340. return
  341. folder = checker._select_sent(conn)
  342. if not folder:
  343. return
  344. typ, _ = conn.append(
  345. _quote_folder(folder), r"(\Seen)",
  346. imaplib.Time2Internaldate(datetime.now().timestamp()),
  347. msg.as_bytes(),
  348. )
  349. if typ == "OK":
  350. log.info("[邮件] 已 APPEND 转发记录到 %s", folder)
  351. else:
  352. log.warning(
  353. "[邮件] %ds 内未确认归档,且 %s 不接受 APPEND(返回 %s)。"
  354. "若服务商稍后完成归档仍可去重,否则该条下线通知下轮会重复转发",
  355. attempts * delay, folder, typ,
  356. )
  357. except Exception as e: # noqa: BLE001
  358. log.warning("[邮件] 确认已发送归档失败: %s", e)
  359. finally:
  360. try:
  361. conn.logout()
  362. except Exception: # noqa: BLE001
  363. pass