4 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
7 changed files with 97 additions and 24 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
# claude
CLAUDE.md
.claude/
# python
__pycache__/
+5 -5
View File
@@ -11,22 +11,22 @@ 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.4
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.4
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.4"
pip install "aiomail[oauth] @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiomail.git@v0.1.4"
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.4` suffix from the line above to install the latest unpinned.
Drop the `@v0.1.6` suffix from the line above to install the latest unpinned.
## Password auth
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
[project]
name = "aiomail"
version = "0.1.4"
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.4"
__version__ = "0.1.5"
+75 -9
View File
@@ -4,6 +4,16 @@ 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
@@ -62,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()
@@ -71,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:
@@ -114,22 +139,60 @@ class IMAPClient:
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()
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"""
@@ -154,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"""
+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)