gpt2api / services /register /mail_provider.py
jiayi.xie
feat(register): read outlook_pool OTP via pool relay, not direct IMAP
b3df5c1
Raw
History Blame Contribute Delete
42.8 kB
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}"
@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()