"""Celery tasks: ticket priority escalation.""" import logging from datetime import datetime, timezone, timedelta from app.worker import celery_app logger = logging.getLogger(__name__) def _make_session(): import os from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker db_url = os.getenv("DATABASE_URL", "postgresql://postgres:postgres@db:5432/taptrack_hub") engine = create_engine(db_url.replace("postgresql+asyncpg://", "postgresql://"), pool_pre_ping=True) return sessionmaker(bind=engine)() @celery_app.task(name="tickets.escalate_stale") def escalate_stale(): """ Escalate stale open tickets: - open > 48h + priority medium → high - open > 72h + no reply yet → urgent """ from sqlalchemy import select, and_ from app.models.ticket import SupportTicket, TicketReply, TicketStatus, TicketPriority db = _make_session() try: now = datetime.now(timezone.utc) threshold_48h = now - timedelta(hours=48) threshold_72h = now - timedelta(hours=72) # Open tickets > 48h, medium priority → high stale_medium = db.execute( select(SupportTicket).where( and_( SupportTicket.status == TicketStatus.open, SupportTicket.priority == TicketPriority.medium, SupportTicket.created_at < threshold_48h, ) ) ).scalars().all() for t in stale_medium: t.priority = TicketPriority.high logger.info("Escalated ticket %s to high priority (48h)", t.ticket_number) # Open tickets > 72h, no reply → urgent stale_72h = db.execute( select(SupportTicket).where( and_( SupportTicket.status == TicketStatus.open, SupportTicket.first_response_at == None, # noqa: E711 SupportTicket.created_at < threshold_72h, ) ) ).scalars().all() for t in stale_72h: if t.priority != TicketPriority.urgent: t.priority = TicketPriority.urgent logger.warning("Escalated ticket %s to urgent (72h no response)", t.ticket_number) db.commit() logger.info("escalate_stale: %d→high, %d→urgent", len(stale_medium), len(stale_72h)) except Exception as e: db.rollback() logger.error("escalate_stale error: %s", e) finally: db.close()