"""邮件模块:IMAP 检测模型下线通知 + 转发给群发列表 + 转发记录去重。 去重机制(本地不落盘,状态在已发送邮件的信头里): - 转发时在转发件上写入指向原邮件的 References 与 X-Forwarded-Msgid 信头 - 每轮先扫"已发送"近 N 天这两个信头,收集所有"已被转发过的原邮件 Message-ID" - 命中即说明这封下线通知已转发过,跳过 这是邮件协议原生的"转发关系"表达,标题可以保持干净,也不怕标题被改动。 """ import imaplib import logging import re import smtplib import ssl import time from datetime import datetime, timedelta from email import message_from_bytes from email.header import Header, decode_header, make_header from email.mime.message import MIMEMessage from email.mime.multipart import MIMEMultipart from email.mime.text import MIMEText from email.utils import formataddr, formatdate, parseaddr log = logging.getLogger("mailer") # 标记转发关系的自定义信头(与标准 References 双保险,防止服务商丢弃其一) FORWARD_HEADER = "X-Forwarded-Msgid" _ANGLE_MSGID_RE = re.compile(r"<([^<>\s]+)>") _FWD_HEADER_RE = re.compile( r"^%s:\s*(.+)$" % re.escape(FORWARD_HEADER), re.IGNORECASE | re.MULTILINE ) def normalize_message_id(message_id): """把 规整成 abc@host。""" return (message_id or "").strip().strip("<>").strip() def build_forward_subject(origin_subject): """转发件标题:干净的转发标识 + 原邮件标题(不再往标题里塞 Message-ID)。""" origin = (origin_subject or "").strip() if len(origin) > 120: origin = origin[:120] + "…" return f"【模型下线通知转发】{origin}" def normalize_domain(value): """把用户可能写成 "@aliyun.com" / "noreply@aliyun.com" / "ALIYUN.COM" 的配置统一成小写裸域名 "aliyun.com"。 """ d = (value or "").strip().lower().rstrip(".") if "@" in d: # 误填了完整邮箱地址,取 @ 后面的部分 d = d.rsplit("@", 1)[1] return d.lstrip("@").strip() def match_sender_domain(from_addr, domains): """发件人域名是否属于 domains 之一(本域或其子域)。 只接受"完全相等"或"以 .域名 结尾"两种情况,**不能用裸 endswith**: aliyun.com -> 命中 aliyun.com、mail.aliyun.com notaliyun.com -> 不命中(裸 endswith 会误放行,等于开后门) aliyun.com.evil.cn-> 不命中 """ addr = (from_addr or "").strip().lower() if "@" not in addr: return False host = addr.rsplit("@", 1)[1].strip().rstrip(".") if not host: return False for d in domains: d = normalize_domain(d) if d and (host == d or host.endswith("." + d)): return True return False def _imap_quote(value): """IMAP 字符串字面量转义(反斜杠和双引号)。""" return '"%s"' % str(value).replace("\\", "\\\\").replace('"', '\\"') def _imap_or_from(domains): """把多个域名拼成 IMAP 的 FROM 或条件。 IMAP 的 OR 只接受两个 key,多个要嵌套: 1 个 -> FROM "a" 2 个 -> OR FROM "a" FROM "b" 3 个 -> OR FROM "a" OR FROM "b" FROM "c" """ terms = ["FROM %s" % _imap_quote(normalize_domain(d)) for d in domains if normalize_domain(d)] if not terms: return "" expr = terms[-1] for term in reversed(terms[:-1]): expr = "OR %s %s" % (term, expr) return expr def _imap_and(*parts): """IMAP 搜索条件默认就是 AND,空条件跳过。""" return " ".join(p for p in parts if p) # ========== IMAP 检测 ========== class MailChecker: def __init__(self, mail_cfg, notice_cfg): self.mail = mail_cfg self.notice = notice_cfg def find_offline_notices(self): """返回需要转发的下线通知列表(已排除"已发送"中转发过的)。""" days = max(1, int(self.notice.get("lookback_days", 3))) since = (datetime.now() - timedelta(days=days)).strftime("%d-%b-%Y") conn = self._connect() try: forwarded = self._forwarded_message_ids(conn, since) log.info("[邮件] 已发送中近 %d 天的转发记录: %d 封原邮件", days, len(forwarded)) notices = self._scan_inbox(conn, since) log.info("[邮件] 收件箱近 %d 天命中下线通知: %d 封", days, len(notices)) fresh = [n for n in notices if not self._already_forwarded(n["message_id"], forwarded)] skipped = len(notices) - len(fresh) if skipped: log.info("[邮件] 其中 %d 封已转发过,跳过", skipped) return fresh finally: try: conn.logout() except Exception: # noqa: BLE001 pass # ----- 内部实现 ----- def _connect(self): host, port = self.mail["imap_host"], int(self.mail["imap_port"]) if port == 993: conn = imaplib.IMAP4_SSL(host, port) else: conn = imaplib.IMAP4(host, port) conn.starttls() conn.login(self.mail["username"], self.mail["password"]) return conn def _scan_inbox(self, conn, since): """扫收件箱:发件人域名属于 sender_domains 且标题/正文含 keyword。""" if conn.select(self.mail["inbox_folder"], readonly=True)[0] != "OK": raise RuntimeError(f"打开收件箱 {self.mail['inbox_folder']} 失败") domains = self.notice.get("sender_domains") or [] criteria = f"SINCE {since}" if domains: # IMAP 的 FROM 是子串匹配,先用它把范围缩小,减少 fetch 量; # 精确的域名归属校验在下面逐封做。 criteria = _imap_and(criteria, _imap_or_from(domains)) else: log.warning("[邮件] offline_notice.sender_domains 未配置,本轮不做发件人过滤") typ, data = conn.search(None, criteria) if typ != "OK" or not data or not data[0]: return [] keyword = self.notice.get("keyword") or "模型下线通知" notices = [] for num in data[0].split(): typ, msg_data = conn.fetch(num, "(RFC822)") if typ != "OK" or not msg_data or not msg_data[0]: continue raw = msg_data[0][1] msg = message_from_bytes(raw) from_addr = parseaddr(msg.get("From", ""))[1] # IMAP 的子串匹配不可信(evil@notaliyun.com 也会命中),必须精确校验域名归属 if domains and not match_sender_domain(from_addr, domains): continue # 告警邮件本身也会进自己的收件箱,必须排除,否则会自我触发 if from_addr.lower() == self.mail["username"].lower(): continue subject = _decode_mime(msg.get("Subject")) body = _body_text(msg) if keyword not in f"{subject}\n{body}": continue message_id = normalize_message_id(msg.get("Message-ID")) if not message_id: # 极少数邮件没有 Message-ID,用发件人+日期+标题兜底出一个稳定标记 message_id = "nomid-%s-%s" % (from_addr, abs(hash(subject + msg.get("Date", "")))) notices.append({ "message_id": message_id, "subject": subject, "from": from_addr, "date": msg.get("Date", ""), "body": body, "raw": raw, # 原始邮件字节,用于转发时附上完整原文 }) return notices def _forwarded_message_ids(self, conn, since): """扫"已发送"近 N 天的转发关系信头,返回已转发过的原邮件 Message-ID 集合。""" folder = self._select_sent(conn) if not folder: return set() typ, data = conn.search(None, f"SINCE {since}") if typ != "OK" or not data or not data[0]: return set() nums = data[0].split() # 一次 fetch 取回转发关系信头,避免逐封往返 typ, msg_data = conn.fetch( b",".join(nums), "(BODY.PEEK[HEADER.FIELDS (REFERENCES %s)])" % FORWARD_HEADER.upper(), ) if typ != "OK" or not msg_data: return set() ids = set() for item in msg_data: if not isinstance(item, tuple) or len(item) < 2 or not item[1]: continue ids |= _extract_msgids(item[1].decode("utf-8", errors="ignore")) return ids def _select_sent(self, conn): """依次尝试候选"已发送"文件夹名,返回成功选中的名字。""" for folder in _sent_candidates(self.mail.get("sent_folder")): try: if conn.select(_quote_folder(folder), readonly=True)[0] == "OK": return folder except Exception: # noqa: BLE001 - 中文文件夹名编码问题等,换下一个候选 pass log.info("[邮件] 已发送文件夹 %s 不可用,尝试下一个", folder) log.warning("[邮件] 未找到可用的已发送文件夹,本轮无法去重(可能重复转发)") return None @staticmethod def _already_forwarded(message_id, forwarded_ids): """该原邮件 Message-ID 出现在已发送的转发关系信头里 -> 已转发过。""" mid = normalize_message_id(message_id) return bool(mid) and mid in forwarded_ids def _sent_candidates(configured): names = [n for n in (configured, "Sent Messages", "Sent", "已发送", "INBOX.Sent") if n] seen, out = set(), [] for n in names: if n not in seen: seen.add(n) out.append(n) return out def _quote_folder(name): name = (name or "").strip() if name.startswith('"') and name.endswith('"'): return name if any(ch.isspace() or ch in '(){%*"\\' for ch in name): return '"%s"' % name return name def _extract_msgids(raw_headers): """从 References / X-Forwarded-Msgid 信头原文里抽出所有 Message-ID(已规整)。""" ids = {m.group(1).strip() for m in _ANGLE_MSGID_RE.finditer(raw_headers)} # 自定义信头可能是不带尖括号的裸值 for m in _FWD_HEADER_RE.finditer(raw_headers): v = normalize_message_id(m.group(1)) if v: ids.add(v) return {i for i in ids if i} def _decode_mime(value): """解码 MIME 编码标题(=?UTF-8?B?...?=)。""" if not value: return "" try: return str(make_header(decode_header(value))) except Exception: # noqa: BLE001 return value def _body_text(msg): """提取纯文本正文;没有 text/plain 时退化为 text/html 原文。""" plain, html = [], [] for part in (msg.walk() if msg.is_multipart() else [msg]): ctype = part.get_content_type() if ctype not in ("text/plain", "text/html"): continue payload = part.get_payload(decode=True) if not payload: continue text = payload.decode(part.get_content_charset() or "utf-8", errors="ignore") (plain if ctype == "text/plain" else html).append(text) return "\n".join(plain or html) # ========== SMTP 群发 ========== def send_mail(mail_cfg, subject, body): """向 notify_emails 群发一封纯文本邮件(用于模型可用性告警,无需去重)。""" _send(mail_cfg, MIMEText(body, "plain", "utf-8"), subject, action="群发告警") def forward_notice(mail_cfg, notice): """把检测到的下线通知原邮件转发给 notify_emails。 转发件结构: text/plain 转发说明 + 原邮件正文(可直接阅读) message/rfc822 原始邮件完整存档(.eml 附件,保留全部头信息) 去重关系写在信头里:References + X-Forwarded-Msgid 指向原邮件 Message-ID, 下轮扫"已发送"这两个信头即可判定该原邮件是否已转发过。 """ mid = normalize_message_id(notice["message_id"]) subject = build_forward_subject(notice["subject"]) intro = "\n".join([ "检测到模型下线通知,以下为原邮件转发。", "", f"原发件人 : {notice['from']}", f"原发送时间: {notice['date']}", f"原邮件标题: {notice['subject']}", f"Message-ID: {mid}", "", "-" * 56, "以下为原邮件正文:", "", notice.get("body") or "(原邮件无纯文本正文,请查看附件 original.eml)", "", "-" * 56, "本邮件由 model-watchdog 自动转发。", ]) msg = MIMEMultipart() msg.attach(MIMEText(intro, "plain", "utf-8")) if notice.get("raw"): try: att = MIMEMessage(message_from_bytes(notice["raw"])) att.add_header("Content-Disposition", "attachment", filename="original.eml") msg.attach(att) except Exception as e: # noqa: BLE001 - 附件失败不该阻断转发 log.warning("[邮件] 原邮件附件构造失败,仅转发正文: %s", e) if notice.get("from"): msg["Reply-To"] = notice["from"] # 转发关系(去重依据) msg["References"] = f"<{mid}>" msg[FORWARD_HEADER] = mid _send(mail_cfg, msg, subject, sent_marker=mid, action="转发下线通知") def _send(mail_cfg, msg, subject, sent_marker=None, action="发送"): """填充信头 -> SMTP 投递 -> 必要时确认已发送归档。""" recipients = mail_cfg["notify_emails"] if not recipients: raise ValueError("notify_emails 为空,无法群发") msg["Subject"] = Header(subject, "utf-8") msg["From"] = formataddr(("model-watchdog", mail_cfg["username"])) msg["To"] = ", ".join(recipients) msg["Date"] = formatdate(localtime=True) if mail_cfg.get("smtp_ssl"): server = smtplib.SMTP_SSL( mail_cfg["smtp_host"], int(mail_cfg["smtp_port"]), context=ssl.create_default_context(), timeout=30, ) else: server = smtplib.SMTP(mail_cfg["smtp_host"], int(mail_cfg["smtp_port"]), timeout=30) server.starttls(context=ssl.create_default_context()) try: server.login(mail_cfg["username"], mail_cfg["password"]) server.sendmail(mail_cfg["username"], recipients, msg.as_string()) log.info("[邮件] 已%s给 %d 人: %s", action, len(recipients), subject) finally: try: server.quit() except Exception: # noqa: BLE001 pass if sent_marker and mail_cfg.get("append_to_sent", True): _ensure_in_sent(mail_cfg, msg, sent_marker) if sent_marker and mail_cfg.get("append_to_sent", True): _ensure_in_sent(mail_cfg, msg, sent_marker) def _ensure_in_sent(mail_cfg, msg, marker, attempts=4, delay=4): """确认转发关系信头已出现在已发送里,否则 APPEND 补档。失败只告警不抛异常。 邮箱服务商自动归档存在延迟(腾讯企业邮实测 >2s),所以轮询重试若干次, 避免"其实已归档成功"却误报为可能重复转发。 """ checker = MailChecker(mail_cfg, {}) since = (datetime.now() - timedelta(days=1)).strftime("%d-%b-%Y") try: conn = checker._connect() except Exception as e: # noqa: BLE001 log.warning("[邮件] 无法确认已发送归档(IMAP 登录失败): %s", e) return try: for _ in range(max(1, attempts)): time.sleep(delay) if checker._already_forwarded(marker, checker._forwarded_message_ids(conn, since)): log.info("[邮件] 转发记录已确认在已发送中(服务商自动归档),去重生效") return folder = checker._select_sent(conn) if not folder: return typ, _ = conn.append( _quote_folder(folder), r"(\Seen)", imaplib.Time2Internaldate(datetime.now().timestamp()), msg.as_bytes(), ) if typ == "OK": log.info("[邮件] 已 APPEND 转发记录到 %s", folder) else: log.warning( "[邮件] %ds 内未确认归档,且 %s 不接受 APPEND(返回 %s)。" "若服务商稍后完成归档仍可去重,否则该条下线通知下轮会重复转发", attempts * delay, folder, typ, ) except Exception as e: # noqa: BLE001 log.warning("[邮件] 确认已发送归档失败: %s", e) finally: try: conn.logout() except Exception: # noqa: BLE001 pass