"""The polling surface for the three background catalogue jobs.

POST /catalog/pair is the vehicle for these tests -- pair, sync and enrich all
go through the same shared runner in app.api.catalog, so exercising one of
them exercises the wiring the other two share.

Every request goes through one TestClient held open for the whole module, so
every HTTP call and every direct jobs.* call runs on the same event loop. The
Redis connection pool in app.services.infra.redis ties a connection to the
loop that created it; a fresh loop per call (the default, portal-per-request
TestClient behaviour) would hand a later call a connection born on an already
closed loop. tests/integration/test_jobs.py hit the same problem and pins its
own module to one loop for the same reason.
"""
import time

import pytest
from fastapi.testclient import TestClient

from app.main import app
from app.services.infra import jobs

TENANT = "org_job_routes_test"
POLL_TIMEOUT_SECONDS = 5


@pytest.fixture(scope="module")
def client():
    with TestClient(app) as c:
        yield c


@pytest.fixture(autouse=True)
def clean(client):
    yield
    client.portal.call(jobs.clear_tenant, TENANT)


def _poll_until_terminal(client, job_id):
    deadline = time.monotonic() + POLL_TIMEOUT_SECONDS
    while time.monotonic() < deadline:
        state = client.get(f"/catalog/jobs/{job_id}").json()
        if state["status"] in ("done", "failed", "lost"):
            return state
        time.sleep(0.02)
    raise AssertionError(f"job {job_id} did not reach a terminal state in time")


def test_pair_returns_202_and_a_job_id(client, monkeypatch):
    from app.api import catalog
    monkeypatch.setattr(catalog, "pair_tenant", lambda t, progress=None: {
        "products": 1, "pairs": 0, "servable": 0, "queued": 0, "by_type": {}})

    resp = client.post("/catalog/pair", data={"tenant_id": TENANT})
    assert resp.status_code == 202
    body = resp.json()
    assert body["status"] == "queued"
    assert body["job_id"]


def test_the_job_can_then_be_polled(client, monkeypatch):
    from app.api import catalog
    report = {"products": 218, "pairs": 900, "servable": 700, "queued": 200,
              "by_type": {"similar": 400}}
    monkeypatch.setattr(catalog, "pair_tenant", lambda t, progress=None: report)

    job_id = client.post("/catalog/pair", data={"tenant_id": TENANT}).json()["job_id"]
    state = _poll_until_terminal(client, job_id)

    assert state["job_id"] == job_id
    assert state["status"] == "done"
    assert state["percent"] == 100
    assert state["result"] == report


def test_wait_true_returns_the_report_inline(client, monkeypatch):
    from app.api import catalog
    monkeypatch.setattr(catalog, "pair_tenant", lambda t: {
        "products": 5, "pairs": 2, "servable": 2, "queued": 0, "by_type": {}})

    resp = client.post("/catalog/pair", data={"tenant_id": TENANT, "wait": "true"})
    assert resp.status_code == 200
    assert resp.json()["pairs"] == 2
    # No job was created for a synchronous call.
    assert client.get("/catalog/jobs", params={"tenant_id": TENANT}).json() == []


def test_a_second_job_of_the_same_kind_is_409_with_the_running_id(client, monkeypatch):
    # Held open with an event rather than a sleep: jobs.create yields to the
    # event loop while it talks to the database, which was long enough for a
    # sleep-based first job to finish before the second request arrived. The
    # lock was never the problem; the timing assumption was.
    import threading

    from app.api import catalog

    release = threading.Event()

    def blocked_pair(t, progress=None):
        release.wait(timeout=10)
        return {"products": 0, "pairs": 0, "servable": 0, "queued": 0, "by_type": {}}

    monkeypatch.setattr(catalog, "pair_tenant", blocked_pair)

    first = client.post("/catalog/pair", data={"tenant_id": TENANT})
    assert first.status_code == 202
    first_job_id = first.json()["job_id"]

    try:
        second = client.post("/catalog/pair", data={"tenant_id": TENANT})
        assert second.status_code == 409
        assert second.json()["detail"]["job_id"] == first_job_id
    finally:
        release.set()

    _poll_until_terminal(client, first_job_id)


def test_an_unknown_job_id_is_404(client):
    resp = client.get("/catalog/jobs/job_does_not_exist")
    assert resp.status_code == 404


def test_recent_jobs_are_listed(client, monkeypatch):
    from app.api import catalog
    monkeypatch.setattr(catalog, "pair_tenant", lambda t, progress=None: {
        "products": 0, "pairs": 0, "servable": 0, "queued": 0, "by_type": {}})

    job_id = client.post("/catalog/pair", data={"tenant_id": TENANT}).json()["job_id"]
    _poll_until_terminal(client, job_id)

    listed = client.get("/catalog/jobs", params={"tenant_id": TENANT}).json()
    assert job_id in [j["job_id"] for j in listed]


def test_a_failed_job_exposes_only_the_exception_class(client, monkeypatch):
    from app.api import catalog

    def boom(t, progress=None):
        raise RuntimeError("embedding service unavailable")

    monkeypatch.setattr(catalog, "pair_tenant", boom)

    job_id = client.post("/catalog/pair", data={"tenant_id": TENANT}).json()["job_id"]
    state = _poll_until_terminal(client, job_id)

    assert state["status"] == "failed"
    assert state["error"] == "RuntimeError"
    assert "unavailable" not in str(state)


# --- the two job routes that live in endpoints.py, not catalog.py ------------
# This module's docstring says exercising /catalog/pair exercises the wiring the
# other two share. They share the RUNNER, not their module's imports: a missing
# `import asyncio` in endpoints.py reached a live run because nothing here ever
# called those two routes.

@pytest.mark.parametrize("path,kind", [("/catalog/enrich", "enrich"),
                                       ("/catalog/sync", "sync")])
def test_the_endpoints_py_jobs_start_and_return_a_job_id(client, path, kind,
                                                         monkeypatch):
    from app.api import endpoints

    async def fake_create(tenant_id, job_kind):
        assert job_kind == kind
        return f"job_{kind}"

    started = {}

    def fake_task(coro):
        coro.close()
        started["ran"] = True

    monkeypatch.setattr(endpoints.jobs, "create", fake_create)
    monkeypatch.setattr(endpoints.asyncio, "create_task", fake_task)
    # /catalog/sync refuses with 404 before creating a job when the tenant has
    # no connected source, which is correct: a job that could only fail should
    # never be started.
    monkeypatch.setattr(endpoints, "get_sources",
                        lambda t: [{"kind": "http_api", "external_ref": "x",
                                    "config": {}}])

    resp = client.post(path, data={"tenant_id": TENANT})
    assert resp.status_code == 202
    assert resp.json()["job_id"] == f"job_{kind}"
    assert started["ran"]


@pytest.mark.parametrize("path", ["/catalog/enrich", "/catalog/sync"])
def test_the_endpoints_py_jobs_refuse_a_duplicate(client, path, monkeypatch):
    from app.api import endpoints

    async def already(tenant_id, kind):
        raise jobs.JobAlreadyRunning("job_existing")

    monkeypatch.setattr(endpoints.jobs, "create", already)
    monkeypatch.setattr(endpoints, "get_sources",
                        lambda t: [{"kind": "http_api", "external_ref": "x",
                                    "config": {}}])
    resp = client.post(path, data={"tenant_id": TENANT})
    assert resp.status_code == 409
    assert resp.json()["detail"]["job_id"] == "job_existing"


# --- the one-button build ----------------------------------------------------

def test_build_returns_a_job_id(client, monkeypatch):
    from app.api import catalog

    async def fake_create(tenant_id, kind):
        assert kind == "build"
        return "job_build"

    monkeypatch.setattr(catalog.jobs, "create", fake_create)
    monkeypatch.setattr(catalog.asyncio, "create_task",
                        lambda coro: coro.close())

    resp = client.post("/catalog/build", data={"tenant_id": TENANT})
    assert resp.status_code == 202
    assert resp.json()["job_id"] == "job_build"


def test_build_refuses_a_second_run(client, monkeypatch):
    # The lock covers the whole build, so a second press cannot start pairing
    # while the first is still enriching.
    from app.api import catalog

    async def already(tenant_id, kind):
        raise jobs.JobAlreadyRunning("job_running")

    monkeypatch.setattr(catalog.jobs, "create", already)
    resp = client.post("/catalog/build", data={"tenant_id": TENANT})
    assert resp.status_code == 409
    assert resp.json()["detail"]["job_id"] == "job_running"


def test_build_can_run_inline(client, monkeypatch):
    from app.api import catalog
    monkeypatch.setattr(catalog, "build_catalog",
                        lambda t, progress=None, force=False:
                        {"products": 212, "pairs": 1247, "stages": {}})

    resp = client.post("/catalog/build",
                       data={"tenant_id": TENANT, "wait": "true"})
    assert resp.status_code == 200
    assert resp.json()["pairs"] == 1247


# --- the job store being unreachable -----------------------------------------

def test_an_unreachable_job_store_is_503_not_500(client, monkeypatch):
    """An unreachable job store surfaced as an opaque 500 on /catalog/build,
    which reads as "this endpoint is broken" and sends people looking in the
    wrong place. 503 says "try again", which is what is actually true.
    """
    from psycopg2 import OperationalError

    from app.api import catalog

    async def unreachable(tenant_id, kind):
        raise OperationalError("could not connect to server")

    monkeypatch.setattr(catalog.jobs, "create", unreachable)

    resp = client.post("/catalog/build", data={"tenant_id": TENANT})
    assert resp.status_code == 503
    assert resp.json()["detail"]["error"] == "job_store_unavailable"
    # The class name is useful; the message can carry connection details.
    assert resp.json()["detail"]["cause"] == "OperationalError"


def test_a_dropped_socket_is_also_503(client, monkeypatch):
    from psycopg2 import OperationalError

    from app.api import catalog

    async def dropped(tenant_id, kind):
        raise OperationalError("server closed the connection unexpectedly")

    monkeypatch.setattr(catalog.jobs, "create", dropped)
    resp = client.post("/catalog/build", data={"tenant_id": TENANT})
    assert resp.status_code == 503


def test_already_running_is_still_409_not_503(client, monkeypatch):
    # The new handler must not swallow the lock signal.
    from app.api import catalog

    async def already(tenant_id, kind):
        raise jobs.JobAlreadyRunning("job_running")

    monkeypatch.setattr(catalog.jobs, "create", already)
    resp = client.post("/catalog/build", data={"tenant_id": TENANT})
    assert resp.status_code == 409


def test_the_job_store_lives_in_postgres_not_redis():
    """Redis was the first choice and the deployment could not reach it. A
    control plane that fails whenever its own store is unreachable is worse
    than one sharing a database already proven to work."""
    import inspect

    from app.services.infra import jobs

    # The docstring explains why the store moved, so check what the module
    # actually imports rather than whether the word appears.
    imports = [line for line in inspect.getsource(jobs).splitlines()
               if line.startswith(("import ", "from "))]
    assert not [line for line in imports if "redis" in line.lower()]
    assert any("get_master_db_connection" in line for line in imports)
