"""
Celery Beat task: poll Global Payments for onboarding status updates.

PRD-HWONB-001 §7.4 — runs daily at 06:00 UTC (see beat_schedule entry
"onboarding-poll-status-daily" in src/worker/celery_app.py).

PRD-HWONB-013 §2 — real implementation. Queries every application still "in
flight" with GP (transmitted/in_process/pending), calls
`services.refresh_status()` for each one (which itself calls
`base_boarding_client.get_status()`/`get_activity()`, maps the result via
`helpers.status_machine._map_gp_status_to_local()`, persists the raw payload
to `tsys_decision_payload` regardless of the mapping outcome, and emits
`onboarding.approved`/`onboarding.rejected` on an actual local-status
transition), and aggregates counts for observability.

Each application is polled inside its own try/except so one merchant's GP
call failing (network blip, GP 5xx, unexpected payload shape) never aborts
the whole nightly run for every other merchant.

Reviewer finding #1 (retry loop for stranded provisioning): after the normal
status-poll pass above, this task ALSO re-attempts `services.provision_transit()`
for any application with `status == "approved"` whose merchant is still
`is_onboarded == False` — i.e. `listener.py`'s `onboarding.approved` handler
ran once (dispatched exactly on the local status transition, so it can never
run again for this application) but `provision_transit()` failed that time
(no `tsys_mid` yet, or the platform TSYS_* fallback credentials weren't
configured). Since that event fires only once, this daily poll is the ONLY
retry path available without adding new retry machinery (a dedicated retry
queue/countdown, a `provisioning_attempts` column, etc.) — reusing the
existing daily-poll infrastructure was the explicit design choice here. The
approval email itself is unaffected by this retry (it already went out from
the listener, guarded by `approved_email_sent_at`, regardless of whether
provisioning succeeded) — this loop only ever flips `application.status` to
"provisioned" and `merchant.is_onboarded` to True once provisioning actually
succeeds.
"""

from __future__ import annotations

import asyncio
from typing import Any, Dict

from celery.utils.log import get_task_logger
from sqlalchemy import select

from src.worker.celery_app import celery_app

logger = get_task_logger(__name__)

# Statuses considered "still in flight with GP" — matches PRD-HWONB-013 §2.3's
# instruction to poll both endpoints for anything not yet at a terminal state.
_POLLABLE_STATUSES = ("transmitted", "in_process", "pending")


@celery_app.task(name="onboarding.poll_status", bind=True, max_retries=3)
def poll_status(self, **kwargs) -> Dict[str, Any]:
    """
    1. Query all MerchantOnboardingApplication rows with
       status IN ('transmitted', 'in_process', 'pending').
    2. For each, call services.refresh_status(db, application).
    3. Return counts of {applications_checked, transitioned_approved,
       transitioned_rejected, errors}.

    Runs synchronously inside the Celery worker via asyncio.run() — the
    underlying services/client calls are async (httpx.AsyncClient +
    EventDispatcher.dispatch), and Celery tasks are themselves sync
    functions, so a single event loop is created for the whole task run
    rather than one per application.
    """
    result = asyncio.run(_poll_status_async())
    result["task_id"] = self.request.id
    logger.info("onboarding.poll_status: run complete — %s", result)
    return result


async def _poll_status_async() -> Dict[str, Any]:
    from src.apps.merchant_onboarding import services
    from src.apps.merchant_onboarding.models.application import MerchantOnboardingApplication
    from src.core.database import SessionCelery

    applications_checked = 0
    transitioned_approved = 0
    transitioned_rejected = 0
    errors = 0

    with SessionCelery() as db:
        stmt = select(MerchantOnboardingApplication).where(
            MerchantOnboardingApplication.status.in_(_POLLABLE_STATUSES),
            MerchantOnboardingApplication.deleted_at.is_(None),
        )
        application_ids = [row.id for row in db.execute(stmt).scalars().all()]

    for application_id in application_ids:
        with SessionCelery() as db:
            application = db.execute(
                select(MerchantOnboardingApplication).where(
                    MerchantOnboardingApplication.id == application_id,
                    MerchantOnboardingApplication.deleted_at.is_(None),
                )
            ).scalar_one_or_none()
            if application is None:
                continue

            applications_checked += 1
            try:
                outcome = await services.refresh_status(db, application)
            except Exception as exc:
                errors += 1
                logger.error(
                    "onboarding.poll_status: refresh_status failed for application %s: %s",
                    application_id,
                    exc,
                )
                continue

            transitioned_to = outcome.get("transitioned_to")
            if transitioned_to == "approved":
                transitioned_approved += 1
            elif transitioned_to == "rejected":
                transitioned_rejected += 1

    provisioning_result = await _retry_pending_provisioning_async()

    return {
        "status": "ok",
        "applications_checked": applications_checked,
        "transitioned_approved": transitioned_approved,
        "transitioned_rejected": transitioned_rejected,
        "errors": errors,
        **provisioning_result,
    }


async def _retry_pending_provisioning_async() -> Dict[str, Any]:
    """
    Reviewer finding #1 retry loop — re-attempt `services.provision_transit()`
    for every application that is locally `status == "approved"` but whose
    merchant is still `is_onboarded == False`. This covers both "never
    provisioned yet" (listener.py's `onboarding.approved` handler hit an
    APIException — no tsys_mid yet, or the platform TSYS_* fallback
    credentials aren't configured) and any future case where provisioning
    needs re-running. Runs after the main status-poll pass above so an
    application that just transitioned to "approved" this very run is
    already eligible to be picked up in the same task execution.

    Each application is retried inside its own try/except (same isolation
    pattern as `_poll_status_async` above) so one merchant's provisioning
    failure never blocks the retry attempt for any other merchant.
    """
    from src.apps.merchant_onboarding import services
    from src.apps.merchant_onboarding.models.application import MerchantOnboardingApplication
    from src.apps.merchants.models.merchant import Merchant
    from src.core.database import SessionCelery

    provisioning_retried = 0
    provisioning_succeeded = 0
    provisioning_errors = 0

    with SessionCelery() as db:
        stmt = (
            select(MerchantOnboardingApplication.id)
            .join(Merchant, Merchant.id == MerchantOnboardingApplication.merchant_id)
            .where(
                MerchantOnboardingApplication.status == "approved",
                MerchantOnboardingApplication.deleted_at.is_(None),
                Merchant.is_onboarded.is_(False),
            )
        )
        application_ids = [row[0] for row in db.execute(stmt).all()]

    for application_id in application_ids:
        with SessionCelery() as db:
            application = db.execute(
                select(MerchantOnboardingApplication).where(
                    MerchantOnboardingApplication.id == application_id,
                    MerchantOnboardingApplication.deleted_at.is_(None),
                )
            ).scalar_one_or_none()
            # Re-check status/is_onboarded against this fresh session — the
            # id list above may be stale by the time we get here (e.g. a
            # concurrent manual `GET /onboarding/status` refresh already
            # provisioned this same application).
            if application is None or application.status != "approved":
                continue

            merchant = db.execute(
                select(Merchant).where(Merchant.id == application.merchant_id)
            ).scalar_one_or_none()
            if merchant is None or merchant.is_onboarded:
                continue

            provisioning_retried += 1
            try:
                await services.provision_transit(db, application)
            except Exception as exc:
                provisioning_errors += 1
                logger.error(
                    "onboarding.poll_status: provisioning retry failed for "
                    "application_id=%s merchant_id=%s tsys_mid=%s: %s",
                    application_id,
                    application.merchant_id,
                    application.tsys_mid,
                    exc,
                )
                continue

            application.status = "provisioned"
            merchant.is_onboarded = True
            db.commit()
            provisioning_succeeded += 1
            logger.info(
                "onboarding.poll_status: provisioning retry succeeded for application_id=%s",
                application_id,
            )

    return {
        "provisioning_retried": provisioning_retried,
        "provisioning_succeeded": provisioning_succeeded,
        "provisioning_errors": provisioning_errors,
    }
