"""Finance Celery beats: variance scan + payment-due reminders.

Scheduled in :mod:`config.celery`.
"""

from __future__ import annotations

import logging
from datetime import timedelta
from decimal import Decimal

from celery import shared_task
from django.core.mail import EmailMessage
from django.utils import timezone

logger = logging.getLogger(__name__)


@shared_task
def dispatch_donor_closeout_report(funding_record_id: int) -> int:
    """PRD §5.1 FRFA-CO019 — at close-out, email a donor close-out report PDF
    to every distinct donor (FundingSource → Partner) that funded the
    award's parent funding record.

    Returns the number of donor reports dispatched. Idempotent within a 7-day
    window: re-running for the same (funding_record, donor) is a no-op.
    """
    from apps.core.audit.models import AuditLog
    from apps.rims.finance.services_donor_closeout import (
        _already_dispatched_recently,
        funding_sources_for_award,
        generate_donor_closeout_pdf,
    )
    from apps.rims.grants.audit_helper import audit_rims
    from apps.rims.grants.models import Award, FundingRecord

    funding_record = FundingRecord.objects.filter(pk=funding_record_id).first()
    if funding_record is None:
        return 0

    closed_awards = list(
        Award.objects.filter(
            application__call__funding_record=funding_record,
            status=Award.Status.CLOSED,
        )
    )
    if not closed_awards:
        return 0

    # Compute the union of donors across every closed award under this funding.
    donors: dict[int, object] = {}
    for award in closed_awards:
        for donor in funding_sources_for_award(award):
            donors[donor.pk] = donor

    # Period: from earliest closed_at to now (broad-by-design — donor reports
    # tend to want the cumulative picture, not a slice).
    period_start = min(a.closed_at.date() for a in closed_awards if a.closed_at)
    period_end = timezone.now().date()

    dispatched = 0
    for donor_id, donor in donors.items():
        if _already_dispatched_recently(funding_record_id, donor_id):
            logger.info(
                "dispatch_donor_closeout_report: skipping donor=%s fr=%s — recent audit",
                donor_id,
                funding_record_id,
            )
            continue
        partner = getattr(donor, "partner", None)
        contact_email = getattr(partner, "contact_email", "") if partner else ""
        if not contact_email:
            logger.info(
                "dispatch_donor_closeout_report: donor=%s has no partner contact email — skipping",
                donor_id,
            )
            continue
        try:
            pdf_bytes = generate_donor_closeout_pdf(
                donor, period_start=period_start, period_end=period_end
            )
        except Exception as exc:  # noqa: BLE001
            logger.exception(
                "dispatch_donor_closeout_report: PDF render failed donor=%s fr=%s: %s",
                donor_id,
                funding_record_id,
                exc,
            )
            continue

        subject = f"Close-out report: {funding_record.title}"
        body = (
            f"Dear {partner.name},\n\nAttached is the close-out report for "
            f"{funding_record.title} covering closed awards funded by your "
            f"contributions through {period_end:%d %B %Y}.\n\nKind regards,\nRUFORUM"
        )
        msg = EmailMessage(subject=subject, body=body, to=[contact_email])
        msg.attach(
            f"close-out-report-{funding_record.pk}-donor-{donor_id}.pdf",
            pdf_bytes,
            "application/pdf",
        )
        msg.send(fail_silently=False)
        audit_rims(
            actor=None,
            action="DONOR_CLOSEOUT_REPORT_SENT",
            target_model="FundingRecord",
            object_id=funding_record.pk,
            object_repr=str(funding_record),
            changes={
                "donor_id": donor_id,
                "donor_name": donor.name,
                "contact_email": contact_email,
                "period_start": period_start.isoformat(),
                "period_end": period_end.isoformat(),
                "closed_award_count": len(closed_awards),
            },
        )
        dispatched += 1
    if dispatched:
        logger.info(
            "dispatch_donor_closeout_report: dispatched %s donor report(s) for fr=%s",
            dispatched,
            funding_record_id,
        )
    return dispatched


@shared_task
def scan_budget_variances() -> int:
    """PRD §5.4 NFRFM004 — alert finance staff when a budget's spend ratio
    diverges from the elapsed-time ratio by more than the configured
    percent_threshold.

    Delegates the actual gap calculation to
    :func:`apps.rims.finance.services_analytics.aggregate_variance_snapshot`
    so the operator dashboard and the alert path agree on what counts as a
    variance.
    """
    from apps.rims.finance.notifications import _finance_alerts_users
    from apps.rims.finance.services_analytics import aggregate_variance_snapshot
    from apps.rims.operations.models import VarianceAlertConfig
    from apps.core.notifications.emailing import rims_email_context
    from apps.core.notifications.models import Notification
    from apps.core.notifications.services import bulk_notify

    cfg = VarianceAlertConfig.objects.first()
    if cfg is None or not cfg.enabled:
        return 0
    threshold_ratio = float(cfg.percent_threshold) / 100.0
    flagged_variances = aggregate_variance_snapshot(threshold_ratio=threshold_ratio)
    finance_users = list(_finance_alerts_users())
    if finance_users:
        for v in flagged_variances:
            msg = (
                f'Variance alert on "{v.budget_name}" (award #{v.award_id}): spent '
                f"{v.spent_ratio*100:.1f}% of budget vs {v.elapsed_ratio*100:.1f}% of "
                f"project window — {v.direction}-running by {abs(v.gap)*100:.1f}%."
            )
            ctx = rims_email_context(
                subject=f"Variance alert: {v.budget_name}",
                headline="Budget variance threshold exceeded",
                action_path=f"/rims/finance/budgets/{v.budget_id}/",
                action_label="Open budget",
                preheader=(
                    "Spend pace differs from elapsed time by more than the configured threshold."
                ),
            )
            bulk_notify(
                finance_users,
                msg,
                verb=Notification.Verb.MEL_INDICATOR_OFF_TRACK,
                email_context=ctx,
            )
    flagged = len(flagged_variances)
    if flagged:
        logger.info("scan_budget_variances: flagged %s budget(s)", flagged)
    return flagged


@shared_task
def remind_payment_due(warn_days_before: int | None = None) -> int:
    """PRD §5.4 FRFM026 — alert when scheduled payments approach their due_on.

    Defaults to a 7-day window unless overridden by ReminderConfig (we reuse
    closeout_warn_days_before since they share semantics for "warn me N days
    early"; if a separate knob is needed it's a one-line follow-up).
    """
    from apps.rims.finance.models import PaymentSchedule
    from apps.rims.finance.notifications import notify_payment_due
    from apps.rims.operations.models import ReminderConfig

    cfg = ReminderConfig.objects.first()
    if warn_days_before is None:
        warn_days_before = cfg.closeout_warn_days_before if cfg else 7

    today = timezone.now().date()
    horizon = today + timedelta(days=warn_days_before)
    schedules = PaymentSchedule.objects.filter(
        paid=False,
        due_on__lte=horizon,
        due_on__gte=today,
    )
    sent = 0
    cutoff = timezone.now() - timedelta(days=1)
    for sch in schedules:
        if sch.last_reminded_at and sch.last_reminded_at >= cutoff:
            continue
        try:
            notify_payment_due(sch)
            sch.last_reminded_at = timezone.now()
            sch.save(update_fields=["last_reminded_at"])
            sent += 1
        except Exception as exc:  # noqa: BLE001
            logger.warning("remind_payment_due: failed for schedule=%s error=%s", sch.pk, exc)
    if sent:
        logger.info("remind_payment_due: sent %s reminder(s)", sent)
    return sent


@shared_task
def recompute_reforecasts() -> int:
    """PRD §5.4 FRFM027 — monthly rolling-budget reforecast for every APPROVED
    or FINALIZED budget on file. Reforecasts are append-only, so re-running on
    the same budget creates a new snapshot rather than overwriting.

    Returns the number of reforecast rows appended.
    """
    from apps.rims.finance.models import Budget
    from apps.rims.finance.services_reforecast import generate_reforecast

    count = 0
    qs = Budget.objects.filter(
        status__in=[Budget.Status.APPROVED, Budget.Status.FINALIZED]
    )
    for budget in qs.iterator(chunk_size=50):
        try:
            generate_reforecast(budget, actor=None)
            count += 1
        except Exception:  # pragma: no cover - keep the beat alive on per-budget errors
            logger.exception(
                "Reforecast failed for budget %s; continuing.", budget.pk
            )
    return count


@shared_task
def process_payment_retry_queue() -> int:
    """PRD §5.4 FRFM043 — sweep queued PaymentRetry rows whose next_retry_at
    is due. The beat does not attempt the actual bank retry (that's manual
    today); it just bumps QUEUED → RETRYING so dashboards reflect activity
    and finance staff see the queue moving.

    Returns the number of retries advanced.
    """
    from apps.rims.finance.models import PaymentRetry

    now = timezone.now()
    qs = PaymentRetry.objects.filter(
        status=PaymentRetry.Status.QUEUED,
    ).filter(next_retry_at__isnull=True) | PaymentRetry.objects.filter(
        status=PaymentRetry.Status.QUEUED,
        next_retry_at__lte=now,
    )
    count = 0
    for retry in qs.iterator():
        retry.status = PaymentRetry.Status.RETRYING
        retry.attempts = (retry.attempts or 0) + 1
        retry.last_attempt_at = now
        retry.save(update_fields=["status", "attempts", "last_attempt_at", "updated_at"])
        count += 1
    return count


@shared_task
def revalidate_gl_entries() -> int:
    """PRD §5.4 FRFM046 / FRFM049 — daily re-validation of recent GL entries.

    Also posts any DisbursementExecution rows whose gl_posting_status is
    "pending" (signal failed earlier). Returns the number of entries
    revalidated plus the number of executions newly posted.
    """
    from apps.rims.finance.models import DisbursementExecution, GLEntry
    from apps.rims.finance.services_gl import post_gl_entries, revalidate_gl_entry

    cutoff = timezone.now() - timedelta(days=14)
    count = 0
    for entry in GLEntry.objects.filter(
        validation_status=GLEntry.ValidationStatus.POSTED,
        posted_at__gte=cutoff,
    ).iterator():
        revalidate_gl_entry(entry)
        count += 1
    # Recover pending executions
    pending = DisbursementExecution.objects.filter(
        status=DisbursementExecution.Status.EXECUTED,
        gl_posting_status="pending",
    )
    for exe in pending.iterator():
        try:
            post_gl_entries(exe)
            DisbursementExecution.objects.filter(pk=exe.pk).update(
                gl_posting_status="posted"
            )
            count += 1
        except Exception:  # pragma: no cover
            logger.exception("Recovery posting failed for execution %s", exe.pk)
    return count
