| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434 |
- """邮件模块: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> 规整成 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
|