Spaces:
Paused
Paused
| 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]*?</p>", 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*<|(?<![#&])\b(\d{6})\b", content): | |
| value = code[0] or code[1] | |
| if value and value != "177010": | |
| return value | |
| return None | |
| def _message_tracking_ref(message: dict[str, Any]) -> 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}" | |
| 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() | |
| 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 | |
| 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() | |