"""A small HTTP server, on a thread.

The daemon's job model suits polling and does not suit a listener: a server
blocks, and a blocking job would stop everything else coming round. So this runs
alongside on its own thread.

`http.server` from the standard library rather than a framework. This serves two
health paths and, later, one webhook — and adding FastAPI plus its dependency
tree to route three URLs would be a large cost for a small job, in a codebase
whose small auditable surface has been worth keeping.

Routes are registered rather than hard-coded, so Phase 3's webhook goes into this
server instead of needing a second one on another port.
"""

from __future__ import annotations

import contextlib
import json
import logging
import threading
from collections.abc import Callable
from dataclasses import dataclass, field
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib.parse import parse_qs

logger = logging.getLogger(__name__)

#: Largest request body accepted, in bytes.
#:
#: THE MEDIA IS IN THE BODY, WHICH THIS DID NOT ASSUME
#:
#: This was two megabytes, on the reasoning that a WhatsApp webhook carries JSON
#: with a media *id* rather than the media itself — so a legitimate body is
#: kilobytes, and two megabytes was generous by three orders of magnitude.
#:
#: That is not how this deployment is wired. Its Evolution instance is
#: configured `webhookBase64: true`, so the image arrives inside the webhook,
#: base64 encoded, which inflates it by about a third. An ordinary phone
#: photograph of a CNIC is three to four megabytes, and every one of them was
#: refused.
#:
#: Found on 2026-08-04: a document sat unprocessed for 45 minutes while
#: Evolution re-posted the same 4,415,811 byte body every few minutes and this
#: server turned it away each time. Nothing was broken and nothing said so — the
#: record simply never moved.
#:
#: Thirty-two megabytes, chosen against the other end rather than picked: the
#: CMS accepts documents up to 20 MB, base64 carries that to ~27 MB, and the
#: remainder is JSON overhead. A body the CMS would refuse anyway is still
#: refused here, cheaply, before it is read.
#:
#: SINCE 2026-08-15 this is a backstop and not a document-size gate. The
#: instance is configured `base64: false`, so media no longer travels inside the
#: webhook at all — the payload is the message, a few kilobytes of it, and the
#: file is fetched afterwards through getBase64FromMediaMessage. Nothing in this
#: codebase ever read the inline copy: `_media_of` keys off `url`/`mediaKey` and
#: `download_media` fetches separately, so the base64 was a third of a document
#: posted, refused and thrown away.
#:
#: What that cost: roughly fifteen client documents between 31 July and 9 August
#: were refused here and silently lost, the largest at 48.8 MB. The refusal is
#: header-only by necessity — nothing is known about the sender or the file at
#: that point, so nobody could be told. Removing the inflation removes the
#: situation rather than trying to report it.
MAX_BODY_BYTES = 32 * 1024 * 1024

#: How much of a rejected body to read and discard before answering.
#:
#: Enough that an ordinary overage — somebody attaching a large file by mistake —
#: finishes writing and receives a clean 413 rather than a connection reset.
#: Beyond this the connection is closed, because politeness stops being worth
#: unbounded reading.
DRAIN_LIMIT_BYTES = 8 * 1024 * 1024

#: How long to spend draining before giving up.
#:
#: Needed alongside the byte cap: a caller claiming ten gigabytes and sending
#: four bytes would otherwise hold a thread until the socket timed out, turning
#: a defence against memory exhaustion into a way to exhaust threads.
DRAIN_TIMEOUT_SECONDS = 2.0


@dataclass(frozen=True, slots=True)
class Request:
    """What a route is given.

    Headers travel here rather than through a thread-local or a shared
    attribute. A webhook has to verify a signature carried in a header, and
    smuggling it in through a side channel would be a hidden dependency between
    the server and one route — and, on a threaded server, a race where one
    request verifies against another's signature.
    """

    body: bytes = b""
    headers: dict[str, str] = field(default_factory=dict)
    query: dict[str, str] = field(default_factory=dict)

    def header(self, name: str) -> str | None:
        return self.headers.get(name.lower())


#: A route: request in, (status, payload) out. A str payload is sent verbatim,
#: which Meta's verification handshake needs — it wants the bare challenge back.
Route = Callable[[Request], tuple[int, object]]


class _Handler(BaseHTTPRequestHandler):
    routes: dict[tuple[str, str], Route] = {}

    # ── Requests ──────────────────────────────────────────────────────────

    def do_GET(self) -> None:  # noqa: N802 - BaseHTTPRequestHandler's contract
        self._dispatch("GET")

    def do_POST(self) -> None:  # noqa: N802
        self._dispatch("POST")

    def _dispatch(self, method: str) -> None:
        raw_path, _, raw_query = self.path.partition("?")
        path = raw_path.rstrip("/") or "/"
        route = self.routes.get((method, path))

        if route is None:
            self._respond(404, {"error": "not_found"})

            return

        try:
            length = int(self.headers.get("Content-Length") or 0)
        except ValueError:
            self._respond(400, {"error": "bad_content_length"})

            return

        if length > MAX_BODY_BYTES:
            # Refused on the header alone. `rfile.read(n)` allocates n, so an
            # attacker claiming ten gigabytes exhausts memory without ever
            # sending ten gigabytes — and this endpoint faces the internet.
            logger.warning("Refused a %d byte request body.", length)
            self._drain(length)
            self._respond(413, {"error": "body_too_large"})

            return

        try:
            request = Request(
                body=self.rfile.read(length) if length else b"",
                headers={k.lower(): v for k, v in self.headers.items()},
                query={k: v[0] for k, v in parse_qs(raw_query).items()},
            )
            status, payload = route(request)
        except Exception as exc:  # noqa: BLE001
            # A failing route must not kill the server thread and take health
            # reporting down with it — which would turn one broken endpoint into
            # a container restart loop.
            logger.exception("Request to %s failed: %s", path, exc)
            self._respond(500, {"error": "internal_error"})

            return

        self._respond(status, payload)

    def _drain(self, length: int) -> None:
        """Read and discard a bounded amount of a rejected body.

        Without this, the server answers 413 and closes while the client is
        still writing — and the client sees a connection reset instead of the
        response. A reset carries no information and reads as a network fault,
        which is a poor way to tell somebody their file is too big.

        Chunked and capped: memory stays at one chunk however much was claimed,
        and beyond the cap the connection is simply closed. Draining without a
        limit would reintroduce exactly the problem the 413 exists to prevent.
        """
        remaining = min(length, DRAIN_LIMIT_BYTES)
        original = self.connection.gettimeout()

        try:
            # Bounded by TIME as well as by bytes, and both are needed. A caller
            # that claims ten gigabytes and sends four bytes would otherwise hold
            # this thread until the socket timed out — turning a defence against
            # memory exhaustion into a way to exhaust threads instead.
            self.connection.settimeout(DRAIN_TIMEOUT_SECONDS)

            while remaining > 0:
                chunk = self.rfile.read(min(remaining, 65536))

                if not chunk:
                    break

                remaining -= len(chunk)
        except OSError:
            # Timed out, or the client gave up first. Either way there is nothing
            # left to be polite about, and the 413 still goes out.
            pass
        finally:
            with contextlib.suppress(OSError):
                self.connection.settimeout(original)

    def _respond(self, status: int, payload: object) -> None:
        # A plain string goes back verbatim: Meta's verification handshake wants
        # the challenge itself, not the challenge wrapped in JSON.
        if isinstance(payload, str):
            body = payload.encode("utf-8")
            content_type = "text/plain; charset=utf-8"
        else:
            body = json.dumps(payload).encode("utf-8")
            content_type = "application/json"

        self.send_response(status)
        self.send_header("Content-Type", content_type)
        self.send_header("Content-Length", str(len(body)))
        # This is a machine endpoint; a cached health answer is a wrong one.
        self.send_header("Cache-Control", "no-store")
        self.send_header("X-Content-Type-Options", "nosniff")
        self.end_headers()
        self.wfile.write(body)

    def log_message(self, format: str, *args) -> None:  # noqa: A002
        """Quiet by default.

        BaseHTTPRequestHandler writes a line to stderr per request. A health
        probe every ten seconds would fill the container's logs with noise and
        bury anything worth reading.
        """
        logger.debug("%s - %s", self.address_string(), format % args)


class HttpServer:
    """Serves registered routes until stopped."""

    def __init__(self, host: str = "127.0.0.1", port: int = 8080) -> None:
        self._host = host
        self._port = port
        self._routes: dict[tuple[str, str], Route] = {}
        self._server: ThreadingHTTPServer | None = None
        self._thread: threading.Thread | None = None

    def route(self, method: str, path: str, handler: Route) -> None:
        self._routes[(method.upper(), path.rstrip("/") or "/")] = handler

    @property
    def port(self) -> int:
        """The bound port, which is the requested one unless 0 was asked for."""
        return self._server.server_address[1] if self._server else self._port

    def start(self) -> None:
        handler = type("Handler", (_Handler,), {"routes": self._routes})

        self._server = ThreadingHTTPServer((self._host, self._port), handler)
        # Daemon thread: a stuck request must never keep the process alive after
        # the loop has stopped and everything else has shut down.
        self._thread = threading.Thread(target=self._server.serve_forever, daemon=True)
        self._thread.start()

        logger.info("Listening on http://%s:%d", self._host, self.port)

    def stop(self) -> None:
        if self._server is not None:
            self._server.shutdown()
            self._server.server_close()
            self._server = None

        if self._thread is not None:
            self._thread.join(timeout=5)
            self._thread = None


def health_routes(server: HttpServer, container, daemon) -> None:
    """Register the two probes.

    Separate paths because they answer different questions and the orchestrator
    does different things with them. One combined endpoint would mean an
    unreachable CMS restarting the container.
    """
    from app.runtime import health

    def live(_request: Request) -> tuple[int, object]:
        report = health.liveness(daemon)

        return (200 if report.is_ready else 503), report.to_dict()

    def ready(_request: Request) -> tuple[int, object]:
        report = health.readiness(container, daemon)

        return (200 if report.is_ready else 503), report.to_dict()

    server.route("GET", "/health/live", live)
    server.route("GET", "/health/ready", ready)


def metrics_route(server: HttpServer, container=None) -> None:
    """Expose telemetry for scraping.

    Prometheus text format, which any monitoring tool reads and which needs no
    library to emit. Bound to the same loopback interface as the health probes —
    this is operational detail about a system handling client documents, and it
    is not a public page.
    """
    from app.runtime.metrics import METRICS

    def scrape(_request: Request) -> tuple[int, object]:
        rendered = METRICS.render() + _build_info()

        if container is not None:
            # Durable figures come from the database rather than from counters,
            # which reset on restart. Appended at scrape time so a monitoring
            # tool sees one surface instead of two.
            rendered += _workflow_series(container)
            rendered += _pending_pdf_series(container)
            rendered += _intake_outbox_series(container)

        return 200, rendered

    server.route("GET", "/metrics", scrape)


def _pending_pdf_series(container) -> str:
    """How many PDFs are waiting to be told whose they are.

    A queue holding client documents needs a number somebody can see. Everything
    in it has been taken out of WhatsApp and is this system's responsibility, and
    a depth that climbs and never falls is the signal that identifiers have
    stopped arriving — which otherwise looks exactly like a quiet evening.
    """
    try:
        waiting = len(container.pending_pdfs)
    except Exception:  # noqa: BLE001
        # Telemetry must never be the reason a scrape fails.
        return ""

    return (
        "# HELP taxpilot_pending_pdfs PDFs downloaded and awaiting an identifying message.\n"
        "# TYPE taxpilot_pending_pdfs gauge\n"
        f"taxpilot_pending_pdfs {waiting}\n"
    )


def _intake_outbox_series(container) -> str:
    """Documents that arrived and the CMS has not been told about (ADR-0010).

    The number that makes holding them acceptable. Anything counted here is a
    document somebody sent which does not yet appear anywhere in the CMS, so a
    non-zero value is not a backlog to watch — it is the firm being unable to see
    their own post, and it should be alerted on rather than trended.

    Zero is the normal reading. It is only ever non-zero while the CMS is
    unreachable.
    """
    try:
        held = len(container.intake_outbox)
        capacity = container.intake_outbox.capacity
    except Exception:  # noqa: BLE001
        # Telemetry must never be the reason a scrape fails.
        return ""

    return (
        "# HELP taxpilot_intake_outbox Arrivals the CMS has not acknowledged yet.\n"
        "# TYPE taxpilot_intake_outbox gauge\n"
        f"taxpilot_intake_outbox {held}\n"
        "# HELP taxpilot_intake_outbox_capacity Where the agent stops accepting documents.\n"
        "# TYPE taxpilot_intake_outbox_capacity gauge\n"
        f"taxpilot_intake_outbox_capacity {capacity}\n"
    )


def _build_info() -> str:
    """Which release is running, as a labelled constant.

    The Prometheus convention for build metadata: a gauge fixed at 1 whose
    labels carry the information. It exists so a dashboard can group by version
    and an alert can say "this started at 1.2.0" — a version reported only in a
    health response cannot be correlated with a spike in anything.
    """
    from app.release.version import VERSION

    return (
        "# HELP taxpilot_build_info The running release, as a label.\n"
        "# TYPE taxpilot_build_info gauge\n"
        f'taxpilot_build_info{{version="{VERSION}"}} 1\n'
    )


def _workflow_series(container) -> str:
    """Workflow history, rendered alongside the live counters."""
    from app.runtime.workflow_stats import workflow_stats

    try:
        stats = workflow_stats(container.connect)
    except Exception:  # noqa: BLE001
        # A scrape that fails takes the monitoring of everything else with it.
        # Losing one series beats losing the endpoint.
        logger.warning("Could not read workflow statistics for /metrics.", exc_info=True)

        return ""

    if stats is None:
        return ""

    lines = [
        "# HELP taxpilot_workflows Workflow runs over the last 7 days, by state.",
        "# TYPE taxpilot_workflows gauge",
        f'taxpilot_workflows{{state="completed"}} {stats.completed}',
        f'taxpilot_workflows{{state="failed"}} {stats.failed}',
        f'taxpilot_workflows{{state="awaiting"}} {stats.awaiting}',
        "# HELP taxpilot_workflows_retried Runs where a step was attempted more than once.",
        "# TYPE taxpilot_workflows_retried gauge",
        f"taxpilot_workflows_retried {stats.retried}",
    ]

    if stats.median_seconds is not None:
        lines += [
            "# HELP taxpilot_workflow_median_seconds Median time from start to finish.",
            "# TYPE taxpilot_workflow_median_seconds gauge",
            f"taxpilot_workflow_median_seconds {stats.median_seconds}",
        ]

    return "\n".join(lines) + "\n"
