hotmail_helper.py 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804
  1. import email
  2. import html
  3. import imaplib
  4. import json
  5. import os
  6. import re
  7. import threading
  8. import time
  9. import traceback
  10. from datetime import datetime, timezone
  11. from email.header import decode_header
  12. from email.utils import parseaddr, parsedate_to_datetime
  13. from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
  14. from urllib.error import HTTPError, URLError
  15. from urllib.parse import urlencode
  16. from urllib.request import Request, urlopen
  17. HOST = "127.0.0.1"
  18. PORT = 17373
  19. LIVE_TOKEN_URL = "https://login.live.com/oauth20_token.srf"
  20. ENTRA_COMMON_TOKEN_URL = "https://login.microsoftonline.com/common/oauth2/v2.0/token"
  21. ENTRA_CONSUMERS_TOKEN_URL = "https://login.microsoftonline.com/consumers/oauth2/v2.0/token"
  22. GRAPH_API_ORIGIN = "https://graph.microsoft.com"
  23. OUTLOOK_API_ORIGIN = "https://outlook.office.com"
  24. GRAPH_SCOPES = "offline_access https://graph.microsoft.com/Mail.Read https://graph.microsoft.com/User.Read"
  25. GRAPH_DEFAULT_SCOPE = "https://graph.microsoft.com/.default"
  26. TOKEN_ENDPOINTS = {
  27. "live": {
  28. "name": "live",
  29. "url": LIVE_TOKEN_URL,
  30. "extra_data": {},
  31. },
  32. "entra-consumers-delegated": {
  33. "name": "entra-consumers-delegated",
  34. "url": ENTRA_CONSUMERS_TOKEN_URL,
  35. "extra_data": {
  36. "scope": GRAPH_SCOPES,
  37. },
  38. },
  39. "entra-common-delegated": {
  40. "name": "entra-common-delegated",
  41. "url": ENTRA_COMMON_TOKEN_URL,
  42. "extra_data": {
  43. "scope": GRAPH_SCOPES,
  44. },
  45. },
  46. "entra-common-default": {
  47. "name": "entra-common-default",
  48. "url": ENTRA_COMMON_TOKEN_URL,
  49. "extra_data": {
  50. "scope": GRAPH_DEFAULT_SCOPE,
  51. },
  52. },
  53. "entra-common-outlook": {
  54. "name": "entra-common-outlook",
  55. "url": ENTRA_COMMON_TOKEN_URL,
  56. "extra_data": {},
  57. },
  58. }
  59. IMAP_HOST = "outlook.office365.com"
  60. IMAP_PORT = 993
  61. REQUEST_TIMEOUT_SECONDS = 45
  62. FETCH_LIMIT_DEFAULT = 5
  63. BASE_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), ".."))
  64. ACCOUNT_LOG_PATH = os.path.join(BASE_DIR, "data", "account-run-history.txt")
  65. ACCOUNT_RECORDS_SNAPSHOT_PATH = os.path.join(BASE_DIR, "data", "account-run-history.json")
  66. ACCOUNT_RECORDS_LOCK = threading.Lock()
  67. def json_response(handler, status, payload):
  68. body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
  69. handler.send_response(status)
  70. handler.send_header("Content-Type", "application/json; charset=utf-8")
  71. handler.send_header("Content-Length", str(len(body)))
  72. handler.send_header("Access-Control-Allow-Origin", "*")
  73. handler.send_header("Access-Control-Allow-Headers", "Content-Type")
  74. handler.send_header("Access-Control-Allow-Methods", "POST, OPTIONS")
  75. handler.end_headers()
  76. handler.wfile.write(body)
  77. def read_json_payload(handler):
  78. length = int(handler.headers.get("Content-Length", "0") or 0)
  79. raw = handler.rfile.read(length) if length > 0 else b"{}"
  80. try:
  81. return json.loads(raw.decode("utf-8"))
  82. except Exception as exc:
  83. raise RuntimeError(f"Invalid JSON payload: {exc}") from exc
  84. def post_form(url, data):
  85. encoded = urlencode(data).encode("utf-8")
  86. request = Request(url, data=encoded, headers={"Content-Type": "application/x-www-form-urlencoded"})
  87. with urlopen(request, timeout=REQUEST_TIMEOUT_SECONDS) as response:
  88. return json.loads(response.read().decode("utf-8"))
  89. def get_json(url, headers=None):
  90. request = Request(url, headers=headers or {})
  91. with urlopen(request, timeout=REQUEST_TIMEOUT_SECONDS) as response:
  92. return response.getcode(), json.loads(response.read().decode("utf-8"))
  93. def mask_secret(value, keep=6):
  94. raw = str(value or "")
  95. if not raw:
  96. return ""
  97. if len(raw) <= keep:
  98. return "*" * len(raw)
  99. return raw[:keep] + "..." + raw[-keep:]
  100. def compact_text(value, limit=400):
  101. text = str(value or "").replace("\r", " ").replace("\n", " ").strip()
  102. return text[:limit]
  103. def log_info(message):
  104. print(f"[HotmailHelper] {message}", flush=True)
  105. def append_account_log(email_addr, password, status, recorded_at="", reason=""):
  106. normalized_email = str(email_addr or "").strip()
  107. normalized_password = str(password or "").strip()
  108. normalized_status = str(status or "").strip().lower()
  109. normalized_recorded_at = str(recorded_at or "").strip() or datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
  110. normalized_reason = str(reason or "").strip().replace("\r", " ").replace("\n", " ")
  111. if not normalized_email or not normalized_password or not normalized_status:
  112. raise RuntimeError("Missing email/password/status for account log append")
  113. os.makedirs(os.path.dirname(ACCOUNT_LOG_PATH), exist_ok=True)
  114. line = f"{normalized_recorded_at}\t{normalized_email}\t{normalized_password}\t{normalized_status}\t{normalized_reason}\n"
  115. with ACCOUNT_RECORDS_LOCK:
  116. with open(ACCOUNT_LOG_PATH, "a", encoding="utf-8") as handle:
  117. handle.write(line)
  118. return ACCOUNT_LOG_PATH
  119. def normalize_account_run_snapshot_record(record):
  120. if not isinstance(record, dict):
  121. return None
  122. email_addr = str(record.get("email") or "").strip()
  123. password = str(record.get("password") or "").strip()
  124. final_status = str(record.get("finalStatus") or "").strip().lower()
  125. if not email_addr or not password or final_status not in {"success", "failed", "stopped"}:
  126. return None
  127. finished_at = str(record.get("finishedAt") or "").strip() or datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
  128. retry_count = max(0, int(record.get("retryCount") or 0))
  129. failed_step_raw = record.get("failedStep")
  130. try:
  131. failed_step = int(failed_step_raw)
  132. except (TypeError, ValueError):
  133. failed_step = None
  134. if failed_step is not None and failed_step <= 0:
  135. failed_step = None
  136. auto_run_context = record.get("autoRunContext") if isinstance(record.get("autoRunContext"), dict) else None
  137. normalized_auto_run_context = None
  138. if auto_run_context:
  139. normalized_auto_run_context = {
  140. "currentRun": max(0, int(auto_run_context.get("currentRun") or 0)),
  141. "totalRuns": max(0, int(auto_run_context.get("totalRuns") or 0)),
  142. "attemptRun": max(0, int(auto_run_context.get("attemptRun") or 0)),
  143. }
  144. if not any(normalized_auto_run_context.values()):
  145. normalized_auto_run_context = None
  146. source = "auto" if str(record.get("source") or "").strip().lower() == "auto" else "manual"
  147. return {
  148. "recordId": str(record.get("recordId") or email_addr).strip() or email_addr,
  149. "email": email_addr,
  150. "password": password,
  151. "finalStatus": final_status,
  152. "finishedAt": finished_at,
  153. "retryCount": retry_count,
  154. "failureLabel": str(record.get("failureLabel") or "").strip(),
  155. "failureDetail": str(record.get("failureDetail") or "").strip(),
  156. "failedStep": failed_step,
  157. "source": source,
  158. "autoRunContext": normalized_auto_run_context,
  159. }
  160. def summarize_account_run_snapshot(records):
  161. summary = {
  162. "total": 0,
  163. "success": 0,
  164. "failed": 0,
  165. "retryTotal": 0,
  166. }
  167. for item in records:
  168. summary["total"] += 1
  169. if item.get("finalStatus") == "success":
  170. summary["success"] += 1
  171. elif item.get("finalStatus") == "failed":
  172. summary["failed"] += 1
  173. summary["retryTotal"] += max(0, int(item.get("retryCount") or 0))
  174. return summary
  175. def normalize_account_run_snapshot_payload(payload):
  176. if not isinstance(payload, dict):
  177. raise RuntimeError("Invalid account run snapshot payload")
  178. normalized_records = []
  179. for item in payload.get("records") if isinstance(payload.get("records"), list) else []:
  180. normalized = normalize_account_run_snapshot_record(item)
  181. if normalized:
  182. normalized_records.append(normalized)
  183. return {
  184. "generatedAt": str(payload.get("generatedAt") or "").strip() or datetime.now(timezone.utc).isoformat().replace("+00:00", "Z"),
  185. "summary": summarize_account_run_snapshot(normalized_records),
  186. "records": normalized_records,
  187. }
  188. def sync_account_run_records(payload):
  189. normalized_payload = normalize_account_run_snapshot_payload(payload)
  190. os.makedirs(os.path.dirname(ACCOUNT_RECORDS_SNAPSHOT_PATH), exist_ok=True)
  191. with ACCOUNT_RECORDS_LOCK:
  192. with open(ACCOUNT_RECORDS_SNAPSHOT_PATH, "w", encoding="utf-8") as handle:
  193. json.dump(normalized_payload, handle, ensure_ascii=False, indent=2)
  194. handle.write("\n")
  195. return ACCOUNT_RECORDS_SNAPSHOT_PATH
  196. def try_refresh_access_token(endpoint, client_id, refresh_token):
  197. request_data = {
  198. "client_id": client_id,
  199. "refresh_token": refresh_token,
  200. "grant_type": "refresh_token",
  201. **(endpoint.get("extra_data") or {}),
  202. }
  203. started_at = time.monotonic()
  204. try:
  205. payload = post_form(endpoint["url"], request_data)
  206. except HTTPError as exc:
  207. detail = exc.read().decode("utf-8", errors="ignore")
  208. return {
  209. "ok": False,
  210. "endpoint": endpoint["name"],
  211. "url": endpoint["url"],
  212. "status": getattr(exc, "code", None),
  213. "error": compact_text(detail or str(exc)),
  214. "elapsed_ms": int((time.monotonic() - started_at) * 1000),
  215. }
  216. except URLError as exc:
  217. return {
  218. "ok": False,
  219. "endpoint": endpoint["name"],
  220. "url": endpoint["url"],
  221. "status": None,
  222. "error": compact_text(f"Token request failed: {exc}"),
  223. "elapsed_ms": int((time.monotonic() - started_at) * 1000),
  224. }
  225. access_token = str(payload.get("access_token") or "").strip()
  226. if not access_token:
  227. return {
  228. "ok": False,
  229. "endpoint": endpoint["name"],
  230. "url": endpoint["url"],
  231. "status": 200,
  232. "error": compact_text(payload.get("error_description") or payload.get("error") or json.dumps(payload, ensure_ascii=False)),
  233. "elapsed_ms": int((time.monotonic() - started_at) * 1000),
  234. }
  235. return {
  236. "ok": True,
  237. "endpoint": endpoint["name"],
  238. "url": endpoint["url"],
  239. "elapsed_ms": int((time.monotonic() - started_at) * 1000),
  240. "payload": {
  241. "access_token": access_token,
  242. "next_refresh_token": str(payload.get("refresh_token") or "").strip(),
  243. },
  244. }
  245. def refresh_access_token(client_id, refresh_token, strategy_names=None):
  246. errors = []
  247. selected_endpoints = [
  248. TOKEN_ENDPOINTS[name]
  249. for name in (strategy_names or ["live", "entra-consumers-delegated", "entra-common-delegated"])
  250. if name in TOKEN_ENDPOINTS
  251. ]
  252. log_info(
  253. "token refresh start "
  254. f"clientId={mask_secret(client_id)} "
  255. f"refreshToken={mask_secret(refresh_token)} "
  256. f"strategies={[item['name'] for item in selected_endpoints]}"
  257. )
  258. for endpoint in selected_endpoints:
  259. result = try_refresh_access_token(endpoint, client_id, refresh_token)
  260. if result["ok"]:
  261. log_info(
  262. "token refresh success "
  263. f"endpoint={result['endpoint']} "
  264. f"elapsedMs={result['elapsed_ms']}"
  265. )
  266. return {
  267. "access_token": result["payload"]["access_token"],
  268. "next_refresh_token": result["payload"]["next_refresh_token"],
  269. "token_endpoint": result["endpoint"],
  270. "token_url": result["url"],
  271. }
  272. errors.append(result)
  273. log_info(
  274. "token refresh failed "
  275. f"endpoint={result['endpoint']} "
  276. f"status={result['status']} "
  277. f"elapsedMs={result['elapsed_ms']} "
  278. f"detail={result['error']}"
  279. )
  280. details = " | ".join(
  281. f"{item['endpoint']}({item['status']}): {item['error']}"
  282. for item in errors
  283. )
  284. raise RuntimeError(f"Token refresh failed on all endpoints: {details}")
  285. def build_xoauth2(email_addr, access_token):
  286. return f"user={email_addr}\x01auth=Bearer {access_token}\x01\x01".encode("utf-8")
  287. def open_mailbox(email_addr, access_token):
  288. client = imaplib.IMAP4_SSL(IMAP_HOST, IMAP_PORT, timeout=REQUEST_TIMEOUT_SECONDS)
  289. client.authenticate("XOAUTH2", lambda _: build_xoauth2(email_addr, access_token))
  290. return client
  291. def decode_mime_header(value):
  292. if not value:
  293. return ""
  294. parts = []
  295. for chunk, charset in decode_header(value):
  296. if isinstance(chunk, bytes):
  297. parts.append(chunk.decode(charset or "utf-8", errors="ignore"))
  298. else:
  299. parts.append(str(chunk))
  300. return "".join(parts).strip()
  301. def extract_text_part(message):
  302. if message.is_multipart():
  303. for part in message.walk():
  304. if part.get_content_maintype() == "multipart":
  305. continue
  306. if "attachment" in str(part.get("Content-Disposition") or "").lower():
  307. continue
  308. payload = part.get_payload(decode=True) or b""
  309. charset = part.get_content_charset() or "utf-8"
  310. text = payload.decode(charset, errors="ignore").strip()
  311. if part.get_content_type() == "text/plain" and text:
  312. return text
  313. if part.get_content_type() == "text/html" and text:
  314. return re.sub(r"\s+", " ", re.sub(r"<[^>]+>", " ", html.unescape(text))).strip()
  315. return ""
  316. payload = message.get_payload(decode=True) or b""
  317. charset = message.get_content_charset() or "utf-8"
  318. text = payload.decode(charset, errors="ignore").strip()
  319. if message.get_content_type() == "text/html":
  320. return re.sub(r"\s+", " ", re.sub(r"<[^>]+>", " ", html.unescape(text))).strip()
  321. return text
  322. def mailbox_candidates(mailbox):
  323. normalized = str(mailbox or "INBOX").strip().lower()
  324. if normalized in {"junk", "junk email", "junk e-mail", "junkemail"}:
  325. return ["Junk", "Junk Email", "Junk E-Mail"]
  326. return ["INBOX"]
  327. def normalize_mailbox_label(mailbox):
  328. normalized = str(mailbox or "INBOX").strip().lower()
  329. if normalized in {"junk", "junk email", "junk e-mail", "junkemail"}:
  330. return "Junk"
  331. return "INBOX"
  332. def normalize_mailbox_id(mailbox):
  333. normalized = str(mailbox or "INBOX").strip().lower()
  334. if normalized in {"junk", "junk email", "junk e-mail", "junkemail"}:
  335. return "junkemail"
  336. return "inbox"
  337. def select_mailbox(client, mailbox):
  338. for candidate in mailbox_candidates(mailbox):
  339. status, _ = client.select(candidate)
  340. if status == "OK":
  341. return candidate
  342. raise RuntimeError(f"Mailbox not found: {mailbox}")
  343. def to_timestamp_ms(raw_date):
  344. if not raw_date:
  345. return 0
  346. try:
  347. parsed = parsedate_to_datetime(raw_date)
  348. if parsed.tzinfo is None:
  349. parsed = parsed.replace(tzinfo=timezone.utc)
  350. return int(parsed.timestamp() * 1000)
  351. except Exception:
  352. return 0
  353. def to_iso_string(timestamp_ms):
  354. if not timestamp_ms:
  355. return ""
  356. return datetime.fromtimestamp(timestamp_ms / 1000, tz=timezone.utc).isoformat().replace("+00:00", "Z")
  357. def normalize_message(message_id, raw_bytes, mailbox):
  358. parsed = email.message_from_bytes(raw_bytes)
  359. sender_name, sender_addr = parseaddr(parsed.get("From", ""))
  360. subject = decode_mime_header(parsed.get("Subject", ""))
  361. body = extract_text_part(parsed)
  362. timestamp_ms = to_timestamp_ms(parsed.get("Date"))
  363. return {
  364. "id": str(message_id),
  365. "mailbox": mailbox,
  366. "subject": subject,
  367. "from": {
  368. "emailAddress": {
  369. "address": sender_addr.strip(),
  370. "name": sender_name.strip(),
  371. }
  372. },
  373. "bodyPreview": body[:500],
  374. "receivedDateTime": to_iso_string(timestamp_ms),
  375. "receivedTimestamp": timestamp_ms,
  376. }
  377. def fetch_messages(email_addr, access_token, mailbox="INBOX", top=FETCH_LIMIT_DEFAULT):
  378. client = None
  379. logical_mailbox = normalize_mailbox_label(mailbox)
  380. try:
  381. client = open_mailbox(email_addr, access_token)
  382. select_mailbox(client, mailbox)
  383. status, data = client.search(None, "ALL")
  384. if status != "OK" or not data or not data[0]:
  385. return {"mailbox": logical_mailbox, "messages": [], "count": 0}
  386. message_ids = data[0].split()
  387. selected_ids = list(reversed(message_ids[-max(1, min(int(top or FETCH_LIMIT_DEFAULT), 30)):]))
  388. messages = []
  389. for message_id in selected_ids:
  390. fetch_status, fetch_data = client.fetch(message_id, "(RFC822)")
  391. if fetch_status != "OK" or not fetch_data:
  392. continue
  393. raw_bytes = b""
  394. for item in fetch_data:
  395. if isinstance(item, tuple) and len(item) >= 2:
  396. raw_bytes = item[1]
  397. break
  398. if not raw_bytes:
  399. continue
  400. messages.append(normalize_message(message_id.decode("utf-8", errors="ignore"), raw_bytes, logical_mailbox))
  401. return {"mailbox": logical_mailbox, "messages": messages, "count": len(messages)}
  402. finally:
  403. if client is not None:
  404. try:
  405. client.logout()
  406. except Exception:
  407. pass
  408. def fetch_messages_for_mailboxes(email_addr, access_token, mailboxes, top):
  409. mailbox_results = []
  410. all_messages = []
  411. for mailbox in mailboxes or ["INBOX"]:
  412. result = fetch_messages(email_addr, access_token, mailbox=mailbox, top=top)
  413. mailbox_results.append(result)
  414. all_messages.extend(result["messages"])
  415. all_messages.sort(key=lambda item: int(item.get("receivedTimestamp") or 0), reverse=True)
  416. return {"mailboxResults": mailbox_results, "messages": all_messages}
  417. def normalize_graph_message(message, mailbox):
  418. sender = message.get("from", {}) or {}
  419. email_addr = sender.get("emailAddress", {}) if isinstance(sender, dict) else {}
  420. received = str(message.get("receivedDateTime") or "").strip()
  421. return {
  422. "id": str(message.get("id") or message.get("internetMessageId") or "").strip(),
  423. "mailbox": mailbox,
  424. "subject": str(message.get("subject") or "").strip(),
  425. "from": {
  426. "emailAddress": {
  427. "address": str(email_addr.get("address") or "").strip(),
  428. "name": str(email_addr.get("name") or "").strip(),
  429. }
  430. },
  431. "bodyPreview": str(message.get("bodyPreview") or "").strip(),
  432. "receivedDateTime": received,
  433. "receivedTimestamp": int(datetime.fromisoformat(received.replace("Z", "+00:00")).timestamp() * 1000) if received else 0,
  434. }
  435. def normalize_outlook_message(message, mailbox):
  436. sender = message.get("From", {}) or message.get("from", {}) or {}
  437. email_addr = sender.get("EmailAddress", {}) if isinstance(sender, dict) else {}
  438. if isinstance(sender, dict) and not email_addr:
  439. email_addr = sender.get("emailAddress", {}) if isinstance(sender, dict) else {}
  440. received = str(message.get("ReceivedDateTime") or message.get("receivedDateTime") or "").strip()
  441. return {
  442. "id": str(message.get("Id") or message.get("id") or "").strip(),
  443. "mailbox": mailbox,
  444. "subject": str(message.get("Subject") or message.get("subject") or "").strip(),
  445. "from": {
  446. "emailAddress": {
  447. "address": str(email_addr.get("Address") or email_addr.get("address") or "").strip(),
  448. "name": str(email_addr.get("Name") or email_addr.get("name") or "").strip(),
  449. }
  450. },
  451. "bodyPreview": str(message.get("BodyPreview") or message.get("bodyPreview") or "").strip(),
  452. "receivedDateTime": received,
  453. "receivedTimestamp": int(datetime.fromisoformat(received.replace("Z", "+00:00")).timestamp() * 1000) if received else 0,
  454. }
  455. def fetch_graph_messages(access_token, mailbox="INBOX", top=FETCH_LIMIT_DEFAULT):
  456. mailbox_id = normalize_mailbox_id(mailbox)
  457. url = (
  458. f"{GRAPH_API_ORIGIN}/v1.0/me/mailFolders/{mailbox_id}/messages"
  459. f"?$top={max(1, min(int(top or FETCH_LIMIT_DEFAULT), 30))}"
  460. f"&$select=id,internetMessageId,subject,from,bodyPreview,receivedDateTime"
  461. f"&$orderby=receivedDateTime desc"
  462. )
  463. try:
  464. _, payload = get_json(url, headers={
  465. "Accept": "application/json",
  466. "Authorization": f"Bearer {access_token}",
  467. })
  468. except HTTPError as exc:
  469. detail = exc.read().decode("utf-8", errors="ignore")
  470. raise RuntimeError(f"Graph request failed: {detail or exc}") from exc
  471. except URLError as exc:
  472. raise RuntimeError(f"Graph request failed: {exc}") from exc
  473. messages = [normalize_graph_message(item, normalize_mailbox_label(mailbox)) for item in (payload.get("value") or [])]
  474. return {"mailbox": normalize_mailbox_label(mailbox), "messages": messages, "count": len(messages)}
  475. def fetch_outlook_api_messages(access_token, mailbox="INBOX", top=FETCH_LIMIT_DEFAULT):
  476. mailbox_id = normalize_mailbox_id(mailbox)
  477. url = (
  478. f"{OUTLOOK_API_ORIGIN}/api/v2.0/me/mailfolders/{mailbox_id}/messages"
  479. f"?$top={max(1, min(int(top or FETCH_LIMIT_DEFAULT), 30))}"
  480. f"&$select=Id,Subject,From,BodyPreview,ReceivedDateTime"
  481. f"&$orderby=ReceivedDateTime desc"
  482. )
  483. try:
  484. _, payload = get_json(url, headers={
  485. "Accept": "application/json",
  486. "Authorization": f"Bearer {access_token}",
  487. })
  488. except HTTPError as exc:
  489. detail = exc.read().decode("utf-8", errors="ignore")
  490. raise RuntimeError(f"Outlook API request failed: {detail or exc}") from exc
  491. except URLError as exc:
  492. raise RuntimeError(f"Outlook API request failed: {exc}") from exc
  493. messages = [normalize_outlook_message(item, normalize_mailbox_label(mailbox)) for item in (payload.get("value") or [])]
  494. return {"mailbox": normalize_mailbox_label(mailbox), "messages": messages, "count": len(messages)}
  495. def collect_imap_messages(email_addr, client_id, refresh_token, mailboxes, top):
  496. token_payload = refresh_access_token(client_id, refresh_token, [
  497. "live",
  498. "entra-consumers-delegated",
  499. "entra-common-delegated",
  500. ])
  501. result = fetch_messages_for_mailboxes(email_addr, token_payload["access_token"], mailboxes, top)
  502. result["transport"] = "imap"
  503. result["token_payload"] = token_payload
  504. return result
  505. def collect_graph_messages(email_addr, client_id, refresh_token, mailboxes, top):
  506. token_payload = refresh_access_token(client_id, refresh_token, [
  507. "entra-common-delegated",
  508. "entra-consumers-delegated",
  509. "entra-common-default",
  510. ])
  511. mailbox_results = [fetch_graph_messages(token_payload["access_token"], mailbox=mailbox, top=top) for mailbox in mailboxes]
  512. messages = []
  513. for item in mailbox_results:
  514. messages.extend(item["messages"])
  515. messages.sort(key=lambda item: int(item.get("receivedTimestamp") or 0), reverse=True)
  516. return {
  517. "transport": "graph",
  518. "token_payload": token_payload,
  519. "mailboxResults": mailbox_results,
  520. "messages": messages,
  521. }
  522. def collect_outlook_messages(email_addr, client_id, refresh_token, mailboxes, top):
  523. token_payload = refresh_access_token(client_id, refresh_token, [
  524. "entra-common-outlook",
  525. "entra-common-delegated",
  526. ])
  527. mailbox_results = [fetch_outlook_api_messages(token_payload["access_token"], mailbox=mailbox, top=top) for mailbox in mailboxes]
  528. messages = []
  529. for item in mailbox_results:
  530. messages.extend(item["messages"])
  531. messages.sort(key=lambda item: int(item.get("receivedTimestamp") or 0), reverse=True)
  532. return {
  533. "transport": "outlook",
  534. "token_payload": token_payload,
  535. "mailboxResults": mailbox_results,
  536. "messages": messages,
  537. }
  538. def collect_messages(email_addr, client_id, refresh_token, mailboxes, top):
  539. errors = []
  540. collectors = [
  541. ("imap", collect_imap_messages),
  542. ("graph", collect_graph_messages),
  543. ("outlook", collect_outlook_messages),
  544. ]
  545. for transport_name, collector in collectors:
  546. try:
  547. log_info(f"message collection start transport={transport_name}")
  548. result = collector(email_addr, client_id, refresh_token, mailboxes, top)
  549. log_info(
  550. f"message collection success transport={transport_name} "
  551. f"tokenEndpoint={result['token_payload'].get('token_endpoint', '')}"
  552. )
  553. return result
  554. except Exception as exc:
  555. message = compact_text(str(exc), 600)
  556. errors.append(f"{transport_name}: {message}")
  557. log_info(f"message collection failed transport={transport_name} detail={message}")
  558. raise RuntimeError(f"Message collection failed on all transports: {' | '.join(errors)}")
  559. def extract_code(text):
  560. source = str(text or "")
  561. patterns = [
  562. r"(?:代码为|验证码[^0-9]*?)[\s::]*(\d{6})",
  563. r"code(?:\s+is|[\s:])+(\d{6})",
  564. r"\b(\d{6})\b",
  565. ]
  566. for pattern in patterns:
  567. match = re.search(pattern, source, flags=re.IGNORECASE)
  568. if match:
  569. return match.group(1)
  570. return ""
  571. def select_latest_code(messages, sender_filters, subject_filters, exclude_codes, filter_after_timestamp):
  572. sender_keywords = [str(item).strip().lower() for item in sender_filters or [] if str(item).strip()]
  573. subject_keywords = [str(item).strip().lower() for item in subject_filters or [] if str(item).strip()]
  574. excluded = {str(item).strip() for item in exclude_codes or [] if str(item).strip()}
  575. def match_message(message, apply_time_filter):
  576. timestamp = int(message.get("receivedTimestamp") or 0)
  577. if apply_time_filter and filter_after_timestamp and timestamp and timestamp < int(filter_after_timestamp):
  578. return None
  579. sender = str(message.get("from", {}).get("emailAddress", {}).get("address", "")).lower()
  580. subject = str(message.get("subject", ""))
  581. preview = str(message.get("bodyPreview", ""))
  582. combined = " ".join([sender, subject.lower(), preview.lower()])
  583. code = extract_code(" ".join([subject, preview, sender]))
  584. if not code or code in excluded:
  585. return None
  586. sender_ok = not sender_keywords or any(keyword in combined for keyword in sender_keywords)
  587. subject_ok = not subject_keywords or any(keyword in combined for keyword in subject_keywords)
  588. if not sender_ok and not subject_ok:
  589. return None
  590. return {"code": code, "message": message}
  591. for use_time_fallback in [False, True]:
  592. matched = []
  593. for message in messages:
  594. result = match_message(message, apply_time_filter=not use_time_fallback)
  595. if result:
  596. matched.append(result)
  597. if matched:
  598. matched.sort(key=lambda item: int(item["message"].get("receivedTimestamp") or 0), reverse=True)
  599. best = matched[0]
  600. return {
  601. "code": best["code"],
  602. "message": best["message"],
  603. "usedTimeFallback": use_time_fallback,
  604. }
  605. return {"code": "", "message": None, "usedTimeFallback": False}
  606. class HotmailHelperHandler(BaseHTTPRequestHandler):
  607. def do_OPTIONS(self):
  608. self.send_response(204)
  609. self.send_header("Access-Control-Allow-Origin", "*")
  610. self.send_header("Access-Control-Allow-Headers", "Content-Type")
  611. self.send_header("Access-Control-Allow-Methods", "POST, OPTIONS")
  612. self.end_headers()
  613. def do_POST(self):
  614. try:
  615. payload = read_json_payload(self)
  616. if self.path == "/sync-account-run-records":
  617. file_path = sync_account_run_records(payload)
  618. json_response(self, 200, {
  619. "ok": True,
  620. "filePath": file_path,
  621. })
  622. return
  623. if self.path == "/append-account-log":
  624. file_path = append_account_log(
  625. payload.get("email"),
  626. payload.get("password"),
  627. payload.get("status"),
  628. payload.get("recordedAt"),
  629. payload.get("reason"),
  630. )
  631. json_response(self, 200, {
  632. "ok": True,
  633. "filePath": file_path,
  634. })
  635. return
  636. email_addr = str(payload.get("email") or "").strip()
  637. client_id = str(payload.get("clientId") or "").strip()
  638. refresh_token = str(payload.get("refreshToken") or "").strip()
  639. if not email_addr or not client_id or not refresh_token:
  640. raise RuntimeError("Missing email/clientId/refreshToken")
  641. top = max(1, min(int(payload.get("top") or FETCH_LIMIT_DEFAULT), 30))
  642. mailboxes = payload.get("mailboxes") if isinstance(payload.get("mailboxes"), list) else [payload.get("mailbox") or "INBOX"]
  643. if self.path == "/messages":
  644. result = collect_messages(email_addr, client_id, refresh_token, mailboxes, top)
  645. json_response(self, 200, {
  646. "ok": True,
  647. "messages": result["messages"],
  648. "mailboxResults": result["mailboxResults"],
  649. "nextRefreshToken": result["token_payload"].get("next_refresh_token") or "",
  650. "tokenEndpoint": result["token_payload"].get("token_endpoint") or "",
  651. "transport": result.get("transport") or "",
  652. })
  653. return
  654. if self.path == "/code":
  655. result = collect_messages(email_addr, client_id, refresh_token, mailboxes, top)
  656. selected = select_latest_code(
  657. result["messages"],
  658. payload.get("senderFilters") or [],
  659. payload.get("subjectFilters") or [],
  660. payload.get("excludeCodes") or [],
  661. int(payload.get("filterAfterTimestamp") or 0),
  662. )
  663. json_response(self, 200, {
  664. "ok": True,
  665. "code": selected["code"],
  666. "message": selected["message"],
  667. "usedTimeFallback": selected["usedTimeFallback"],
  668. "nextRefreshToken": result["token_payload"].get("next_refresh_token") or "",
  669. "tokenEndpoint": result["token_payload"].get("token_endpoint") or "",
  670. "transport": result.get("transport") or "",
  671. })
  672. return
  673. json_response(self, 404, {"ok": False, "error": f"Unsupported path: {self.path}"})
  674. except Exception as exc:
  675. traceback.print_exc()
  676. json_response(self, 500, {"ok": False, "error": str(exc)})
  677. def main():
  678. server = ThreadingHTTPServer((HOST, PORT), HotmailHelperHandler)
  679. print(f"Hotmail helper listening on http://{HOST}:{PORT}", flush=True)
  680. print(f"Account log file: {ACCOUNT_LOG_PATH}", flush=True)
  681. print(f"Account snapshot file: {ACCOUNT_RECORDS_SNAPSHOT_PATH}", flush=True)
  682. try:
  683. server.serve_forever()
  684. except KeyboardInterrupt:
  685. pass
  686. finally:
  687. server.server_close()
  688. if __name__ == "__main__":
  689. main()