from __future__ import annotations import hashlib import os import random import re import string import time from datetime import datetime, timezone from email import message_from_string, policy from email.utils import parsedate_to_datetime from threading import Lock from typing import Any, Callable, TypeVar import requests from curl_cffi import requests as curl_requests from services.register import domain_reputation, pool_mail ResultT = TypeVar("ResultT") domain_lock = Lock() provider_lock = Lock() domain_index = 0 provider_index = 0 def _config(mail_config: dict) -> dict: return { "request_timeout": float(mail_config.get("request_timeout") or 30), "wait_timeout": float(mail_config.get("wait_timeout") or 30), "wait_interval": float(mail_config.get("wait_interval") or 2), "user_agent": str(mail_config.get("user_agent") or "Mozilla/5.0"), } def _random_mailbox_name() -> str: return f"{''.join(random.choices(string.ascii_lowercase, k=5))}{''.join(random.choices(string.digits, k=random.randint(1, 3)))}{''.join(random.choices(string.ascii_lowercase, k=random.randint(1, 3)))}" def _random_subdomain_label() -> str: return "".join(random.choices(string.ascii_lowercase + string.digits, k=random.randint(4, 10))) def _next_domain(domains: list[str]) -> str: global domain_index domains = [str(item).strip() for item in domains if str(item).strip()] if not domains: raise RuntimeError("mail.domain 不能为空") if len(domains) == 1: return domains[0] with domain_lock: value = domains[domain_index % len(domains)] domain_index = (domain_index + 1) % len(domains) return value def _parse_received_at(value: Any) -> datetime | None: if isinstance(value, (int, float)): try: return datetime.fromtimestamp(float(value), tz=timezone.utc) except Exception: return None text = str(value or "").strip() if not text: return None try: date = datetime.fromisoformat(text[:-1] + "+00:00" if text.endswith("Z") else text) return date if date.tzinfo else date.replace(tzinfo=timezone.utc) except Exception: pass try: date = parsedate_to_datetime(text) return date if date.tzinfo else date.replace(tzinfo=timezone.utc) except Exception: return None def _extract_content(data: dict[str, Any]) -> tuple[str, str]: text_content = str(data.get("text_content") or data.get("text") or data.get("body") or data.get("content") or "") html_content = str(data.get("html_content") or data.get("html") or data.get("html_body") or data.get("body_html") or "") if text_content or html_content: return text_content, html_content raw = data.get("raw") if not isinstance(raw, str) or not raw.strip(): return "", "" try: parsed = message_from_string(raw, policy=policy.default) except Exception: return raw, "" plain: list[str] = [] html: list[str] = [] for part in parsed.walk() if parsed.is_multipart() else [parsed]: if part.get_content_maintype() == "multipart": continue try: payload = part.get_content() except Exception: payload = "" if not payload: continue if part.get_content_type() == "text/html": html.append(str(payload)) else: plain.append(str(payload)) return "\n".join(plain).strip(), "\n".join(html).strip() def _extract_text_candidates(value: Any) -> list[str]: if isinstance(value, str): return [value] if isinstance(value, dict): out: list[str] = [] for key in ("address", "email", "name", "value"): if value.get(key): out.extend(_extract_text_candidates(value.get(key))) return out if isinstance(value, list): out: list[str] = [] for item in value: out.extend(_extract_text_candidates(item)) return out return [] def _message_matches_email(data: dict[str, Any], email: str) -> bool: target = str(email or "").strip().lower() candidates: list[str] = [] for key in ("to", "mailTo", "receiver", "receivers", "address", "email", "envelope_to"): if key in data: candidates.extend(_extract_text_candidates(data.get(key))) return not target or not candidates or any(target in str(item).strip().lower() for item in candidates if str(item).strip()) def _extract_code(message: dict[str, Any]) -> str | None: content = f"{message.get('subject', '')}\n{message.get('text_content', '')}\n{message.get('html_content', '')}".strip() if not content: return None match = re.search(r"background-color:\s*#F3F3F3[^>]*>[\s\S]*?(\d{6})[\s\S]*?

", content, re.I) if match: return match.group(1) match = re.search(r"(?:Verification code|code is|代码为|验证码)[:\s]*(\d{6})", content, re.I) if match and match.group(1) != "177010": return match.group(1) for code in re.findall(r">\s*(\d{6})\s*<|(? str: provider = str(message.get("provider") or "").strip() mailbox = str(message.get("mailbox") or "").strip() message_id = str(message.get("message_id") or "").strip() if message_id: return f"id:{provider}:{mailbox}:{message_id}" received_at = message.get("received_at") received_value = received_at.isoformat() if isinstance(received_at, datetime) else str(received_at or "") content = "\n".join(str(message.get(key) or "") for key in ("subject", "sender", "text_content", "html_content")) digest = hashlib.sha256(content.encode("utf-8", errors="replace")).hexdigest() return f"content:{provider}:{mailbox}:{received_value}:{digest}" class BaseMailProvider: name = "unknown" def __init__(self, conf: dict, provider_ref: str = ""): self.conf = conf self.provider_ref = provider_ref def wait_for(self, mailbox: dict[str, Any], on_message: Callable[[dict[str, Any]], ResultT | None]) -> ResultT | None: deadline = time.monotonic() + self.conf["wait_timeout"] while time.monotonic() < deadline: message = self.fetch_latest_message(mailbox) if message: result = on_message(message) if result is not None: return result time.sleep(max(0.2, self.conf["wait_interval"])) return None def wait_for_code(self, mailbox: dict[str, Any]) -> str | None: seen_value = mailbox.setdefault("_seen_code_message_refs", []) if not isinstance(seen_value, list): seen_value = [] mailbox["_seen_code_message_refs"] = seen_value seen_refs = {str(item) for item in seen_value} def extract_unseen_code(message: dict[str, Any]) -> str | None: ref = _message_tracking_ref(message) if ref in seen_refs: return None code = _extract_code(message) if code: seen_value.append(ref) seen_refs.add(ref) return code return self.wait_for(mailbox, extract_unseen_code) def close(self) -> None: pass class CloudflareTempMailProvider(BaseMailProvider): name = "cloudflare_temp_email" retry_statuses = {429, 500, 502, 503, 504} def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_base = str(entry["api_base"]).rstrip("/") self.admin_password = str(entry["admin_password"]).strip() self.domain = entry.get("domain") or [] self.session = requests.Session() def _request(self, method: str, path: str, headers: dict | None = None, params: dict | None = None, payload: dict | None = None, expected: tuple[int, ...] = (200,)): last_error = "" for attempt in range(3): try: resp = self.session.request(method.upper(), f"{self.api_base}{path}", headers={"Content-Type": "application/json", "User-Agent": self.conf["user_agent"], **(headers or {})}, params=params, json=payload, timeout=self.conf["request_timeout"], verify=False) except Exception as exc: last_error = str(exc) if attempt < 2: time.sleep(0.5 * (attempt + 1)) continue raise RuntimeError(f"CloudflareTempMail 请求异常: {method} {path}, error={last_error}") from exc if resp.status_code in expected: return {} if resp.status_code == 204 else resp.json() if resp.status_code in self.retry_statuses and attempt < 2: time.sleep(0.5 * (attempt + 1)) continue raise RuntimeError(f"CloudflareTempMail 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") raise RuntimeError(f"CloudflareTempMail 请求异常: {method} {path}, error={last_error}") def create_mailbox(self, username: str | None = None) -> dict[str, Any]: data = self._request("POST", "/admin/new_address", headers={"x-admin-auth": self.admin_password}, payload={"enablePrefix": True, "name": username or _random_mailbox_name(), "domain": _next_domain(self.domain)}) address = str(data.get("address") or "").strip() token = str(data.get("jwt") or "").strip() if not address or not token: raise RuntimeError("CloudflareTempMail 缺少 address 或 jwt") return {"provider": self.name, "provider_ref": self.provider_ref, "address": address, "token": token} def close(self) -> None: self.session.close() def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: data = self._request("GET", "/api/mails", headers={"Authorization": f"Bearer {mailbox['token']}"}, params={"limit": 10, "offset": 0}) raw = list(data.get("results") or []) if isinstance(data, dict) else data if isinstance(data, list) else [] messages = [item for item in raw if isinstance(item, dict) and _message_matches_email(item, str(mailbox.get("address") or ""))] if not messages: return None item = messages[0] text_content, html_content = _extract_content(item) sender = item.get("from") or item.get("sender") or "" if isinstance(sender, dict): sender = sender.get("address") or sender.get("email") or sender.get("name") or "" return {"provider": self.name, "mailbox": mailbox["address"], "message_id": str(item.get("id") or item.get("_id") or ""), "subject": str(item.get("subject") or ""), "sender": str(sender), "text_content": text_content, "html_content": html_content, "received_at": _parse_received_at(item.get("createdAt") or item.get("created_at") or item.get("receivedAt") or item.get("date") or item.get("timestamp")), "raw": item} def close(self) -> None: self.session.close() class TempMailLolProvider(BaseMailProvider): name = "tempmail_lol" def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_key = str(entry.get("api_key") or "").strip() self.domain = [str(item).strip() for item in (entry.get("domain") or []) if str(item).strip()] self.session = requests.Session() self.session.trust_env = False self.session.headers.update({"User-Agent": conf["user_agent"], "Accept": "application/json", "Content-Type": "application/json"}) if self.api_key: self.session.headers["Authorization"] = f"Bearer {self.api_key}" @staticmethod def _resolve_domain(domain: str) -> tuple[str, bool]: text = str(domain or "").strip().lower() if text.startswith("*.") and len(text) > 2: return f"{_random_subdomain_label()}.{text[2:]}", True return text, False def _request(self, method: str, path: str, params: dict | None = None, payload: dict | None = None, expected: tuple[int, ...] = (200,)): resp = self.session.request(method.upper(), f"https://api.tempmail.lol/v2{path}", params=params, json=payload, timeout=self.conf["request_timeout"], verify=False) if resp.status_code not in expected: raise RuntimeError(f"TempMail.lol 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") data = resp.json() if not isinstance(data, dict): raise RuntimeError(f"TempMail.lol {method} {path} 返回结构不是对象") return data def create_mailbox(self, username: str | None = None) -> dict[str, Any]: payload: dict[str, Any] = {} if self.domain: domain, force_random_prefix = self._resolve_domain(random.choice(self.domain)) payload["domain"] = domain if force_random_prefix: payload["prefix"] = _random_mailbox_name() if username and "prefix" not in payload: payload["prefix"] = username data = self._request("POST", "/inbox/create", payload=payload, expected=(200, 201)) address = str(data.get("address") or "").strip() token = str(data.get("token") or "").strip() if not address or not token: raise RuntimeError("TempMail.lol 缺少 address 或 token") return {"provider": self.name, "provider_ref": self.provider_ref, "address": address, "token": token} def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: data = self._request("GET", "/inbox", params={"token": mailbox["token"]}) items = data.get("emails") or data.get("messages") or [] messages = [item for item in items if isinstance(item, dict)] if isinstance(items, list) else [] if not messages: return None item = max(messages, key=lambda value: ((_parse_received_at(value.get("created_at") or value.get("createdAt") or value.get("date") or value.get("received_at") or value.get("timestamp")) or datetime.fromtimestamp(0, tz=timezone.utc)).timestamp(), str(value.get("id") or value.get("token") or ""))) text_content, html_content = _extract_content(item) return {"provider": self.name, "mailbox": mailbox["address"], "message_id": str(item.get("id") or item.get("token") or ""), "subject": str(item.get("subject") or ""), "sender": str(item.get("from") or item.get("from_address") or ""), "text_content": text_content, "html_content": html_content, "received_at": _parse_received_at(item.get("created_at") or item.get("createdAt") or item.get("date") or item.get("received_at") or item.get("timestamp")), "raw": item} def close(self) -> None: self.session.close() class DuckMailProvider(BaseMailProvider): name = "duckmail" def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_key = str(entry["api_key"]).strip() self.default_domain = str(entry.get("default_domain") or "duckmail.sbs").strip() or "duckmail.sbs" self.session = requests.Session() self.session.trust_env = False self.session.headers.update({"User-Agent": conf["user_agent"], "Accept": "application/json", "Content-Type": "application/json"}) def _request(self, method: str, path: str, token: str = "", use_api_key: bool = False, params: dict | None = None, payload: dict | None = None, expected: tuple[int, ...] = (200, 201, 204)): headers = {"Authorization": f"Bearer {self.api_key if use_api_key else token}"} if use_api_key or token else {} resp = self.session.request(method.upper(), f"https://api.duckmail.sbs{path}", headers=headers, params=params, json=payload, timeout=self.conf["request_timeout"], verify=False) if resp.status_code not in expected: raise RuntimeError(f"DuckMail 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") return {} if resp.status_code == 204 else resp.json() @staticmethod def _items(data): return data if isinstance(data, list) else data.get("hydra:member") or data.get("member") or data.get("data") or [] def create_mailbox(self, username: str | None = None) -> dict[str, Any]: password = "".join(random.choices(string.ascii_letters + string.digits, k=12)) address = f"{username or _random_mailbox_name()}@{self.default_domain}" payload = {"address": address, "password": password} account = self._request("POST", "/accounts", use_api_key=True, payload=payload) token_data = self._request("POST", "/token", use_api_key=True, payload=payload) return {"provider": self.name, "provider_ref": self.provider_ref, "address": address, "token": str(token_data.get("token") or ""), "password": password, "account_id": str(account.get("id") or "")} def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: data = self._request("GET", "/messages", token=str(mailbox.get("token") or ""), params={"page": 1}) items = self._items(data) if not items: return None item = items[0] message_id = str(item.get("id") or item.get("@id") or "").replace("/messages/", "") if message_id: item = self._request("GET", f"/messages/{message_id}", token=str(mailbox.get("token") or "")) sender = item.get("from") or "" if isinstance(sender, dict): sender = sender.get("address") or sender.get("name") or "" html_content = item.get("html") or "" if isinstance(html_content, list): html_content = "".join(str(value) for value in html_content) return {"provider": self.name, "mailbox": mailbox["address"], "message_id": message_id, "subject": str(item.get("subject") or ""), "sender": str(sender), "text_content": str(item.get("text") or item.get("text_content") or ""), "html_content": str(html_content), "received_at": _parse_received_at(item.get("createdAt") or item.get("created_at") or item.get("receivedAt") or item.get("date")), "raw": item} def close(self) -> None: self.session.close() class GptMailProvider(BaseMailProvider): name = "gptmail" def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_key = str(entry["api_key"]).strip() self.default_domain = str(entry.get("default_domain") or "").strip() self.session = requests.Session() self.session.trust_env = False self.session.headers.update({"User-Agent": conf["user_agent"], "Accept": "application/json", "Content-Type": "application/json", "X-API-Key": self.api_key}) def _request(self, method: str, path: str, params: dict | None = None, payload: dict | None = None): query = dict(params or {}) resp = self.session.request(method.upper(), f"https://mail.chatgpt.org.uk{path}", params=query, json=payload, timeout=self.conf["request_timeout"], verify=False) if resp.status_code != 200: raise RuntimeError(f"GPTMail 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") data = resp.json() return data["data"] if isinstance(data, dict) and "data" in data else data def create_mailbox(self, username: str | None = None) -> dict[str, Any]: payload = {key: value for key, value in {"prefix": username, "domain": self.default_domain}.items() if value} data = self._request("POST" if payload else "GET", "/api/generate-email", payload=payload or None) return {"provider": self.name, "provider_ref": self.provider_ref, "address": str(data["email"])} def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: data = self._request("GET", "/api/emails", params={"email": mailbox["address"]}) emails = data if isinstance(data, list) else data.get("emails") or [] if not emails: return None item = max(emails, key=lambda value: (float(value.get("timestamp") or 0), str(value.get("id") or ""))) if item.get("id"): item = self._request("GET", f"/api/email/{item['id']}") return {"provider": self.name, "mailbox": mailbox["address"], "message_id": str(item.get("id") or ""), "subject": str(item.get("subject") or ""), "sender": str(item.get("from_address") or ""), "text_content": str(item.get("content") or ""), "html_content": str(item.get("html_content") or ""), "received_at": _parse_received_at(item.get("timestamp") or item.get("created_at")), "raw": item} def close(self) -> None: self.session.close() class MoEmailProvider(BaseMailProvider): name = "moemail" def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_base = str(entry["api_base"]).rstrip("/") self.api_key = str(entry["api_key"]).strip() raw_domains = entry.get("domain") or [] if isinstance(raw_domains, list): self.domain = [str(item).strip() for item in raw_domains if str(item).strip()] else: self.domain = [str(raw_domains).strip()] if str(raw_domains).strip() else [] self.expiry_time = int(entry.get("expiry_time") or 0) self.session = curl_requests.Session(impersonate="chrome") def _request(self, method: str, path: str, params: dict | None = None, payload: dict | None = None, expected: tuple[int, ...] = (200,)): resp = self.session.request(method.upper(), f"{self.api_base}{path}", headers={"X-API-Key": self.api_key, "Content-Type": "application/json", "User-Agent": self.conf["user_agent"]}, params=params, json=payload, timeout=self.conf["request_timeout"], verify=False) if resp.status_code not in expected: raise RuntimeError(f"MoEmail 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") data = resp.json() if not isinstance(data, dict): raise RuntimeError(f"MoEmail {method} {path} 返回结构不是对象") return data def create_mailbox(self, username: str | None = None) -> dict[str, Any]: data = self._request("POST", "/api/emails/generate", payload={"name": username or _random_mailbox_name(), "expiryTime": self.expiry_time, "domain": _next_domain(self.domain)}, expected=(200, 201)) address = str(data.get("email") or "").strip() email_id = str(data.get("id") or data.get("email_id") or "").strip() if not address or not email_id: raise RuntimeError("MoEmail 缺少 email 或 id") return {"provider": self.name, "provider_ref": self.provider_ref, "address": address, "email_id": email_id} def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: email_id = str(mailbox.get("email_id") or "").strip() if not email_id: raise RuntimeError("MoEmail 缺少 email_id") data = self._request("GET", f"/api/emails/{email_id}") items = data.get("messages") or [] messages = [item for item in items if isinstance(item, dict)] if isinstance(items, list) else [] if not messages: return None _, item = max(enumerate(messages), key=lambda pair: (((_parse_received_at(pair[1].get("createdAt") or pair[1].get("created_at") or pair[1].get("receivedAt") or pair[1].get("date") or pair[1].get("timestamp")) or datetime.fromtimestamp(0, tz=timezone.utc)).timestamp()), pair[0])) message_id = str(item.get("id") or item.get("message_id") or item.get("_id") or "").strip() detail = self._request("GET", f"/api/emails/{email_id}/{message_id}") if message_id else {"message": item} message = detail.get("message") if isinstance(detail.get("message"), dict) else detail text_content, html_content = _extract_content(message) sender = message.get("from") or message.get("sender") or "" if isinstance(sender, dict): sender = sender.get("address") or sender.get("email") or sender.get("name") or "" return {"provider": self.name, "mailbox": mailbox["address"], "message_id": message_id, "subject": str(message.get("subject") or item.get("subject") or ""), "sender": str(sender), "text_content": text_content, "html_content": html_content, "received_at": _parse_received_at(message.get("createdAt") or message.get("created_at") or message.get("receivedAt") or message.get("date") or message.get("timestamp") or item.get("createdAt") or item.get("created_at") or item.get("receivedAt") or item.get("date") or item.get("timestamp")), "raw": detail} def close(self) -> None: self.session.close() class InbucketMailProvider(BaseMailProvider): name = "inbucket" def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_base = str(entry["api_base"]).rstrip("/") raw_domains = entry.get("domain") or [] if isinstance(raw_domains, list): self.domain = [str(item).strip() for item in raw_domains if str(item).strip()] else: self.domain = [str(raw_domains).strip()] if str(raw_domains).strip() else [] self.random_subdomain = bool(entry.get("random_subdomain", True)) self.session = requests.Session() self.session.trust_env = False self.session.headers.update({ "User-Agent": conf["user_agent"], "Accept": "application/json", }) def _request(self, method: str, path: str, expected: tuple[int, ...] = (200,)): resp = self.session.request( method.upper(), f"{self.api_base}{path}", timeout=self.conf["request_timeout"], verify=False, ) if resp.status_code not in expected: raise RuntimeError(f"Inbucket 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") if resp.status_code == 204: return {} content_type = str(resp.headers.get("content-type") or "").lower() if "application/json" in content_type: return resp.json() return resp.text def _resolve_domain(self) -> str: if self.domain: return _next_domain(self.domain) raise RuntimeError("Inbucket 需要至少配置一个 domain") def _mailbox_name(self, address: str) -> str: local_part, _, _ = str(address or "").partition("@") return local_part.strip() def create_mailbox(self, username: str | None = None) -> dict[str, Any]: local_part = username or _random_mailbox_name() base_domain = self._resolve_domain() domain = f"{_random_subdomain_label()}.{base_domain}" if self.random_subdomain else base_domain address = f"{local_part}@{domain}" mailbox_name = self._mailbox_name(address) return { "provider": self.name, "provider_ref": self.provider_ref, "address": address, "base_domain": base_domain, "mailbox_name": mailbox_name, } def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: mailbox_name = str(mailbox.get("mailbox_name") or self._mailbox_name(str(mailbox.get("address") or ""))).strip() if not mailbox_name: raise RuntimeError("Inbucket 缺少 mailbox_name") data = self._request("GET", f"/api/v1/mailbox/{mailbox_name}") items = [item for item in data if isinstance(item, dict)] if isinstance(data, list) else [] if not items: return None items.sort( key=lambda value: ( (_parse_received_at(value.get("date")) or datetime.fromtimestamp(0, tz=timezone.utc)).timestamp(), str(value.get("id") or ""), ), reverse=True, ) address = str(mailbox.get("address") or "").strip() for item in items: message_id = str(item.get("id") or "").strip() if not message_id: continue detail = self._request("GET", f"/api/v1/mailbox/{mailbox_name}/{message_id}") if not isinstance(detail, dict): continue header = detail.get("header") if isinstance(detail.get("header"), dict) else {} body = detail.get("body") if isinstance(detail.get("body"), dict) else {} normalized = { "provider": self.name, "mailbox": mailbox_name, "message_id": message_id, "subject": str(detail.get("subject") or item.get("subject") or ""), "sender": str(detail.get("from") or item.get("from") or ""), "text_content": str(body.get("text") or ""), "html_content": str(body.get("html") or ""), "received_at": _parse_received_at(detail.get("date") or item.get("date")), "to": header.get("To") if isinstance(header, dict) else None, "raw": detail, } if _message_matches_email(normalized, address): return normalized return None def close(self) -> None: self.session.close() class YydsMailProvider(BaseMailProvider): name = "yyds_mail" retry_statuses = {429, 500, 502, 503, 504} def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.api_base = str(entry.get("api_base") or "https://maliapi.215.im/v1").rstrip("/") self.api_key = str(entry["api_key"]).strip() self.domain = [str(item).strip() for item in (entry.get("domain") or []) if str(item).strip()] self.subdomain = str(entry.get("subdomain") or "").strip() self.wildcard = bool(entry.get("wildcard")) self.session = requests.Session() self.session.trust_env = False self.session.headers.update({"User-Agent": conf["user_agent"], "Accept": "application/json", "Content-Type": "application/json"}) def _request(self, method: str, path: str, token: str = "", params: dict | None = None, payload: dict | None = None, expected: tuple[int, ...] = (200, 201, 204)): headers = {"Authorization": f"Bearer {token}"} if token else {"X-API-Key": self.api_key} last_error = "" for attempt in range(3): try: resp = self.session.request(method.upper(), f"{self.api_base}{path}", headers=headers, params=params, json=payload, timeout=self.conf["request_timeout"], verify=False) except Exception as exc: last_error = str(exc) if attempt < 2: time.sleep(0.5 * (attempt + 1)) continue raise RuntimeError(f"YYDSMail 请求异常: {method} {path}, error={last_error}") from exc if resp.status_code in expected: break if resp.status_code in self.retry_statuses and attempt < 2: time.sleep(0.5 * (attempt + 1)) continue raise RuntimeError(f"YYDSMail 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") else: raise RuntimeError(f"YYDSMail 请求异常: {method} {path}, error={last_error}") if resp.status_code not in expected: raise RuntimeError(f"YYDSMail 请求失败: {method} {path}, HTTP {resp.status_code}, body={resp.text[:300]}") if resp.status_code == 204: return {} data = resp.json() if isinstance(data, dict) and data.get("success") is False: raise RuntimeError(f"YYDSMail 请求失败: {data.get('errorCode') or data.get('error')}") return data.get("data") if isinstance(data, dict) and isinstance(data.get("data"), (dict, list)) else data @staticmethod def _items(data): return data if isinstance(data, list) else data.get("items") or data.get("messages") or data.get("data") or [] def create_mailbox(self, username: str | None = None) -> dict[str, Any]: payload = {"localPart": username or _random_mailbox_name()} domains = domain_reputation.store.filter_domains(self.name, self.domain) if not domains: domains = domain_reputation.store.good_domains(self.name) if domains: payload["domain"] = _next_domain(domains) if self.subdomain: payload["subdomain"] = self.subdomain data = self._request("POST", "/accounts/wildcard" if self.wildcard else "/accounts", payload=payload) address = str(data.get("address") or data.get("email") or "").strip() token = str(data.get("token") or data.get("temp_token") or data.get("tempToken") or data.get("access_token") or "").strip() if not address or not token: raise RuntimeError("YYDSMail 缺少 address 或 token") domain = address.rsplit("@", 1)[-1].strip().lower() if "@" in address else str(payload.get("domain") or "").strip().lower() return {"provider": self.name, "provider_ref": self.provider_ref, "address": address, "domain": domain, "token": token, "account_id": str(data.get("id") or "")} def fetch_latest_message(self, mailbox: dict[str, Any]) -> dict[str, Any] | None: data = self._request("GET", "/messages", token=str(mailbox.get("token") or ""), params={"address": mailbox["address"]}) messages = [item for item in self._items(data) if isinstance(item, dict)] if not messages: return None item = max(messages, key=lambda value: ((_parse_received_at(value.get("createdAt") or value.get("created_at") or value.get("receivedAt") or value.get("date") or value.get("timestamp")) or datetime.fromtimestamp(0, tz=timezone.utc)).timestamp(), str(value.get("id") or ""))) message_id = str(item.get("id") or item.get("message_id") or "").strip() if message_id: item = self._request("GET", f"/messages/{message_id}", token=str(mailbox.get("token") or ""), params={"address": mailbox["address"]}) text_content, html_content = _extract_content(item) sender = item.get("from") or item.get("sender") or "" if isinstance(sender, dict): sender = sender.get("address") or sender.get("email") or sender.get("name") or "" return {"provider": self.name, "mailbox": mailbox["address"], "message_id": message_id, "subject": str(item.get("subject") or ""), "sender": str(sender), "text_content": text_content, "html_content": html_content, "received_at": _parse_received_at(item.get("createdAt") or item.get("created_at") or item.get("receivedAt") or item.get("date") or item.get("timestamp")), "raw": item} def close(self) -> None: self.session.close() class OutlookPoolProvider(BaseMailProvider): """Lease a real Outlook/Gmail mailbox from the rg-gpt account pool and read the OTP over IMAP. Unlike the temp-mail providers this hands out durable credentials (not an HTTP OTP feed), so ``wait_for_code`` logs into the mailbox itself: Outlook via XOAUTH2 (refresh_token→access_token), Gmail via app password. A base mailbox is shared by its ``+1..+5`` aliases, so messages are filtered by exact recipient. The (rotated) refresh_token is written back into ``mailbox`` so the caller can report it to the pool and keep the stored credential valid.""" name = "outlook_pool" def __init__(self, entry: dict, conf: dict): super().__init__(conf, str(entry.get("provider_ref") or "")) self.base_url = str(entry.get("base_url") or entry.get("api_base") or "").strip().rstrip("/") # Secret may live in the provider config or (preferred on HF) an env var. self.api_key = str(entry.get("api_key") or os.getenv("ACCOUNT_POOL_API_KEY") or os.getenv("POOL_API_KEY") or "").strip() self.kind = str(entry.get("kind") or "").strip().lower() # "outlook" | "gmail" | "" (any) self.client_id = str(entry.get("client_id") or pool_mail.THUNDERBIRD_CLIENT_ID).strip() if not self.base_url or not self.api_key: raise RuntimeError("outlook_pool 需要配置 base_url 和 api_key(或 POOL_API_KEY 环境变量)") self._client = pool_mail.PoolClient(self.base_url, self.api_key, timeout=int(conf["request_timeout"])) def create_mailbox(self, username: str | None = None) -> dict[str, Any]: leased = self._client.lease(count=1, leased_by="gpt2api", kind=self.kind) if not leased: raise RuntimeError(f"账号池已空:pool 无可租用账号(kind={self.kind or 'any'})") acct = leased[0] address = str(acct.get("email") or "").strip().lower() if not address: raise RuntimeError("账号池 lease 未返回 email") domain = address.rsplit("@", 1)[-1] if "@" in address else "" is_gmail = domain in pool_mail.GMAIL_DOMAINS # base_url/api_key are stashed for the (in-process only) report_result call. return { "provider": self.name, "provider_ref": self.provider_ref, "address": address, "domain": domain, "kind": "gmail" if is_gmail else "outlook", "pool_id": acct.get("id"), "lease_token": str(acct.get("lease_token") or ""), "password": str(acct.get("password") or ""), "refresh_token": str(acct.get("refresh_token") or ""), "client_id": str(acct.get("client_id") or self.client_id), "pool_base_url": self.base_url, "pool_api_key": self.api_key, } def wait_for_code(self, mailbox: dict[str, Any]) -> str | None: # OTP is read server-side by the pool (its China egress reaches office365/gmail; this # consumer's does not). We just poll the pool's /otp relay until a fresh code appears. pool_id = mailbox.get("pool_id") if not pool_id: return None lease_token = str(mailbox.get("lease_token") or "") seen = mailbox.setdefault("_seen_codes", []) if not isinstance(seen, list): seen = [] mailbox["_seen_codes"] = seen since = time.time() # OTP was just requested; ignore stale codes from earlier attempts deadline = time.monotonic() + self.conf["wait_timeout"] while time.monotonic() < deadline: try: code = self._client.fetch_otp( int(pool_id), lease_token=lease_token, since=since, exclude=",".join(str(c) for c in seen), ) except Exception as exc: # network / pool error — surface it, keep retrying mailbox["_last_otp_error"] = str(exc)[:200] code = None if code and code not in seen: seen.append(code) return code time.sleep(max(0.5, self.conf["wait_interval"])) return None def report_result(mail_config: dict, mailbox: dict, status: str, **fields: Any) -> None: """Best-effort report of a signup outcome back to the account pool. No-op unless the mailbox came from ``outlook_pool``. Never raises — reporting must not break the register worker.""" if not isinstance(mailbox, dict) or str(mailbox.get("provider") or "") != OutlookPoolProvider.name: return pool_id = mailbox.get("pool_id") base_url = str(mailbox.get("pool_base_url") or "").strip() api_key = str(mailbox.get("pool_api_key") or os.getenv("ACCOUNT_POOL_API_KEY") or os.getenv("POOL_API_KEY") or "").strip() if not pool_id or not base_url or not api_key: return try: client = pool_mail.PoolClient(base_url, api_key) client.report( int(pool_id), status, lease_token=str(mailbox.get("lease_token") or ""), refresh_token=str(mailbox.get("refresh_token") or ""), **fields, ) except Exception: pass def _entries(mail_config: dict) -> list[dict]: return [{**item, "provider_ref": f"{item['type']}#{index + 1}"} for index, item in enumerate(mail_config["providers"])] def _enabled_entries(mail_config: dict) -> list[dict]: items = [item for item in _entries(mail_config) if item.get("enable")] if not items: raise RuntimeError("mail.providers 没有启用的 provider") return items def _next_entry(mail_config: dict) -> dict: global provider_index items = _enabled_entries(mail_config) if len(items) == 1: return dict(items[0]) with provider_lock: value = dict(items[provider_index % len(items)]) provider_index = (provider_index + 1) % len(items) return value def _create_provider(mail_config: dict, provider: str = "", provider_ref: str = "") -> BaseMailProvider: entry = next((dict(item) for item in _entries(mail_config) if provider_ref and item["provider_ref"] == provider_ref), None) entry = entry or next((dict(item) for item in _enabled_entries(mail_config) if provider and item["type"] == provider), None) or _next_entry(mail_config) conf = _config(mail_config) if entry["type"] == "cloudflare_temp_email": return CloudflareTempMailProvider(entry, conf) if entry["type"] == "tempmail_lol": return TempMailLolProvider(entry, conf) if entry["type"] == "duckmail": return DuckMailProvider(entry, conf) if entry["type"] == "gptmail": return GptMailProvider(entry, conf) if entry["type"] == "moemail": return MoEmailProvider(entry, conf) if entry["type"] == "inbucket": return InbucketMailProvider(entry, conf) if entry["type"] == "yyds_mail": return YydsMailProvider(entry, conf) if entry["type"] == "outlook_pool": return OutlookPoolProvider(entry, conf) raise RuntimeError(f"不支持的 mail.provider: {entry['type']}") def create_mailbox(mail_config: dict, username: str | None = None) -> dict: provider = _create_provider(mail_config) try: return provider.create_mailbox(username) finally: provider.close() def wait_for_code(mail_config: dict, mailbox: dict) -> str | None: provider = _create_provider(mail_config, str(mailbox.get("provider") or ""), str(mailbox.get("provider_ref") or "")) try: return provider.wait_for_code(mailbox) finally: provider.close()