tests rtemoved
parent
29c0362517
commit
f1597792ed
|
|
@ -62,3 +62,4 @@ Utopia-ai-hr-ats-portal 1.pem
|
|||
|
||||
# Local-only Compose overrides (never deployed)
|
||||
docker.local.env
|
||||
tests/**
|
||||
|
|
@ -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()
|
||||
|
|
@ -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)],
|
||||
)
|
||||
Loading…
Reference in New Issue