"""Health probes.

The important tests here are about what liveness does NOT check. A liveness probe
that considers external services turns a ten-minute CMS outage into a restart
loop: in-flight work discarded, nothing fixed, recovery slowed. That mistake is
easy to make and expensive, so it is pinned from several directions.
"""

from __future__ import annotations

import json
import urllib.error
import urllib.request
from datetime import UTC, datetime, timedelta

import pytest

from app.runtime import health
from app.runtime.daemon import Daemon
from app.runtime.health import Component, Report, Status
from app.runtime.server import HttpServer, health_routes

BASE_ENV = {
    "TAXPILOT_CMS_URL": "https://cms.test",
    "TAXPILOT_AGENT_KEY": "tpa_test",
    "TAXPILOT_AGENT_SECRET": "s" * 64,
}


class FakeContainer:
    """A container whose every dependency can be made to fail."""

    def __init__(self, *, database=True, cms=True, cms_describes=None, ocr="paddleocr", memory=True, inbox=True,
                 permissions=("clients.read", "proposals.submit")):
        self._database, self._cms_ok, self._memory_ok = database, cms, memory
        self._cms_describes = cms_describes if cms_describes is not None else {
            "version": "1.1.4",
            "api_contract": "1.2",
            "minimum_agent_version": "0.1.0",
        }
        self._permissions = permissions
        self.settings = object()
        self.ocr = type("Engine", (), {"name": ocr})()
        self.cms = self._make_cms()
        self.memory = self._make_memory()
        self.inbox = self._make_inbox(inbox)

        # Empty is the healthy reading. Anything held here is a document that
        # arrived and the CMS has not been told about, so it is reported the
        # moment it is non-zero rather than trended (ADR-0010).
        from app.intake import InMemoryOutbox

        self.intake_outbox = InMemoryOutbox()

    @property
    def connect(self):
        if not self._database:
            return None

        outer = self

        class Cursor:
            def __enter__(self): return self
            def __exit__(self, *a): return False
            def execute(self, *a):
                if outer._database == "unreachable":
                    raise RuntimeError("connection refused")
            def fetchone(self): return (1,)

        class Connection:
            def __enter__(self): return self
            def __exit__(self, *a): return False
            def cursor(self): return Cursor()

        return Connection

    def _make_cms(self):
        outer = self

        class Cms:
            """Matches CmsClient's surface, including the version handshake.

            `handshake` rather than a bare `whoami` because that is what the
            real client exposes and what readiness calls — a double that is
            merely close enough proves the code works against the double.
            """

            def __init__(self):
                self.identity = None
                self.compatibility = None

            def whoami(self):
                if not outer._cms_ok:
                    raise RuntimeError("connection refused")

                return {
                    "name": "Intake",
                    "permissions": list(outer._permissions),
                    "cms": dict(outer._cms_describes),
                }

            def handshake(self):
                from app.api.compatibility import evaluate

                self.identity = self.whoami()
                self.compatibility = evaluate(self.identity)

                return self.compatibility

        return Cms()

    def _make_memory(self):
        outer = self

        class Memory:
            def recall(self, *a, **kw):
                if not outer._memory_ok:
                    raise RuntimeError("Postgres is down")

                return []

        return Memory()

    def _make_inbox(self, configured):
        from app.whatsapp.inbox import InboxFilter, SelfChatPolicy

        # "Configured" now means one thing only: the connected number is known.
        # There is nothing else the inbox can be missing.
        policy = SelfChatPolicy.for_owner("923001234567") if configured else SelfChatPolicy()

        return InboxFilter(policy)


class TestLivenessIgnoresTheWorldOutside:
    """The property that stops a dependency outage becoming a restart loop."""

    def test_a_ticking_loop_is_alive(self):
        daemon = Daemon()
        daemon.started_at = datetime.now(UTC)
        daemon.last_tick_at = datetime.now(UTC)

        assert health.liveness(daemon).is_ready

    def test_an_unreachable_cms_does_not_affect_liveness(self):
        daemon = Daemon()
        daemon.started_at = daemon.last_tick_at = datetime.now(UTC)

        # A CMS down for maintenance must not have the orchestrator kill a
        # perfectly healthy process over and over.
        assert health.liveness(daemon).is_ready
        assert health.readiness(FakeContainer(cms=False), daemon).is_ready is False

    def test_a_missing_database_does_not_affect_liveness(self):
        daemon = Daemon()
        daemon.started_at = daemon.last_tick_at = datetime.now(UTC)

        assert health.liveness(daemon).is_ready

    def test_liveness_reports_only_the_loop(self):
        daemon = Daemon()
        daemon.started_at = daemon.last_tick_at = datetime.now(UTC)

        assert [c.name for c in health.liveness(daemon).components] == ["loop"]

    def test_a_stalled_loop_is_not_alive(self):
        daemon = Daemon()
        daemon.started_at = datetime.now(UTC) - timedelta(minutes=10)
        daemon.last_tick_at = datetime.now(UTC) - timedelta(minutes=10)

        # A job hung on a socket with no timeout. The process exists, answers
        # signals, and does nothing — the one condition a restart genuinely fixes.
        assert health.liveness(daemon).is_ready is False

    def test_one_slow_pass_is_not_a_stall(self):
        daemon = Daemon()   # max_wait_seconds 5, tolerance 30
        daemon.started_at = datetime.now(UTC) - timedelta(minutes=1)
        daemon.last_tick_at = datetime.now(UTC) - timedelta(seconds=7)

        # A restart throws away in-flight work, so the bar for declaring death
        # is higher than one iteration running long.
        assert health.liveness(daemon).is_ready

    def test_a_daemon_shutting_down_is_still_alive(self):
        daemon = Daemon()
        daemon.started_at = datetime.now(UTC) - timedelta(minutes=10)
        daemon.last_tick_at = datetime.now(UTC) - timedelta(minutes=10)
        daemon.stop()

        # It was told to exit. Reporting that as a failure would have an
        # orchestrator restart something that is deliberately stopping.
        assert health.liveness(daemon).is_ready

    def test_a_daemon_that_has_not_started_is_not_yet_dead(self):
        assert health.liveness(Daemon()).is_ready


class TestReadinessChecksTheWorldOutside:
    def test_everything_working_is_ready(self):
        report = health.readiness(FakeContainer())

        assert report.status is Status.OK
        assert report.is_ready

    def test_a_cms_speaking_another_contract_takes_the_service_offline(self):
        """The whole readiness path, not just the component.

        The CMS updates itself over the air, so this is what happens when the
        installation moves to a major release under a running deployment: every
        CMS call is already being refused by the gate, and an orchestrator that
        went on routing here would be routing to a service that cannot file
        anything.
        """
        container = FakeContainer(cms_describes={"version": "2.0.0", "api_contract": "2.0"})

        report = health.readiness(container)

        assert report.is_ready is False
        assert _component(report, "compatibility").status is Status.FAILING

    def test_a_cms_that_predates_the_handshake_is_not_ready(self):
        """The reverse of what this asserted until the floor rose to 1.1.

        It used to say that refusing these would be an outage caused by the
        check meant to prevent one, and that was right while the agent could
        work against contract 1.0. Since ADR-0010 it cannot: every arriving
        document is registered with the CMS before anything else happens to it,
        and a pre-handshake CMS has no endpoint to register it with.

        Not ready is the honest answer, and an orchestrator holding work back
        from this deployment is doing the right thing. The alternative is
        routing documents to an agent that will hold every one of them and then
        begin refusing new ones.
        """
        container = FakeContainer(cms_describes={"version": "1.1.4"})

        report = health.readiness(container)

        assert not report.is_ready
        assert _component(report, "compatibility").status is Status.FAILING

    def test_the_running_version_is_reported(self):
        """What lets an installer's readiness probe tell one release from
        another — see docs/distribution.md."""
        from app.release.version import VERSION

        assert health.readiness(FakeContainer()).to_dict()["version"] == VERSION

    def test_an_unreachable_database_is_not_ready(self):
        report = health.readiness(FakeContainer(database="unreachable"))

        assert not report.is_ready
        assert _component(report, "database").status is Status.FAILING

    def test_no_database_configured_is_not_ready(self):
        assert not health.readiness(FakeContainer(database=False)).is_ready

    def test_an_unreachable_cms_is_not_ready(self):
        assert not health.readiness(FakeContainer(cms=False)).is_ready

    def test_a_cms_missing_grants_is_not_ready(self):
        report = health.readiness(FakeContainer(permissions=("clients.read",)))

        # Reachable but useless: every document read would have nowhere to go.
        assert not report.is_ready
        assert "proposals.submit" in _component(report, "cms").detail

    def test_every_component_is_reported_even_when_one_fails(self):
        report = health.readiness(FakeContainer(database=False, cms=False))

        # An operator wants the whole picture, not the first thing that broke.
        names = {c.name for c in report.components}
        assert {"configuration", "database", "cms", "ocr", "memory", "whatsapp"} <= names


class TestDegradedIsStillReady:
    def test_missing_ocr_does_not_take_the_service_offline(self):
        report = health.readiness(FakeContainer(ocr="null"))

        # Documents still reach a reviewer, unread. Failing readiness here would
        # remove a working service over a component that only improves it.
        assert report.is_ready
        assert report.status is Status.DEGRADED
        assert _component(report, "ocr").status is Status.DEGRADED

    def test_unresponsive_memory_is_degraded(self):
        report = health.readiness(FakeContainer(memory=False))

        assert report.is_ready
        assert _component(report, "memory").status is Status.DEGRADED

    def test_an_inbox_with_no_linked_account_is_degraded(self):
        report = health.readiness(FakeContainer(inbox=False))

        # It still finishes workflows already in flight; it simply accepts
        # nothing new, which is the only safe answer while nobody can say whose
        # messages these would be.
        assert report.is_ready
        assert _component(report, "whatsapp").status is Status.DEGRADED


class TestTheResponseLeaksNothing:
    def test_a_database_failure_does_not_name_the_host(self):
        body = json.dumps(health.readiness(FakeContainer(database="unreachable")).to_dict())

        # "connection to 10.0.0.4:5432 failed" is a map of the infrastructure for
        # whoever reads it, and this is the endpoint most likely to be exposed
        # by accident.
        assert "connection refused" not in body
        assert "5432" not in body

    def test_the_watched_numbers_are_masked(self):
        body = json.dumps(health.readiness(FakeContainer()).to_dict())

        assert "923001234567" not in body


class TestHttpEndpoints:
    @pytest.fixture
    def server(self):
        daemon = Daemon()
        daemon.started_at = daemon.last_tick_at = datetime.now(UTC)

        server = HttpServer("127.0.0.1", 0)   # 0: let the OS choose a free port
        health_routes(server, FakeContainer(), daemon)
        server.start()

        yield server

        server.stop()

    def test_liveness_answers_200(self, server):
        status, body = _get(server, "/health/live")

        assert status == 200
        assert body["status"] == "ok"

    def test_readiness_answers_200_when_everything_works(self, server):
        status, body = _get(server, "/health/ready")

        assert status == 200
        assert {c["name"] for c in body["components"]} >= {"database", "cms", "ocr"}

    def test_an_unknown_path_is_404(self, server):
        assert _get(server, "/nope")[0] == 404

    def test_readiness_answers_503_when_a_dependency_is_down(self):
        daemon = Daemon()
        daemon.started_at = daemon.last_tick_at = datetime.now(UTC)
        server = HttpServer("127.0.0.1", 0)
        health_routes(server, FakeContainer(cms=False), daemon)
        server.start()

        try:
            # 503, not 500: the service is not broken, it is not ready — and an
            # orchestrator withholds traffic rather than restarting.
            assert _get(server, "/health/ready")[0] == 503
        finally:
            server.stop()

    def test_liveness_still_answers_200_while_dependencies_are_down(self):
        daemon = Daemon()
        daemon.started_at = daemon.last_tick_at = datetime.now(UTC)
        server = HttpServer("127.0.0.1", 0)
        health_routes(server, FakeContainer(cms=False, database=False), daemon)
        server.start()

        try:
            # The whole point. Docker must not restart this.
            assert _get(server, "/health/live")[0] == 200
        finally:
            server.stop()

    def test_responses_are_not_cached(self, server):
        request = urllib.request.Request(f"http://127.0.0.1:{server.port}/health/live")

        with urllib.request.urlopen(request, timeout=5) as response:
            # A cached health answer is a wrong one.
            assert response.headers["Cache-Control"] == "no-store"

    def test_a_failing_route_does_not_kill_the_server(self, server):
        def explode(_body):
            raise RuntimeError("boom")

        server.route("GET", "/explode", explode)

        assert _get(server, "/explode")[0] == 500
        # One broken endpoint must not take health reporting down with it and
        # turn itself into a restart loop.
        assert _get(server, "/health/live")[0] == 200


class TestReportShape:
    def test_failing_beats_degraded(self):
        report = Report([
            Component("a", Status.DEGRADED),
            Component("b", Status.FAILING),
        ])

        assert report.status is Status.FAILING
        assert not report.is_ready

    def test_degraded_alone_is_ready(self):
        report = Report([Component("a", Status.DEGRADED)])

        assert report.status is Status.DEGRADED
        assert report.is_ready

    def test_an_empty_report_is_ok(self):
        assert Report().status is Status.OK


def _component(report: Report, name: str) -> Component:
    return next(c for c in report.components if c.name == name)


def _get(server: HttpServer, path: str) -> tuple[int, dict]:
    url = f"http://127.0.0.1:{server.port}{path}"

    try:
        with urllib.request.urlopen(url, timeout=5) as response:
            return response.status, json.loads(response.read())
    except urllib.error.HTTPError as exc:
        return exc.code, json.loads(exc.read())


class TestWhatsAppHealthTestsFlowNotLinkage:
    """An account being linked is not messages being able to flow.

    On 2026-08-03 the old check said "ok" through three outages — Evolution's
    process died twice and its stream died once — because it only ever looked
    at the inbox policy. The transport is the deployment's only way to receive
    a document; these pin the new judgement and its deliberate limits.
    """

    def _provider(self, state: str, detail: str = ""):
        from app.whatsapp.session import SessionState, SessionStatus

        class Provider:
            def session_status(self):
                return SessionStatus(SessionState(state), detail)

        return Provider()

    def test_a_connected_session_is_ok(self):
        container = FakeContainer()
        container.whatsapp = self._provider("connected")

        report = health.readiness(container)

        assert _component(report, "whatsapp").status is Status.OK

    def test_an_unreachable_provider_is_failing_and_not_ready(self):
        container = FakeContainer()

        class Broken:
            def session_status(self):
                raise RuntimeError("connection refused")

        container.whatsapp = Broken()

        report = health.readiness(container)

        assert _component(report, "whatsapp").status is Status.FAILING
        assert not report.is_ready

    def test_a_logged_out_session_is_failing(self):
        # The account was unlinked from the phone, or banned. Either way no
        # document can arrive, and "ok" here is the lie the old check told.
        container = FakeContainer()
        container.whatsapp = self._provider("logged_out")

        report = health.readiness(container)

        assert _component(report, "whatsapp").status is Status.FAILING

    def test_a_provider_that_cannot_report_keeps_the_old_judgement(self):
        # Meta, and every stub: "we cannot tell" must not condemn a working
        # deployment. The old policy-based answer stands for them.
        container = FakeContainer()
        container.whatsapp = object()

        report = health.readiness(container)

        assert _component(report, "whatsapp").status is Status.OK


class TestTheStatusReportCarriesTheFlowSignal:
    """`last_event_seconds`: when a message last proved the pipe works.

    Surfaced, not alerted on — silence is ambiguous, but a person reading
    "connected" beside "last event 4h ago" holds exactly the contradiction
    that found the dead-stream outage.
    """

    def _container(self, last_event_at):
        class Policy:
            is_configured = True

        class Inbox:
            policy = Policy()
            allowed = 3
            ignored = 1
            forwards_accepted = 0
            forwards_rejected = 0
            forwards_failed = 0

        Inbox.last_event_at = last_event_at

        class C:
            whatsapp = None
            inbox = Inbox()

        return C()

    def test_present_once_anything_has_arrived(self):
        from datetime import UTC, datetime, timedelta

        from app.runtime.status_report import _whatsapp

        block = _whatsapp(self._container(datetime.now(UTC) - timedelta(seconds=90)))

        assert 85 <= block["last_event_seconds"] <= 150

    def test_absent_before_the_first_event(self):
        from app.runtime.status_report import _whatsapp

        # "0 seconds ago" on a fresh start would be a claim, not a measurement.
        assert "last_event_seconds" not in _whatsapp(self._container(None))

    def test_an_ignored_message_still_stamps_the_pipe_as_alive(self):
        from app.whatsapp.inbox import InboxFilter, SelfChatPolicy
        from app.whatsapp.messages import InboundMessage, MessageType

        inbox = InboxFilter(SelfChatPolicy.for_owner("923001234567"))

        stranger = InboundMessage(
            provider_message_id="m1",
            sender="447700900000",
            type=MessageType.TEXT,
            text="hello",
            from_me=False,
        )

        assert inbox.permits(stranger) is False
        # Delivery is the fact being recorded; the verdict is separate.
        assert inbox.last_event_at is not None
