"""Intake and settlement services.

ops-contract: the ops runbook and the nightly reconciliation job call
`adjustments.services.reconciliation_report(account_id)` and
`redrive_reconciliation(account_id)` from outside the app, so those two names
and signatures are fixed. The recovery harness also passes ``crash_after`` into
`deliver` and `issue_statement_run`; keep those stage names pointing at the same
moments they point at now.

`deliver` is the partner-facing entry point. It accepts a correction, fixes its
place in the account's acceptance order, and hands the economic work to a
worker so the caller is not blocked behind resolution.

Pricing arithmetic is not reimplemented here -- see `pricing/engine.py`.
"""

from __future__ import annotations

import os
from threading import Lock

from django.db import IntegrityError, transaction

from pricing.engine import PricingEngine

from adjustments import hydrate, idempotency, monitoring, persist
from adjustments.intake import AdjustmentIntake, InjectedCrash
from adjustments.models import (
    Account,
    Adjustment,
    AdjustmentRecord,
    IntakeCheckpoint,
    Invoice,
    StatementRun,
)
from adjustments.tasks import apply_correction


class RefusedError(ValueError):
    """The delivery violates a billing contract and was not accepted."""


_delivery_lock = Lock()
_run_lock = Lock()


def deliver(
    account_id: str,
    payload: dict,
    crash_after: str | None = None,
) -> Adjustment | None:
    with _delivery_lock:
        return _deliver(account_id, payload, crash_after)


@transaction.atomic
def _deliver(
    account_id: str,
    payload: dict,
    crash_after: str | None = None,
) -> Adjustment | None:
    """Accept one delivered correction.

    Re-delivery of an identifier we have already accepted is absorbed: the
    partner feed is at-least-once and must keep moving.
    """

    adjustment_id = payload["adjustment_id"]
    if idempotency.already_accepted(adjustment_id):
        existing = Adjustment.objects.filter(
            account__account_id=account_id,
            adjustment_id=adjustment_id,
        ).first()
        if existing is not None:
            checkpoint = IntakeCheckpoint.objects.filter(
                account=existing.account,
                adjustment_id=adjustment_id,
            ).first()
            if checkpoint is None or checkpoint.stage != "posted":
                apply_correction.delay(existing.pk)
        return None

    account, _ = Account.objects.get_or_create(
        account_id=account_id,
        defaults={
            "settlement_policy": os.environ.get(
                "INVOICE_CLOSED_POLICY", "credit_forward"
            )
        },
    )
    checkpoint, _ = IntakeCheckpoint.objects.get_or_create(
        account=account,
        adjustment_id=adjustment_id,
        defaults={"stage": "received"},
    )
    if crash_after == "received":
        raise InjectedCrash("received")
    value = hydrate.payload_adjustment(account_id, payload)
    if (
        value.subject_key is not None
        or value.replaces_adjustment_id is not None
        or value.withdraws_adjustment_id is not None
    ):
        try:
            AdjustmentIntake(PricingEngine()).validate(
                hydrate.load(account_id), value
            )
        except ValueError as error:
            monitoring.acceptance_refused(account_id, adjustment_id, str(error))
            raise RefusedError(str(error)) from error
    ordinal = idempotency.next_ordinal(account_id)
    promotion = value.promotion

    try:
        with transaction.atomic():
            adjustment = Adjustment.objects.create(
                adjustment_id=adjustment_id,
                account=account,
                period_key=payload["period_key"],
                kind=payload["kind"],
                effective_at=payload["effective_at"],
                received_at=payload["received_at"],
                acceptance_ordinal=ordinal,
                amount=payload.get("amount"),
                usage_key=payload.get("usage_key"),
                quantity=payload.get("quantity"),
                promotion_key=promotion.key if promotion else None,
                promotion_percent=promotion.percent if promotion else None,
                promotion_kind=promotion.kind if promotion else None,
                promotion_applies_to=(
                    list(promotion.applies_to_usage_keys or ())
                    if promotion
                    else []
                ),
                target_period_key=payload.get("target_period_key"),
                subject_key=payload.get("subject_key"),
                replaces_adjustment_id=payload.get("replaces_adjustment_id"),
                withdraws_adjustment_id=payload.get("withdraws_adjustment_id"),
            )
    except IntegrityError:
        existing = Adjustment.objects.get(adjustment_id=adjustment_id)
        existing_checkpoint = IntakeCheckpoint.objects.filter(
            account=existing.account,
            adjustment_id=adjustment_id,
        ).first()
        if existing_checkpoint is None or existing_checkpoint.stage != "posted":
            apply_correction.delay(existing.pk)
        return None
    idempotency.mark_accepted(adjustment_id)
    checkpoint.stage = "accepted"
    checkpoint.save(update_fields=["stage", "updated_at"])
    if crash_after == "accepted":
        raise InjectedCrash("accepted")

    # Resolution is the expensive part; a worker picks it up.
    apply_correction.delay(adjustment.pk)
    return adjustment


def issue_statement_run(
    account_id: str,
    operation_id: str,
    crash_after: str | None = None,
) -> StatementRun:
    with _run_lock:
        return _issue_statement_run(account_id, operation_id, crash_after)


def _issue_statement_run(
    account_id: str,
    operation_id: str,
    crash_after: str | None = None,
) -> StatementRun:
    """Start and issue a statement run for the account.

    A repeated `operation_id` returns the same run rather than issuing a new
    one.
    """

    idempotency.acquire_run_lock(account_id)
    try:
        domain = hydrate.load(account_id)
        intake = AdjustmentIntake(PricingEngine())
        started = intake.start_statement_run(domain, operation_id)
        persist.resolution(domain)
        if crash_after == "started":
            raise InjectedCrash("started")
        for adjustment in Adjustment.objects.filter(
            account__account_id=account_id,
            acceptance_ordinal__lte=started.cutoff_ordinal,
        ).order_by("acceptance_ordinal"):
            apply_correction(adjustment.pk)
        domain = hydrate.load(account_id)
        intake.issue_statement_run(domain, operation_id)
        persist.resolution(domain)
        return StatementRun.objects.get(run_id=started.run_id)
    finally:
        idempotency.release_run_lock(account_id)


def adopt_run_summaries(account_id: str, adoption_id: str) -> None:
    """Per-account rollout of run reconciliation summaries.

    Requested in the 2026-07-06 thread and specified in the attachment.
    """

    raise NotImplementedError("run summaries are not implemented yet")


def _current_ordinal(account: Account) -> int:
    latest = (
        Adjustment.objects.filter(account=account)
        .order_by("-acceptance_ordinal")
        .first()
    )
    return latest.acceptance_ordinal if latest and latest.acceptance_ordinal else 0


def amount_due(account_id: str, period_key: str) -> float:
    invoice = Invoice.objects.get(account__account_id=account_id, period_key=period_key)
    base = invoice.issued_total if invoice.issued_total is not None else invoice.total
    return base + sum(record.delta for record in invoice.records.all())
