"""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)], )