"""The only way this platform reaches TaxPilot CMS (ADR-0005).

No database driver, no browser automation, no direct file access — every
capability arrives through the signed Agent API. That is what keeps business
rules in one place, in one language, behind one authorization seam.
"""

from __future__ import annotations

import base64
import json
import logging
import ssl
import time
import urllib.error
import urllib.parse
import urllib.request
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any

from app.api.compatibility import Verdict, evaluate
from app.config.settings import Settings
from app.runtime.metrics import METRICS
from app.security.signing import sign
from app.support.http import user_agent

logger = logging.getLogger(__name__)


#: How many times to wait out a rate limit before giving up.
_RATE_LIMIT_RETRIES = 4

#: The longest single pause. A daemon job is blocked while this waits, so a
#: server suggesting sixty seconds gets a shorter answer than it asked for.
_RATE_LIMIT_MAX_WAIT = 8.0


def _retry_after(headers) -> float | None:
    """Seconds the CMS asked us to wait, when it said."""
    raw = headers.get("Retry-After") if headers else None

    try:
        return float(raw) if raw is not None else None
    except (TypeError, ValueError):
        return None


class AgentApiError(RuntimeError):
    """A call the CMS refused, or could not be completed."""

    def __init__(self, message: str, *, status: int | None = None, reason: str | None = None):
        super().__init__(message)
        self.status = status
        # The CMS's machine-readable reason ("signature_invalid",
        # "permission_denied", …). Carried separately so callers can branch on
        # it without parsing prose.
        self.reason = reason


class RateLimited(AgentApiError):
    """The CMS asked this agent to slow down.

    Its own type because it is not a refusal. Everything else in this family
    means the call will not succeed; this one means it will, shortly — and the
    difference decides whether a caller retries or gives up on a document.
    """

    def __init__(self, message: str, *, retry_after: float | None = None):
        super().__init__(message, status=429, reason="rate_limited")
        self.retry_after = retry_after


class PermissionDenied(AgentApiError):
    """This agent was not granted the permission the endpoint requires."""


class NotFound(AgentApiError):
    """No such record — or none this agent may see. Deliberately the same answer."""


#: The only call allowed through when the CMS is known to be incompatible.
#:
#: It is the handshake itself. Blocking it would make the incompatibility
#: permanent for the lifetime of the process — nothing could ever re-ask, so a
#: CMS that was rolled back to a working version would never be noticed.
EXEMPT_FROM_COMPATIBILITY = frozenset({"whoami"})


class IncompatibleCms(AgentApiError):
    """The CMS speaks a contract this release does not.

    Raised in place of making the request, not in response to one. An
    AgentApiError so callers that already handle "the CMS refused" do not crash
    on it; retrying is harmless, because the gate is local and the message says
    what to do.
    """


@dataclass(slots=True)
class CmsClient:
    """A signed client for one TaxPilot installation."""

    settings: Settings

    _identity: dict[str, Any] | None = field(default=None, repr=False)
    """What whoami said, kept so one call answers both questions at boot."""

    _verdict: Verdict | None = field(default=None, repr=False)
    """What the last handshake concluded, or None if there has not been one.

    None never blocks. The gate refuses on knowledge, not on ignorance: a client
    that has not handshaken yet is not known to be incompatible, and refusing
    would break every caller that has no reason to care.
    """

    # ── The handshake ─────────────────────────────────────────────────────

    def handshake(self) -> Verdict:
        """Ask the CMS what it is, and decide whether to work with it.

        Called at boot, and again periodically — the CMS updates itself over the
        air, so the installation this is attached to can change contract under a
        running process. A verdict reached only once at start-up would go stale
        at exactly the moment it mattered.
        """
        try:
            identity = self.whoami()
        except AgentApiError:
            # Unreachable or unauthenticated is not incompatible. Leaving the
            # previous verdict in place keeps a network blip from silently
            # clearing a real incompatibility, and preflight reports the
            # connection failure on its own.
            raise

        self._identity = identity
        self._verdict = evaluate(identity)

        return self._verdict

    @property
    def compatibility(self) -> Verdict | None:
        """The last verdict, for health reporting. None until the first handshake."""
        return self._verdict

    @property
    def identity(self) -> dict[str, Any] | None:
        """What the last handshake learned. None until the first handshake."""
        return self._identity

    # ── Endpoints ─────────────────────────────────────────────────────────

    def whoami(self) -> dict[str, Any]:
        """What this agent is and may do.

        Needs no permission by design, which makes it the right call for a
        deployment to verify its own configuration on startup.
        """
        return self._get("whoami")["data"]

    def report_status(self, status: dict[str, Any]) -> dict[str, Any]:
        """Tell the CMS how this deployment is doing.

        The counterpart to whoami. The CMS cannot discover any of it — ADR-0005
        makes this client the only channel and the agent the caller — so
        component health, WhatsApp state and workflow counts arrive only because
        they are sent.

        Needs no permission, for whoami's reason: the agent most worth hearing
        from is a misconfigured one, and that is exactly the agent a grant
        requirement would silence.
        """
        return self._request("POST", "status", body=status)

    def claim_command(self) -> dict[str, Any] | None:
        """Take the next command the CMS has queued, if any.

        None on an empty queue, which is almost every poll. The CMS answers 204
        for that, so there is no body to interpret and no ambiguity between
        "nothing waiting" and "something malformed".
        """
        response = self._request("POST", "commands/claim")

        return (response or {}).get("data") or None

    def report_command(
        self,
        command_id: str,
        status: str,
        result: dict[str, Any] | None = None,
        error: str = "",
        retryable: bool = True,
    ) -> dict[str, Any]:
        """Say what happened.

        Safe to repeat: the CMS ignores a result for a command it has already
        decided, so a retry after a dropped response is ordinary rather than a
        way to corrupt the record.
        """
        return self._request(
            "POST",
            f"commands/{command_id}/result",
            body={
                "status": status,
                "result": result or {},
                "error": error,
                "retryable": retryable,
            },
        )

    def document_types(self) -> list[dict[str, str]]:
        """The document vocabulary *this* installation uses.

        Fetched rather than hard-coded: a classifier that invents a type the CMS
        does not recognise produces a failed upload instead of a filed document.
        """
        return self._get("document-types")["data"]

    def tax_years(self) -> list[str]:
        return self._get("tax-years")["data"]

    def search_clients(self, **criteria: Any) -> dict[str, Any]:
        """Find clients. The first question when a document arrives: whose is this?"""
        return self._get("clients", params={k: v for k, v in criteria.items() if v is not None})

    def get_client(self, client_id: int) -> dict[str, Any]:
        return self._get(f"clients/{client_id}")["data"]

    def identify_proposal(self, reference: str, identifier: str) -> dict[str, Any]:
        """Answer "whose document is this?" for a proposal nobody could match.

        Addressed by the DOC- reference rather than the id, because that is what
        the sender was asked to type back and what their reply carries.

        The CMS decides whether the identifier is usable — it owns the client
        list, and it is the only side that can refuse one belonging to two people
        or to nobody. This does not file anything: the proposal stays pending and
        moves into the Approval Queue, where a human still approves the write
        (ADR-0008). An answer typed into a phone can direct a document, never
        commit it unseen.
        """
        return self._request("POST", f"proposals/{reference}/identify", body={"identifier": identifier})["data"]

    def submit_proposal(self, submission: dict[str, Any], path: str | Path | None = None) -> dict[str, Any]:
        """Ask a human for permission to change a client's record (ADR-0008).

        The ordinary path for anything that writes. This does not file anything;
        it puts the request in front of a reviewer, and the CMS performs the work
        itself once they approve.

        Submitting is idempotent on the caller's key: a retry returns the same
        proposal rather than creating a second one, which matters because a
        network timeout tells you nothing about whether the request arrived.
        """
        payload = dict(submission)

        if path is not None:
            source = Path(path)

            if not source.exists():
                raise AgentApiError(f"No document at {source}.")

            payload["file_base64"] = base64.b64encode(source.read_bytes()).decode("ascii")
            payload.setdefault("file_name", source.name)

        return self._request("POST", "proposals", body=payload)["data"]

    def register_intake(
        self,
        submission: dict[str, Any],
        path: str | Path | None = None,
    ) -> dict[str, Any]:
        """Tell the CMS a document exists, before knowing anything about it (ADR-0010).

        Called earlier than every other method here. The rest report a conclusion;
        this reports an arrival, and it runs before classification, before reading,
        before there is any opinion to submit — because the gap between "bytes
        arrived" and "the CMS was told" is the gap documents used to vanish into.

        Idempotent on (source, provider_message_id): WhatsApp redelivers, and this
        is retried after a restart. A repeat returns the record that already
        exists, so calling it twice is safe and calling it again after a crash is
        the intended recovery.

        `failure_reason` registers something this agent refused to process at all —
        the wrong file type, too large, not what it claimed to be. Those still get
        a record, because a client who sent something deserves better than silence.
        """
        payload = dict(submission)

        if path is not None:
            source = Path(path)

            if not source.exists():
                raise AgentApiError(f"No document at {source}.")

            payload["file_base64"] = base64.b64encode(source.read_bytes()).decode("ascii")
            payload.setdefault("original_filename", source.name)

        return self._request("POST", "intake", body=payload)["data"]["intake"]

    def update_intake(self, reference: str, **fields: Any) -> dict[str, Any]:
        """Move a document between states, or say where its workflow has got to.

        Pushing a state the record is already in is a no-op rather than an error,
        so an agent that restarted mid-flight can re-announce where it thinks it is
        without needing to know what the CMS already heard.

        A move that is not allowed raises, carrying the state the CMS believes the
        record is in — which is the information needed to reconcile, rather than to
        retry blindly against a record that has moved on.
        """
        body = {key: value for key, value in fields.items() if value is not None}

        return self._request("PATCH", f"intake/{reference}/status", body=body)["data"]["intake"]

    def get_intake(self, reference: str) -> dict[str, Any]:
        return self._get(f"intake/{reference}")["data"]["intake"]

    def pending_intake(
        self,
        status: str | None = None,
        per_page: int | None = None,
    ) -> tuple[list[dict[str, Any]], bool]:
        """What the CMS believes is unfinished, for this agent.

        Filtered by state, because the caller that matters is looking for one.
        Asking for everything open returned the *oldest* page of it, so a
        document stranded mid-read behind a backlog was never in the window and
        never recovered — a recovery path that worked on an empty installation
        and failed silently on a busy one.

        Returns the records and whether the answer was cut short, so a caller can
        tell "that is all of them" from "that is as many as you asked for".
        """
        params: dict[str, Any] = {}

        if status:
            params["status"] = status
        if per_page:
            params["per_page"] = per_page

        payload = self._get("intake/pending", params=params or None)["data"]

        return list(payload.get("intake") or []), bool(payload.get("truncated"))

    def get_proposal(self, proposal_id: int) -> dict[str, Any]:
        return self._get(f"proposals/{proposal_id}")["data"]

    def decided_proposals(self, since: str | int | None = None) -> list[dict[str, Any]]:
        """Everything decided since we last looked.

        One call rather than one per outstanding workflow: an agent restarting
        needs to catch up on decisions made in its absence, and it does not
        necessarily remember what it was waiting for.
        """
        params: dict[str, Any] = {
            "status": "approved,executed,failed,rejected,changes_requested",
        }

        if since is not None:
            params["since"] = since

        return self._get("proposals", params=params).get("data", [])

    def upload_document(
        self,
        *,
        client_id: int,
        path: str | Path,
        document_type: str,
        title: str | None = None,
        tax_year: str | None = None,
        notes: str | None = None,
        source: str | None = None,
    ) -> dict[str, Any]:
        """File a document against a client — the one call that writes.

        The file is base64'd into the JSON body rather than sent as multipart.
        The signature covers the raw body, and multipart boundaries are generated
        per client, so a multipart body is not something both sides can reproduce
        byte-for-byte. The 33% size cost buys a signature that actually verifies.

        Reached only after the workflow engine has recorded a human's approval;
        nothing here re-checks that, because a check the caller could skip is not
        a control. The gate is ExecutionPolicy, and the CMS's own
        ``documents.write`` grant is the backstop.
        """
        source_path = Path(path)

        if not source_path.exists():
            raise AgentApiError(f"No document at {source_path}.")

        payload: dict[str, Any] = {
            "client_id": client_id,
            "document_type": document_type,
            "file_name": source_path.name,
            "file_base64": base64.b64encode(source_path.read_bytes()).decode("ascii"),
        }

        # Omitted rather than sent as null: an explicit null would overwrite a
        # value the CMS may already hold.
        for key, value in (("title", title), ("tax_year", tax_year),
                           ("notes", notes), ("source", source)):
            if value:
                payload[key] = value

        return self._request("POST", "documents", body=payload)["data"]

    # ── Transport ─────────────────────────────────────────────────────────

    def _refuse_if_incompatible(self, path: str) -> None:
        """One gate, on the one way out.

        ADR-0005 makes this client the only channel to the CMS, which makes
        `_request` the only place this check has to exist. Putting it in each
        caller instead would mean the one that gets forgotten is the one that
        writes.
        """
        if path.strip("/") in EXEMPT_FROM_COMPATIBILITY:
            return

        verdict = self._verdict

        if verdict is None or verdict.is_usable:
            return

        raise IncompatibleCms(f"Refusing to call the CMS: {verdict.detail}")

    def _get(self, path: str, params: dict[str, Any] | None = None) -> dict[str, Any]:
        return self._request("GET", path, params=params)

    def _send(
        self,
        method: str,
        path: str,
        params: dict[str, Any] | None = None,
        body: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        self._refuse_if_incompatible(path)

        raw_body = json.dumps(body, separators=(",", ":")) if body is not None else ""

        # The signature covers the routed path WITHOUT the query string, because
        # that is what PHP's $request->path() returns. Signing the query would
        # produce a signature the CMS can never reproduce.
        signed_path = f"api/agent/v1/{path.lstrip('/')}"

        url = f"{self.settings.agent_api_root}/{path.lstrip('/')}"
        if params:
            url = f"{url}?{urllib.parse.urlencode(params, doseq=True)}"

        headers = sign(
            self.settings.agent_api_secret,
            method,
            signed_path,
            raw_body,
            api_key=self.settings.agent_api_key,
        ).as_dict()
        headers["Accept"] = "application/json"

        # Not cosmetic. A CMS on shared hosting sits behind a firewall that
        # judges the client by this header, and urllib's default is one it
        # resets the connection on — before Laravel is reached, so the
        # signature, the key and the URL are all irrelevant to the outcome.
        # Measured on Life Associate's cPanel host: without this, every call
        # failed with "remote end closed connection without response".
        headers["User-Agent"] = user_agent()

        if raw_body:
            headers["Content-Type"] = "application/json"

        request = urllib.request.Request(  # noqa: S310 - fixed scheme, validated in Settings
            url,
            data=raw_body.encode("utf-8") if raw_body else None,
            headers=headers,
            method=method.upper(),
        )

        context = None if self.settings.verify_tls else ssl._create_unverified_context()  # noqa: S323

        # Labelled by endpoint rather than full path: `clients/42` and
        # `clients/77` are the same operation, and one series per client id
        # would be an unbounded number of series — the classic way to bring
        # down a monitoring system.
        endpoint = path.split("/")[0] or "root"

        try:
            with METRICS.time("taxpilot_cms_request_seconds",
                              "Time spent calling the CMS.", endpoint=endpoint):
                with urllib.request.urlopen(  # noqa: S310
                    request, timeout=self.settings.request_timeout, context=context
                ) as response:
                    decoded = self._decode(response.read(), response.status)

            METRICS.counter("taxpilot_cms_requests_total", "Calls to the CMS.",
                            endpoint=endpoint, outcome="ok")

            return decoded
        except urllib.error.HTTPError as exc:
            payload = self._safe_decode(exc.read())
            reason = payload.get("error") if isinstance(payload, dict) else None

            if exc.code == 429:
                METRICS.counter("taxpilot_cms_requests_total", "Calls to the CMS.",
                                endpoint=endpoint, outcome="rate_limited")

                # Not a refusal. The CMS is telling us to slow down, and the
                # caller is the only party that can — so this is raised as its
                # own type and retried by the wrapper rather than handed to a
                # workflow as a failure it can do nothing about.
                raise RateLimited(
                    "The CMS is rate limiting this agent.",
                    retry_after=_retry_after(exc.headers),
                ) from exc

            METRICS.counter("taxpilot_cms_requests_total", "Calls to the CMS.",
                            endpoint=endpoint, outcome="refused")
            raise self._error_for(exc.code, reason) from exc
        except urllib.error.URLError as exc:
            METRICS.counter("taxpilot_cms_requests_total", "Calls to the CMS.",
                            endpoint=endpoint, outcome="unreachable")
            # The CMS being unreachable is an operational fact, not a bug in the
            # caller — surfaced as our own error type so agents handle one family.
            raise AgentApiError(f"TaxPilot CMS is unreachable: {exc.reason}") from exc

    def _request(
        self,
        method: str,
        path: str,
        params: dict[str, Any] | None = None,
        body: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        """Send, and wait rather than fail when the CMS says to slow down.

        A rate limit is an instruction, not a refusal — the CMS is asking for a
        pause, and the caller is the only party able to give one. Treating 429
        as an error made a burst of documents look like a broken pipeline:
        registrations fell into the outbox, status pushes went stale, and
        proposal submissions failed outright, all while the CMS was healthy and
        simply asking for less.

        Found by the stress harness at 100 documents against a 120-per-minute
        limit, which is roughly twenty-four documents a minute sustained — well
        within a morning's forwarding, and exactly what the outbox draining after
        an outage looks like.

        Bounded on purpose. This blocks a daemon job, so it waits a few seconds
        rather than however long a server suggests, and gives up in a way the
        caller already knows how to handle.
        """
        for attempt in range(_RATE_LIMIT_RETRIES + 1):
            try:
                return self._send(method, path, params, body)
            except RateLimited as limited:
                if attempt == _RATE_LIMIT_RETRIES:
                    raise AgentApiError(
                        "The CMS is rate limiting this agent and did not let up.",
                        status=429,
                        reason="rate_limited",
                    ) from limited

                delay = min(limited.retry_after or (2 ** attempt), _RATE_LIMIT_MAX_WAIT)
                logger.info(
                    "The CMS asked this agent to slow down; waiting %.1fs before retrying %s.",
                    delay, path,
                )
                time.sleep(delay)

        raise AgentApiError("Unreachable.")  # pragma: no cover

    def _decode(self, raw: bytes, status: int) -> dict[str, Any]:
        """204 carries no body, and that is a success rather than a malformed reply.

        The command queue answers 204 on an empty queue — which is almost every
        poll. Without this the agent logged a warning every five seconds saying
        the CMS had "refused" a request it had in fact answered correctly, and a
        genuinely broken queue would have been indistinguishable from an idle
        one.
        """
        if status == 204 or not raw.strip():
            return {"ok": True}

        payload = self._safe_decode(raw)

        if not isinstance(payload, dict):
            raise AgentApiError(f"Malformed response from the CMS (HTTP {status}).", status=status)

        if not payload.get("ok", False):
            raise self._error_for(status, payload.get("error"))

        return payload

    @staticmethod
    def _safe_decode(raw: bytes) -> Any:
        try:
            return json.loads(raw.decode("utf-8"))
        except (ValueError, UnicodeDecodeError):
            # An HTML error page rather than JSON usually means the request never
            # reached the Agent API at all — a proxy, a maintenance page, a
            # wrong base URL.
            return {}

    @staticmethod
    def _error_for(status: int, reason: str | None) -> AgentApiError:
        message = f"The CMS refused the request ({reason or 'unknown'}, HTTP {status})."

        if status == 403:
            return PermissionDenied(message, status=status, reason=reason)
        if status == 404:
            return NotFound(message, status=status, reason=reason)

        return AgentApiError(message, status=status, reason=reason)
