8 Commits
Author SHA1 Message Date
dsql a340067048 fix: reconnect-on-stale re-selects folder; lock connect() against races
ensure_connection() reconnected via connect(), which only reaches IMAP
state AUTH, not SELECTED - every search/fetch/store after a mid-pass
link drop silently failed (swallowed into []/None/False) because
aioimaplib rejects those commands outside SELECTED. Track the
currently-selected folder and re-select it after reconnecting; note
pre-reconnect sequence-number ids are invalid post-reconnect unless
use_uid=True.

Also add an asyncio.Lock around connect()/close()/ensure_connection():
concurrent callers on one instance used to race inside connect(),
where task B's `await self.close()` tore down task A's mid-handshake
connection, leaking sockets and orphaning server sessions. Superseded
connections are now always logged out instead of silently overwritten.

Bump to 0.1.6, update README install pins.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-07-02 16:40:50 -04:00
dsql b00f122b74 chore: ignore .claude/ dir (CLAUDE.md now lives under .claude/)
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 21:55:13 -04:00
dsql 3da833f2fc fix: revert OTP-in-logs (spent on arrival, not a secret); F1 ContentTypeError, F5 callable None-guard
revert the M-1 log change — a single-use OTP is consumed on arrival, not a live secret,
so log the code value again. keep the oauth error-body truncate.

F1: oauth token fetch uses resp.json(content_type=None) so a 200 with text/plain doesn't
ContentTypeError and discard a valid token. F5: as_predicate coalesces None for the
callable branch like the string/regex branches. drop a redundant digits.isdigit().

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 21:34:35 -04:00
dsql f940641a5a fix: never log the OTP code value (secret-in-logs); correct false test claim (v0.1.5)
M-1: retrieve.py logged the live single-use code at INFO ('found code %s', 'code %s
skipped too old'), shipping the secret to any aggregation/retention sink the host wires
(our /srv/logs -> loki/grafana path). drop the code value from both lines — log that a
code was found/retrieved and where, never the value. also truncate the oauth token-endpoint
error body to 200 chars so a token response can't be dumped whole.

aiomail-F3: CLAUDE.md claimed an '8-case tested' suite that does not exist in the repo;
corrected to describe the manual throwaway-venv exercise + the real flake8 check.

verified by execution: code retrieved, value absent from logs; control confirms the old
line carried it.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 20:46:51 -04:00
dsql 0cf23805dd docs: pin install line to release, note unpinned-latest option
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 18:13:31 -04:00
dsql 75e6550311 docs: show unpinned install line; note tag-pinning for reproducibility
Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 18:07:16 -04:00
dsql a44bf11be6 fix: dead is_throttled, orphan connect-task, server-defined folder delimiter (v0.1.4)
- remove is_throttled(): read a non-existent .resp -> always False (dead) (L2)
- cancel/await aioimaplib's fire-and-forget create_connection task on a failed connect
  so a refused host doesn't log 'Task exception was never retrieved' per retry (L3)
- get_folders() parses the server-announced LIST delimiter instead of hardcoding '/',
  so '.'/NIL-delimited servers (Gmail/Dovecot) return correct names (L4)
- mark the dead aioimaplib-2.0.x tuple branch + the non-aioimaplib authenticate
  fallback as cross-version escape hatches (nits).

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 17:57:37 -04:00
dsql e349638700 fix: fetch() selects message body by structure, not length (v0.1.3)
select the literal payload by isinstance bytearray instead of len>20. aioimaplib
stores the message body as the only bytearray in the response; every other line
(including the '<id> FETCH (...' header) is plain bytes. the length heuristic
matched the header line first for any 2+ digit message id or BODY[]/UID fetch,
returning a blank Message and silently breaking OTP retrieval on real mailboxes.

Signed-off-by: disqualifier <dev@disqualifier.me>
2026-06-29 17:09:08 -04:00
8 changed files with 149 additions and 37 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
# claude
CLAUDE.md
.claude/
# python
__pycache__/
+6 -4
View File
@@ -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.2
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.2
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.2"
pip install "aiomail[oauth] @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.2"
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
View File
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
[project]
name = "aiomail"
version = "0.1.2"
version = "0.1.6"
description = "async IMAP one-time-code retrieval with password/OAuth2 auth and dynamic matching"
requires-python = ">=3.10"
dependencies = [
+1 -1
View File
@@ -29,4 +29,4 @@ __all__ = [
"DEFAULT_FOLDERS",
]
__version__ = "0.1.2"
__version__ = "0.1.5"
+3
View File
@@ -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
View File
@@ -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
+4 -2
View File
@@ -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()
+7 -2
View File
@@ -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)