Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a340067048 | ||
|
|
b00f122b74 | ||
|
|
3da833f2fc | ||
|
|
f940641a5a | ||
|
|
0cf23805dd | ||
|
|
75e6550311 | ||
|
|
a44bf11be6 | ||
|
|
e349638700 | ||
|
|
a4abe354eb |
+1
-1
@@ -1,5 +1,5 @@
|
||||
# claude
|
||||
CLAUDE.md
|
||||
.claude/
|
||||
|
||||
# python
|
||||
__pycache__/
|
||||
|
||||
@@ -11,21 +11,23 @@ This reads codes from email; it does not generate them (that is `pyotp`'s job).
|
||||
`requirements.txt`:
|
||||
|
||||
```
|
||||
aiomail @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.1
|
||||
aiomail @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.6
|
||||
# OAuth token providers (Microsoft / Google) need the extra:
|
||||
aiomail[oauth] @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.1
|
||||
aiomail[oauth] @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.6
|
||||
```
|
||||
|
||||
Direct:
|
||||
|
||||
```bash
|
||||
pip install "aiomail @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.1"
|
||||
pip install "aiomail[oauth] @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.1"
|
||||
pip install "aiomail @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.6"
|
||||
pip install "aiomail[oauth] @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.6"
|
||||
```
|
||||
|
||||
Requires `aioimaplib` and `beautifulsoup4` (pulled transitively). The `oauth`
|
||||
extra adds `aiohttp` for the refresh-token providers.
|
||||
|
||||
Drop the `@v0.1.6` suffix from the line above to install the latest unpinned.
|
||||
|
||||
## Password auth
|
||||
|
||||
```python
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "aiomail"
|
||||
version = "0.1.1"
|
||||
version = "0.1.6"
|
||||
description = "async IMAP one-time-code retrieval with password/OAuth2 auth and dynamic matching"
|
||||
requires-python = ">=3.10"
|
||||
dependencies = [
|
||||
|
||||
@@ -29,4 +29,4 @@ __all__ = [
|
||||
"DEFAULT_FOLDERS",
|
||||
]
|
||||
|
||||
__version__ = "0.1.1"
|
||||
__version__ = "0.1.5"
|
||||
|
||||
@@ -94,6 +94,9 @@ class OAuth2Auth:
|
||||
if xoauth2 is not None:
|
||||
result, data = await xoauth2(self.user, token)
|
||||
elif hasattr(mail, "authenticate"):
|
||||
# escape hatch for a non-aioimaplib client: the shipped aioimaplib IMAP4
|
||||
# always has .xoauth2 and never .authenticate, so this branch never runs
|
||||
# for it; the SASL-callback signature here is untested against any driver
|
||||
result, data = await mail.authenticate(
|
||||
"XOAUTH2", lambda _: _sasl_xoauth2(self.user, token)
|
||||
)
|
||||
|
||||
+122
-22
@@ -4,11 +4,22 @@ a thin, provider-agnostic client: it owns the connection lifecycle (connect with
|
||||
retries, reconnect-on-stale, close) and exposes the handful of operations the OTP
|
||||
flow needs (folders, search, fetch, mark-seen). auth is injected, so the same
|
||||
client serves password and OAuth accounts.
|
||||
|
||||
reconnect-on-stale re-selects whatever folder was selected before the drop, so
|
||||
search/fetch/store keep working against a reconnected session. sequence-number ids
|
||||
from before a reconnect are not valid afterward (a fresh SELECT can renumber the
|
||||
mailbox) — pass `use_uid=True` if ids need to survive a reconnect.
|
||||
|
||||
concurrency: one `IMAPClient` instance is not safe to drive from multiple
|
||||
concurrent tasks/coroutines without external serialization; connect/reconnect
|
||||
internally uses a lock to avoid corrupting `_mail`, but overlapping calls to the
|
||||
same instance are not the intended usage pattern.
|
||||
"""
|
||||
import asyncio
|
||||
import email
|
||||
import email.message
|
||||
import logging
|
||||
import re
|
||||
from typing import List, Optional
|
||||
|
||||
from aioimaplib import IMAP4, IMAP4_SSL
|
||||
@@ -17,6 +28,22 @@ from .auth import Auth
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# IMAP LIST reply: (flags) "<delim>" <name> — delim is server-defined (often "/" or
|
||||
# "." or NIL); capture the trailing name regardless, quoted or bare
|
||||
_LIST_RE = re.compile(rb'^\([^)]*\)\s+(?:"[^"]*"|NIL)\s+(.+)$')
|
||||
|
||||
|
||||
def _folder_name(raw: bytes) -> str:
|
||||
"""extract the folder name from a LIST reply line, delimiter-agnostic
|
||||
|
||||
parses the real reply form `(flags) "<delim>" <name>` so any server hierarchy
|
||||
delimiter works (not just "/"); falls back to the last quoted/space token if the
|
||||
line doesn't match the canonical shape.
|
||||
"""
|
||||
match = _LIST_RE.match(raw.strip())
|
||||
name = match.group(1).decode() if match else raw.decode().rsplit(" ", 1)[-1]
|
||||
return name.strip().strip('"')
|
||||
|
||||
|
||||
class IMAPClient:
|
||||
"""connection-managing IMAP client driven by an injected auth mechanism
|
||||
@@ -45,6 +72,8 @@ class IMAPClient:
|
||||
self.timeout = timeout
|
||||
self.max_retries = max_retries
|
||||
self._mail = None
|
||||
self._selected_folder: Optional[str] = None
|
||||
self._lock = asyncio.Lock()
|
||||
|
||||
async def __aenter__(self) -> "IMAPClient":
|
||||
await self.ensure_connection()
|
||||
@@ -54,8 +83,21 @@ class IMAPClient:
|
||||
await self.close()
|
||||
|
||||
async def connect(self) -> bool:
|
||||
"""open a connection and authenticate, retrying with linear backoff"""
|
||||
await self.close()
|
||||
"""open a connection and authenticate, retrying with linear backoff
|
||||
|
||||
serialized by an internal lock: concurrent callers on the same instance queue
|
||||
up rather than tearing down each other's in-progress handshake. a connection
|
||||
superseded while a caller waited is logged out, never silently overwritten.
|
||||
"""
|
||||
async with self._lock:
|
||||
return await self._connect_locked()
|
||||
|
||||
async def _connect_locked(self) -> bool:
|
||||
"""connect()'s body; caller must hold self._lock"""
|
||||
superseded, self._mail = self._mail, None
|
||||
self._selected_folder = None
|
||||
if superseded is not None:
|
||||
await self._discard_mail(superseded)
|
||||
for attempt in range(self.max_retries):
|
||||
try:
|
||||
if self.use_ssl:
|
||||
@@ -68,40 +110,89 @@ class IMAPClient:
|
||||
except Exception as exc:
|
||||
log.warning("connect attempt %d/%d failed: %s", attempt + 1, self.max_retries, exc)
|
||||
if self._mail is not None:
|
||||
try:
|
||||
await self._mail.logout()
|
||||
except Exception as teardown:
|
||||
log.debug("logout error ignored during failed connect: %s", teardown)
|
||||
await self._discard_mail(self._mail)
|
||||
self._mail = None
|
||||
await asyncio.sleep(2 * (attempt + 1))
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
async def _discard_mail(mail) -> None:
|
||||
"""tear down a half-built IMAP4 without leaking its connect task
|
||||
|
||||
aioimaplib's IMAP4 schedules `create_connection` as a fire-and-forget task it
|
||||
never retrieves; on a refused connection that task raises and asyncio logs a
|
||||
noisy "Task exception was never retrieved" traceback. cancel/await it here (and
|
||||
retrieve its exception) before discarding, so a failed connect stays quiet.
|
||||
"""
|
||||
task = getattr(mail, "_client_task", None)
|
||||
if task is not None and not task.done():
|
||||
task.cancel()
|
||||
if task is not None:
|
||||
try:
|
||||
await task
|
||||
except (asyncio.CancelledError, Exception):
|
||||
pass
|
||||
try:
|
||||
await mail.logout()
|
||||
except Exception as teardown:
|
||||
log.debug("logout error ignored during failed connect: %s", teardown)
|
||||
|
||||
async def close(self) -> None:
|
||||
"""log out and drop the connection, swallowing teardown errors"""
|
||||
async with self._lock:
|
||||
await self._close_locked()
|
||||
|
||||
async def _close_locked(self) -> None:
|
||||
"""close()'s body; caller must hold self._lock"""
|
||||
if self._mail is not None:
|
||||
mail, self._mail = self._mail, None
|
||||
try:
|
||||
await self._mail.logout()
|
||||
await mail.logout()
|
||||
except Exception as exc:
|
||||
log.debug("logout error ignored: %s", exc)
|
||||
self._mail = None
|
||||
self._selected_folder = None
|
||||
|
||||
async def ensure_connection(self) -> bool:
|
||||
"""return a live connection, reconnecting if the link is stale"""
|
||||
if self._mail is None:
|
||||
return await self.connect()
|
||||
"""return a live, SELECTED-if-applicable connection, reconnecting if the link is stale
|
||||
|
||||
a reconnect only re-authenticates (state AUTH); if a folder was selected before
|
||||
the drop, it is re-selected here so search/fetch/store keep working afterward.
|
||||
note: sequence-number ids from before the reconnect are NOT valid against the
|
||||
new session (a fresh SELECT can renumber/re-EXISTS the mailbox) unless
|
||||
use_uid=True, in which case UIDs remain stable across the reconnect.
|
||||
|
||||
serialized by an internal lock, so concurrent callers on the same instance
|
||||
never race each other's connect/reconnect; a caller that waits behind another
|
||||
rechecks liveness first instead of tearing down a connection that just came up.
|
||||
"""
|
||||
async with self._lock:
|
||||
if self._mail is not None:
|
||||
try:
|
||||
await self._mail.noop()
|
||||
return True
|
||||
except Exception:
|
||||
return await self.connect()
|
||||
pass
|
||||
return await self._connect_and_reselect_locked()
|
||||
|
||||
def is_throttled(self) -> bool:
|
||||
"""best-effort detection of a provider throttling response"""
|
||||
return bool(
|
||||
self._mail is not None
|
||||
and getattr(self._mail, "resp", None)
|
||||
and "THROTTLED" in str(self._mail.resp)
|
||||
)
|
||||
async def _connect_and_reselect_locked(self) -> bool:
|
||||
"""connect() then re-select the previously-selected folder; caller must hold self._lock"""
|
||||
folder = self._selected_folder
|
||||
if not await self._connect_locked():
|
||||
return False
|
||||
if folder is None:
|
||||
return True
|
||||
try:
|
||||
result, _ = await self._mail.select(f'"{folder}"')
|
||||
except Exception as exc:
|
||||
log.debug("re-select %s after reconnect failed: %s", folder, exc)
|
||||
self._selected_folder = None
|
||||
return False
|
||||
if result != "OK":
|
||||
log.debug("re-select %s after reconnect failed: %s", folder, result)
|
||||
self._selected_folder = None
|
||||
return False
|
||||
self._selected_folder = folder
|
||||
return True
|
||||
|
||||
async def get_folders(self) -> List[str]:
|
||||
"""list mailbox folder names"""
|
||||
@@ -115,7 +206,7 @@ class IMAPClient:
|
||||
folders: List[str] = []
|
||||
for folder in folder_list or []:
|
||||
try:
|
||||
folders.append(folder.decode().split(' "/" ')[-1].strip('"'))
|
||||
folders.append(_folder_name(folder))
|
||||
except Exception:
|
||||
continue
|
||||
return folders
|
||||
@@ -126,10 +217,13 @@ class IMAPClient:
|
||||
return False
|
||||
try:
|
||||
result, _ = await self._mail.select(f'"{folder}"')
|
||||
return result == "OK"
|
||||
except Exception as exc:
|
||||
log.debug("select %s failed: %s", folder, exc)
|
||||
return False
|
||||
if result == "OK":
|
||||
self._selected_folder = folder
|
||||
return True
|
||||
return False
|
||||
|
||||
async def search(self, query: str) -> List[int]:
|
||||
"""search the selected folder, returning ids newest-first"""
|
||||
@@ -169,8 +263,14 @@ class IMAPClient:
|
||||
if result != "OK" or not data:
|
||||
return None
|
||||
for item in data:
|
||||
if isinstance(item, (bytes, bytearray)) and len(item) > 20:
|
||||
# aioimaplib stores the literal message payload as the only bytearray in
|
||||
# the response; every other line (including the `<id> FETCH (...` header)
|
||||
# is plain bytes. select by structure, not length — a length heuristic
|
||||
# mismatches the header line for any 2+ digit id or a BODY[]/UID fetch.
|
||||
if isinstance(item, bytearray):
|
||||
return email.message_from_bytes(bytes(item))
|
||||
# cross-version fallback: aioimaplib 2.0.x never yields tuples here, but an
|
||||
# imaplib-style (header, payload) tuple is handled if a future/alt driver does
|
||||
if isinstance(item, tuple) and len(item) > 1:
|
||||
return email.message_from_bytes(item[1])
|
||||
return None
|
||||
|
||||
@@ -76,7 +76,7 @@ def _scan(text: str, patterns: list[Pattern], lengths: set[int]) -> Optional[str
|
||||
return m.group(1) if m.groups() else m.group(0)
|
||||
for token in re.split(r"\s+", text):
|
||||
digits = "".join(c for c in token if c.isdigit())
|
||||
if digits and len(digits) in lengths and digits.isdigit():
|
||||
if digits and len(digits) in lengths:
|
||||
return digits
|
||||
return None
|
||||
|
||||
@@ -121,6 +121,8 @@ def as_predicate(spec: MatchSpec) -> Callable[[Optional[str]], bool]:
|
||||
if isinstance(spec, re.Pattern):
|
||||
return lambda value: bool(spec.search(value or ""))
|
||||
if callable(spec):
|
||||
return spec
|
||||
# coalesce None like the string/regex branches so the documented Optional[str]
|
||||
# predicate contract holds even if a caller's callable assumes a real string
|
||||
return lambda value: bool(spec(value or ""))
|
||||
needle = str(spec).lower()
|
||||
return lambda value: needle in (value or "").lower()
|
||||
|
||||
@@ -76,12 +76,17 @@ class _RefreshTokenProvider:
|
||||
async with aiohttp.ClientSession(timeout=timeout) as session:
|
||||
async with session.post(endpoint, data=data) as resp:
|
||||
if resp.status == 200:
|
||||
token = (await resp.json()).get("access_token")
|
||||
# content_type=None: some token endpoints return a 200 with
|
||||
# text/plain or text/javascript; default json() would raise
|
||||
# ContentTypeError and discard a valid token body
|
||||
token = (await resp.json(content_type=None)).get("access_token")
|
||||
if token:
|
||||
self._failures = 0
|
||||
return token
|
||||
else:
|
||||
body = await resp.text()
|
||||
# log a truncated error body only — a token-endpoint
|
||||
# response can carry sensitive material; never dump it whole
|
||||
body = (await resp.text())[:200]
|
||||
log.warning("token endpoint %s -> %s: %s", endpoint, resp.status, body)
|
||||
except Exception as exc:
|
||||
log.warning("token request to %s failed: %s", endpoint, exc)
|
||||
|
||||
+22
-10
@@ -20,16 +20,18 @@ log = logging.getLogger(__name__)
|
||||
DEFAULT_FOLDERS: Sequence[str] = ("INBOX", "Junk", "Spam", "Archive", "All Mail")
|
||||
|
||||
|
||||
def _server_query(sender: MatchSpec, subject: MatchSpec) -> str:
|
||||
def _server_query(sender: MatchSpec, subject: MatchSpec, match_field: str = "from") -> str:
|
||||
"""build a narrowing IMAP query from plain-string specs only
|
||||
|
||||
only plain strings translate to server-side FROM/SUBJECT filters; regex and
|
||||
callable specs fall back to ALL and are filtered client-side, so dynamic
|
||||
matching always works even when the server cannot express it.
|
||||
only plain strings translate to server-side filters; regex and callable specs
|
||||
fall back to ALL and are filtered client-side, so dynamic matching always works
|
||||
even when the server cannot express it. `match_field` selects which header the
|
||||
`sender` spec searches: "from" filters by the sender address (default), "to"
|
||||
filters by the recipient address (the per-user alias the code was sent to).
|
||||
"""
|
||||
parts: List[str] = []
|
||||
if isinstance(sender, str):
|
||||
parts.append(f'FROM "{sender}"')
|
||||
parts.append(f'TO "{sender}"' if match_field == "to" else f'FROM "{sender}"')
|
||||
if isinstance(subject, str):
|
||||
parts.append(f'SUBJECT "{subject}"')
|
||||
return f"({' '.join(parts)})" if parts else "ALL"
|
||||
@@ -52,6 +54,7 @@ async def retrieve_otp(
|
||||
*,
|
||||
sender: MatchSpec = None,
|
||||
subject: MatchSpec = None,
|
||||
match_field: str = "from",
|
||||
folders: Optional[Iterable[str]] = None,
|
||||
patterns: Sequence[Union[str, Pattern]] = DEFAULT_PATTERNS,
|
||||
lengths: Iterable[int] = DEFAULT_LENGTHS,
|
||||
@@ -64,14 +67,18 @@ async def retrieve_otp(
|
||||
) -> Optional[str]:
|
||||
"""return the newest OTP matching the filters, or None
|
||||
|
||||
sender/subject accept a substring, a compiled regex, or a callable. folders,
|
||||
patterns, code lengths, max age and retry behavior are all tunable. set
|
||||
`max_age=None` to disable the freshness check.
|
||||
sender/subject accept a substring, a compiled regex, or a callable. `match_field`
|
||||
selects which header the `sender` spec is matched against: "from" (default)
|
||||
matches the sender address; "to" matches the recipient address (the per-user
|
||||
alias the code was sent to) and additionally accepts a forwarded match on the
|
||||
From header, so a forwarded code still resolves. folders, patterns, code lengths,
|
||||
max age and retry behavior are all tunable. set `max_age=None` to disable the
|
||||
freshness check.
|
||||
"""
|
||||
folders = list(folders) if folders is not None else list(DEFAULT_FOLDERS)
|
||||
sender_ok = as_predicate(sender)
|
||||
subject_ok = as_predicate(subject)
|
||||
query = _server_query(sender, subject)
|
||||
query = _server_query(sender, subject, match_field)
|
||||
|
||||
for attempt in range(retries + 1):
|
||||
for folder in folders:
|
||||
@@ -91,7 +98,12 @@ async def retrieve_otp(
|
||||
|
||||
from_hdr = message.get("From", "")
|
||||
subj_hdr = message.get("Subject", "")
|
||||
if not sender_ok(from_hdr) or not subject_ok(subj_hdr):
|
||||
if match_field == "to":
|
||||
to_hdr = message.get("To", "")
|
||||
matched = sender_ok(to_hdr) or sender_ok(from_hdr)
|
||||
else:
|
||||
matched = sender_ok(from_hdr)
|
||||
if not matched or not subject_ok(subj_hdr):
|
||||
continue
|
||||
|
||||
code = extract_code(message, patterns=patterns, lengths=lengths)
|
||||
|
||||
Reference in New Issue
Block a user