83 lines
2.6 KiB
Python
83 lines
2.6 KiB
Python
"""Inbox Taskiq tasks — Outlook read-status delta sweep."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import os
|
|
|
|
import httpx
|
|
import redis.asyncio as redis
|
|
from dotenv import load_dotenv
|
|
|
|
from db_setup import session_scope
|
|
from inbox.models import Inbox_Messages
|
|
from inbox.plugins import fetch_read_status_delta
|
|
from taskiq_management.broker_setup import broker
|
|
|
|
load_dotenv()
|
|
|
|
logger=logging.getLogger("inbox.sync")
|
|
|
|
EMAIL_SYNC_FOLDER=os.getenv("EMAIL_SYNC_FOLDER","inbox")
|
|
EMAIL_SYNC_SINCE=os.getenv("EMAIL_SYNC_SINCE") or None
|
|
EMAIL_SYNC_CRON=os.getenv("EMAIL_SYNC_CRON","* * * * *")
|
|
REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0")
|
|
|
|
_LOCK_KEY="inbox:sync_read_status:lock"
|
|
_LOCK_TTL=300
|
|
_MAX_ROUNDS=10
|
|
|
|
|
|
@broker.task(task_name="inbox.sync_read_status",schedule=[{"cron":EMAIL_SYNC_CRON}])
|
|
async def sync_read_status() -> dict:
|
|
client=redis.from_url(REDIS_URL,decode_responses=True)
|
|
try:
|
|
acquired=await client.set(_LOCK_KEY,"1",nx=True,ex=_LOCK_TTL)
|
|
if not acquired:
|
|
logger.info("sync_read_status skipped — lock held")
|
|
return {"skipped":"locked"}
|
|
|
|
try:
|
|
rounds=0
|
|
applied_total=0
|
|
removed_total=0
|
|
since=EMAIL_SYNC_SINCE
|
|
|
|
while rounds<_MAX_ROUNDS:
|
|
rounds+=1
|
|
try:
|
|
round_data=await fetch_read_status_delta(
|
|
EMAIL_SYNC_FOLDER,
|
|
since=since if rounds==1 else None,
|
|
limit=1000,
|
|
max_pages=10,
|
|
)
|
|
except httpx.HTTPStatusError as e:
|
|
if e.response.status_code==401:
|
|
logger.warning("sync_read_status 401 — device-code sign-in required")
|
|
return {"error":"unauthorized","status_code":401}
|
|
raise
|
|
|
|
changes=round_data.get("value") or []
|
|
removed=round_data.get("removed") or []
|
|
removed_total+=len(removed)
|
|
if removed:
|
|
logger.info("sync_read_status removed=%s",len(removed))
|
|
|
|
async with session_scope() as session:
|
|
applied=await Inbox_Messages.apply_read_status(session,changes)
|
|
applied_total+=applied
|
|
|
|
if round_data.get("complete",True):
|
|
break
|
|
|
|
return {
|
|
"rounds":rounds,
|
|
"applied":applied_total,
|
|
"removed":removed_total,
|
|
}
|
|
finally:
|
|
await client.delete(_LOCK_KEY)
|
|
finally:
|
|
await client.aclose()
|