From f1597792ed15a441bff648206c2f22a98e23c0e9 Mon Sep 17 00:00:00 2001 From: "ahmed.mujtaba" Date: Mon, 24 Aug 2026 20:40:08 +0500 Subject: [PATCH] tests rtemoved --- .gitignore | 1 + backend/inbox/mailbox_sync_tasks.py | 97 +++++++++++++++++++ .../mailbox_sync_broker_setup.py | 59 +++++++++++ 3 files changed, 157 insertions(+) create mode 100644 backend/inbox/mailbox_sync_tasks.py create mode 100644 backend/taskiq_management/mailbox_sync_broker_setup.py diff --git a/.gitignore b/.gitignore index a043f7a..2275f4b 100644 --- a/.gitignore +++ b/.gitignore @@ -62,3 +62,4 @@ Utopia-ai-hr-ats-portal 1.pem # Local-only Compose overrides (never deployed) docker.local.env +tests/** \ No newline at end of file diff --git a/backend/inbox/mailbox_sync_tasks.py b/backend/inbox/mailbox_sync_tasks.py new file mode 100644 index 0000000..e4f90d9 --- /dev/null +++ b/backend/inbox/mailbox_sync_tasks.py @@ -0,0 +1,97 @@ +"""Mailbox sync Taskiq tasks — Outlook pull + triage + ingest on own stream.""" + +from __future__ import annotations + +import logging +import os +from datetime import datetime,timezone + +import redis.asyncio as redis +from dotenv import load_dotenv + +from db_setup import session_scope +from inbox.models import MailboxSyncRun +from inbox.views import Email +from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY +from taskiq_management.mailbox_sync_broker_setup import mailbox_sync_broker +from taskiq_management.middleware import PermanentTaskError + +load_dotenv() + +logger=logging.getLogger("inbox.mailbox_sync") +REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0") +_LOCK_KEY="inbox:mailbox_sync:lock" +_LOCK_TTL=900 + + +async def _fail(run_id:str,error:str) -> dict: + async with session_scope() as session: + await MailboxSyncRun.update_run(session,run_id,{ + "status":"failed", + "error":error, + "finished_at":datetime.now(timezone.utc), + }) + return {"status":"failed","error":error} + + +@mailbox_sync_broker.task( + task_name="inbox.sync_mailbox", + retry_on_error=True, + max_retries=MAX_RETRIES, + delay=RETRY_DELAY, +) +async def sync_mailbox(run_id:str) -> dict: + if not run_id or not str(run_id).strip(): + raise PermanentTaskError("run_id is required") + run_id=str(run_id).strip() + + client=redis.from_url(REDIS_URL,decode_responses=True) + try: + acquired=await client.set(_LOCK_KEY,run_id,nx=True,ex=_LOCK_TTL) + if not acquired: + return await _fail(run_id,"another mailbox sync is already running") + + try: + async with session_scope() as session: + row=await MailboxSyncRun.get_by_id(session,run_id) + if not row: + raise PermanentTaskError(f"sync run {run_id} not found") + await MailboxSyncRun.update_run(session,run_id,{ + "status":"running", + "started_at":datetime.now(timezone.utc), + "error":None, + }) + top=row.top or 100 + skip=row.skip or 0 + test_on=True if row.test_on is None else bool(row.test_on) + + async with session_scope() as session: + service=Email(session=session) + if not service.token: + return await _fail(run_id,"EMAIL_API_TOKEN is not configured") + try: + summary=await service.run_mailbox_sync_page( + top=top,skip=skip,test_on=test_on, + ) + except Exception as e: + logger.exception("mailbox sync failed for run %s",run_id) + return await _fail(run_id,str(e)) + + await MailboxSyncRun.update_run(session,run_id,{ + "status":"completed", + "entries":summary["entries"], + "triage":summary["triage"], + "error":None, + "finished_at":datetime.now(timezone.utc), + }) + return { + "status":"completed", + "triage":summary["triage"], + "entries":len(summary["entries"]), + } + finally: + current=await client.get(_LOCK_KEY) + if current==run_id: + await client.delete(_LOCK_KEY) + finally: + await client.aclose() diff --git a/backend/taskiq_management/mailbox_sync_broker_setup.py b/backend/taskiq_management/mailbox_sync_broker_setup.py new file mode 100644 index 0000000..2ab9a1f --- /dev/null +++ b/backend/taskiq_management/mailbox_sync_broker_setup.py @@ -0,0 +1,59 @@ +"""Taskiq mailbox-sync broker — isolated Redis stream so Outlook pull/triage +never blocks inbox match/ATS or the cv_upload queue. + +Worker: taskiq worker taskiq_management.mailbox_sync_broker_setup:mailbox_sync_broker inbox.mailbox_sync_tasks +""" + +from __future__ import annotations + +import os + +from dotenv import load_dotenv +from taskiq import TaskiqScheduler +from taskiq.middlewares import SmartRetryMiddleware +from taskiq.schedule_sources import LabelScheduleSource +from taskiq_redis import ( + ListRedisScheduleSource, + RedisAsyncResultBackend, + RedisStreamBroker, +) + +from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY +from taskiq_management.middleware import DeadLetterMiddleware + +load_dotenv() + +REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0") +MAILBOX_SYNC_QUEUE_NAME=os.getenv("TASKIQ_MAILBOX_SYNC_QUEUE_NAME","mailbox_sync") + +result_backend=RedisAsyncResultBackend(redis_url=REDIS_URL) +mailbox_sync_schedule_source=ListRedisScheduleSource( + url=REDIS_URL,prefix="taskiq:schedule:mailbox_sync", +) + +mailbox_sync_broker=( + RedisStreamBroker( + url=REDIS_URL, + queue_name=MAILBOX_SYNC_QUEUE_NAME, + consumer_group_name=os.getenv("TASKIQ_CONSUMER_GROUP","taskiq"), + idle_timeout=int(os.getenv("TASKIQ_IDLE_TIMEOUT_MS","600000")), + ) + .with_result_backend(result_backend) + .with_middlewares( + DeadLetterMiddleware(redis_url=REDIS_URL), + SmartRetryMiddleware( + default_retry_count=MAX_RETRIES, + default_retry_label=True, + default_delay=RETRY_DELAY, + use_jitter=True, + use_delay_exponent=True, + max_delay_exponent=float(os.getenv("TASKIQ_MAX_DELAY","120")), + schedule_source=mailbox_sync_schedule_source, + ), + ) +) + +mailbox_sync_scheduler=TaskiqScheduler( + broker=mailbox_sync_broker, + sources=[mailbox_sync_schedule_source,LabelScheduleSource(mailbox_sync_broker)], +)