Fix UID-based tracking for already-Seen messages
This commit is contained in:
+123
-32
@@ -422,6 +422,55 @@ async def forward_to_poke(
|
|||||||
return False
|
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
|
# IDLE watcher
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
@@ -459,13 +508,15 @@ async def watch_folder(
|
|||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
# Record existing unseen UIDs so we only forward truly new ones
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
existing_unseen = set(await asyncio.to_thread(client.search, ["UNSEEN"]))
|
last_seen_uid = mailbox_uids[-1] if mailbox_uids else 0
|
||||||
logger.info(
|
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"],
|
account["id"],
|
||||||
folder,
|
folder,
|
||||||
len(existing_unseen),
|
last_seen_uid,
|
||||||
|
len(mailbox_uids),
|
||||||
|
account.get("mark_as_read", False),
|
||||||
)
|
)
|
||||||
backoff = 5
|
backoff = 5
|
||||||
|
|
||||||
@@ -485,25 +536,36 @@ async def watch_folder(
|
|||||||
"[%s/%s] IDLE responses: %s", account["id"], folder, responses
|
"[%s/%s] IDLE responses: %s", account["id"], folder, responses
|
||||||
)
|
)
|
||||||
|
|
||||||
uids = await asyncio.to_thread(client.search, ["UNSEEN"])
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
# Only forward emails that arrived after we started watching
|
new_uids = [uid for uid in mailbox_uids if uid > last_seen_uid]
|
||||||
new_uids = [u for u in uids if u not in existing_unseen]
|
|
||||||
if not new_uids:
|
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
|
continue
|
||||||
|
|
||||||
raw_messages = await asyncio.to_thread(
|
logger.info(
|
||||||
client.fetch, new_uids, ["RFC822"]
|
"[%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:
|
except asyncio.CancelledError:
|
||||||
break
|
break
|
||||||
@@ -534,22 +596,51 @@ async def _poll_folder(
|
|||||||
stop_event: asyncio.Event,
|
stop_event: asyncio.Event,
|
||||||
):
|
):
|
||||||
"""Fallback polling for servers without IDLE support. Checks every 60 seconds."""
|
"""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():
|
while not stop_event.is_set():
|
||||||
try:
|
try:
|
||||||
uids = await asyncio.to_thread(client.search, ["UNSEEN"])
|
mailbox_uids = await asyncio.to_thread(client.search, ["ALL"])
|
||||||
if uids:
|
new_uids = [uid for uid in mailbox_uids if uid > last_seen_uid]
|
||||||
raw_messages = await asyncio.to_thread(client.fetch, uids, ["RFC822"])
|
if new_uids:
|
||||||
for uid, data in raw_messages.items():
|
logger.info(
|
||||||
raw = data.get(b"RFC822", b"")
|
"[%s/%s] Found %d new UID(s) above %d: %s",
|
||||||
if not raw:
|
account["id"],
|
||||||
continue
|
folder,
|
||||||
email_data = parse_email_message(raw)
|
len(new_uids),
|
||||||
await forward_to_poke(email_data, account, webhook_url, api_key)
|
last_seen_uid,
|
||||||
if account.get("mark_as_read", False):
|
_format_uid_list(new_uids),
|
||||||
await asyncio.to_thread(client.set_flags, uids, [b"\\Seen"])
|
)
|
||||||
|
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:
|
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
|
raise # reconnect via outer loop
|
||||||
await asyncio.sleep(60)
|
await asyncio.sleep(60)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user