From 3d4a1cb791b31d4cf17627d6338bd1248d904696 Mon Sep 17 00:00:00 2001 From: 0xK Date: Mon, 30 Mar 2026 16:37:37 +0200 Subject: [PATCH] Fix UID-based tracking for already-Seen messages --- src/server.py | 155 +++++++++++++++++++++++++++++++++++++++----------- 1 file changed, 123 insertions(+), 32 deletions(-) diff --git a/src/server.py b/src/server.py index 09b9708..a4ac41f 100644 --- a/src/server.py +++ b/src/server.py @@ -422,6 +422,55 @@ async def forward_to_poke( return False +def _format_uid_list(uids: list[int], limit: int = 10) -> str: + if not uids: + return "[]" + shown = ", ".join(str(uid) for uid in uids[:limit]) + if len(uids) > limit: + shown += f", ... (+{len(uids) - limit} more)" + return f"[{shown}]" + + +async def _forward_uid_batch( + client: IMAPClient, + account: dict, + folder: str, + webhook_url: str, + api_key: str, + uids: list[int], +) -> None: + logger.info( + "[%s/%s] Forwarding %d message(s) for UIDs %s", + account["id"], + folder, + len(uids), + _format_uid_list(uids), + ) + raw_messages = await asyncio.to_thread(client.fetch, uids, ["RFC822"]) + for uid in uids: + data = raw_messages.get(uid, {}) + raw = data.get(b"RFC822", b"") + if not raw: + logger.debug( + "[%s/%s] Skipping UID %s because RFC822 payload was empty", + account["id"], + folder, + uid, + ) + continue + email_data = parse_email_message(raw) + await forward_to_poke(email_data, account, webhook_url, api_key) + + if account.get("mark_as_read", False): + await asyncio.to_thread(client.set_flags, uids, [b"\\Seen"]) + logger.debug( + "[%s/%s] Marked UIDs %s as Seen", + account["id"], + folder, + _format_uid_list(uids), + ) + + # --------------------------------------------------------------------------- # IDLE watcher # --------------------------------------------------------------------------- @@ -459,13 +508,15 @@ async def watch_folder( ) return - # Record existing unseen UIDs so we only forward truly new ones - existing_unseen = set(await asyncio.to_thread(client.search, ["UNSEEN"])) + mailbox_uids = await asyncio.to_thread(client.search, ["ALL"]) + last_seen_uid = mailbox_uids[-1] if mailbox_uids else 0 logger.info( - "[%s/%s] Watching for new emails via IDLE (%d existing unseen skipped)", + "[%s/%s] Watching for new emails via IDLE from UID %d (%d existing message(s), mark_as_read=%s)", account["id"], folder, - len(existing_unseen), + last_seen_uid, + len(mailbox_uids), + account.get("mark_as_read", False), ) backoff = 5 @@ -485,25 +536,36 @@ async def watch_folder( "[%s/%s] IDLE responses: %s", account["id"], folder, responses ) - uids = await asyncio.to_thread(client.search, ["UNSEEN"]) - # Only forward emails that arrived after we started watching - new_uids = [u for u in uids if u not in existing_unseen] + mailbox_uids = await asyncio.to_thread(client.search, ["ALL"]) + new_uids = [uid for uid in mailbox_uids if uid > last_seen_uid] if not new_uids: + logger.debug( + "[%s/%s] No new UIDs above %d (mailbox latest UID %d)", + account["id"], + folder, + last_seen_uid, + mailbox_uids[-1] if mailbox_uids else last_seen_uid, + ) continue - raw_messages = await asyncio.to_thread( - client.fetch, new_uids, ["RFC822"] + logger.info( + "[%s/%s] Found %d new UID(s) above %d: %s", + account["id"], + folder, + len(new_uids), + last_seen_uid, + _format_uid_list(new_uids), + ) + await _forward_uid_batch( + client, account, folder, webhook_url, api_key, new_uids + ) + last_seen_uid = new_uids[-1] + logger.debug( + "[%s/%s] Advanced UID cursor to %d", + account["id"], + folder, + last_seen_uid, ) - for uid, data in raw_messages.items(): - raw = data.get(b"RFC822", b"") - if not raw: - continue - email_data = parse_email_message(raw) - await forward_to_poke(email_data, account, webhook_url, api_key) - - if account.get("mark_as_read", False): - await asyncio.to_thread(client.set_flags, new_uids, [b"\\Seen"]) - existing_unseen.update(new_uids) except asyncio.CancelledError: break @@ -534,22 +596,51 @@ async def _poll_folder( stop_event: asyncio.Event, ): """Fallback polling for servers without IDLE support. Checks every 60 seconds.""" - logger.info("[%s/%s] Polling for new emails every 60s", account["id"], folder) + mailbox_uids = await asyncio.to_thread(client.search, ["ALL"]) + last_seen_uid = mailbox_uids[-1] if mailbox_uids else 0 + logger.info( + "[%s/%s] Polling for new emails every 60s from UID %d (%d existing message(s), mark_as_read=%s)", + account["id"], + folder, + last_seen_uid, + len(mailbox_uids), + account.get("mark_as_read", False), + ) while not stop_event.is_set(): try: - uids = await asyncio.to_thread(client.search, ["UNSEEN"]) - if uids: - raw_messages = await asyncio.to_thread(client.fetch, uids, ["RFC822"]) - for uid, data in raw_messages.items(): - raw = data.get(b"RFC822", b"") - if not raw: - continue - email_data = parse_email_message(raw) - await forward_to_poke(email_data, account, webhook_url, api_key) - if account.get("mark_as_read", False): - await asyncio.to_thread(client.set_flags, uids, [b"\\Seen"]) + mailbox_uids = await asyncio.to_thread(client.search, ["ALL"]) + new_uids = [uid for uid in mailbox_uids if uid > last_seen_uid] + if new_uids: + logger.info( + "[%s/%s] Found %d new UID(s) above %d: %s", + account["id"], + folder, + len(new_uids), + last_seen_uid, + _format_uid_list(new_uids), + ) + await _forward_uid_batch( + client, account, folder, webhook_url, api_key, new_uids + ) + last_seen_uid = new_uids[-1] + logger.debug( + "[%s/%s] Advanced UID cursor to %d", + account["id"], + folder, + last_seen_uid, + ) + else: + logger.debug( + "[%s/%s] No new UIDs above %d (mailbox latest UID %d)", + account["id"], + folder, + last_seen_uid, + mailbox_uids[-1] if mailbox_uids else last_seen_uid, + ) except Exception as e: - logger.warning("[%s/%s] Poll error: %s", account["id"], folder, e) + logger.warning( + "[%s/%s] Poll error: %s", account["id"], folder, e + ) raise # reconnect via outer loop await asyncio.sleep(60)