"""CRUD для очереди отложенных SMS (pending_sms)."""

from __future__ import annotations

import logging
from datetime import datetime

from sqlalchemy import func, select, update
from sqlalchemy.exc import IntegrityError

from app.database.engine import async_session
from app.database.models import PendingSms

_log = logging.getLogger(__name__)

STATUS_PENDING = "pending"
STATUS_SENT = "sent"
STATUS_FAILED = "failed"
STATUS_CANCELLED = "cancelled"


async def enqueue_pending_sms(
    *,
    user_id: int,
    phone_number: str,
    ad_id: str | None = None,
    ad_name: str | None = None,
    link: str | None = None,
) -> PendingSms | None:
    """
    Добавляет SMS в очередь, если такого pending ещё нет.
    Возвращает запись или None, если уже есть pending на этот номер.
    Защита: SELECT + unique partial index (гонки) + IntegrityError.
    """
    phone_number = (phone_number or "").strip()
    if not phone_number:
        return None

    async with async_session() as session:
        try:
            async with session.begin():
                existing = await session.execute(
                    select(PendingSms).where(
                        PendingSms.phone_number == phone_number,
                        PendingSms.status == STATUS_PENDING,
                    )
                )
                if existing.scalar_one_or_none() is not None:
                    _log.info(
                        "Pending SMS уже есть для %s — не дублируем",
                        phone_number,
                    )
                    return None

                row = PendingSms(
                    user_id=int(user_id),
                    phone_number=phone_number,
                    ad_id=str(ad_id) if ad_id is not None else None,
                    ad_name=ad_name,
                    link=link,
                    status=STATUS_PENDING,
                    created_at=datetime.utcnow(),
                    updated_at=datetime.utcnow(),
                )
                session.add(row)
                await session.flush()
                _log.info(
                    "Pending SMS поставлена в очередь: phone=%s user_id=%s ad_id=%s",
                    phone_number,
                    user_id,
                    ad_id,
                )
                return row
        except IntegrityError:
            _log.info(
                "Pending SMS concurrent dup для %s — не дублируем",
                phone_number,
            )
            return None


async def list_pending_sms(limit: int = 50) -> list[PendingSms]:
    """Активные pending, старые первыми."""
    async with async_session() as session:
        result = await session.execute(
            select(PendingSms)
            .where(PendingSms.status == STATUS_PENDING)
            .order_by(PendingSms.created_at.asc())
            .limit(max(1, int(limit)))
        )
        return list(result.scalars().all())


async def mark_pending_sms_status(
    pending_id: int,
    status: str,
) -> bool:
    async with async_session() as session:
        async with session.begin():
            result = await session.execute(
                update(PendingSms)
                .where(PendingSms.id == int(pending_id))
                .values(status=status, updated_at=datetime.utcnow())
            )
            return bool(result.rowcount)


async def count_pending_sms(user_id: int | None = None) -> int:
    async with async_session() as session:
        stmt = (
            select(func.count())
            .select_from(PendingSms)
            .where(PendingSms.status == STATUS_PENDING)
        )
        if user_id is not None:
            stmt = stmt.where(PendingSms.user_id == int(user_id))
        result = await session.execute(stmt)
        return int(result.scalar_one() or 0)


async def cancel_pending_for_phone(phone_number: str) -> int:
    """Отменяет pending на номер (например, после успешной прямой отправки)."""
    async with async_session() as session:
        async with session.begin():
            result = await session.execute(
                update(PendingSms)
                .where(
                    PendingSms.phone_number == phone_number,
                    PendingSms.status == STATUS_PENDING,
                )
                .values(status=STATUS_CANCELLED, updated_at=datetime.utcnow())
            )
            return int(result.rowcount or 0)
