File size: 5,057 Bytes
e26316b
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
"""El trabajo síncrono y caro no puede correr en el bucle de eventos.

`interpretar` es `async`, pero sus dos operaciones más caras son SÍNCRONAS: la recuperación RAG
(embedding bge-m3 + LanceDB + cross-encoder sobre hasta `rag_candidatos` filas, segundos en
cpu-basic) y todo SQLite/scrypt en los endpoints de auth. Llamadas directamente desde una
corrutina, bloquean el proceso ENTERO mientras duran: ni health check, ni el login de otro
veterinario, ni una ingesta del puente. La concurrencia efectiva era 1 en una app que se
anuncia multiusuario.

Estas pruebas miden la propiedad, no la implementación: si alguien sustituye un
`await asyncio.to_thread(...)` por la llamada directa, vuelven a fallar. Los márgenes son
holgados a propósito —miden concurrencia, no latencia— para que no parpadeen en CI.
"""

from __future__ import annotations

import asyncio
import time

import pytest

from app.ai import service
from app.ai.base import ErrorModelo  # noqa: F401  (documenta la superficie que se falsea)
from app.schemas import InterpretacionClinica, PeticionInterpretacion

BLOQUEO_S = 0.4

PETICION = {
    "paciente": {"especie": "canino"},
    "hallazgos": [
        {
            "clave": "hct",
            "nombre": "Hematocrito",
            "valor": 22.0,
            "unidad": "%",
            "direccion": "bajo",
            "gravedad": "grave",
        }
    ],
    "patrones": [{"nombre": "Anemia", "descripcion": "…", "gravedad": "grave"}],
    "imagenes": [],
}


class ClienteInstantaneo:
    """Cliente de modelo que responde sin coste: aquí se mide la recuperación, no la generación."""

    nombre = "medgemma-hf"
    prosa = True
    modelo = "hf-space"

    async def interpretar(self, *_a, **_k):
        return InterpretacionClinica(interpretacion="ok " * 20, requiere_derivacion=True)


@pytest.fixture
def rag_lento(monkeypatch):
    """Retriever síncrono y lento, como el real: `time.sleep` bloquea el hilo que lo ejecute."""

    def _recuperar_bloqueante(*_a, **_k):
        time.sleep(BLOQUEO_S)
        return []

    for nombre in ("recuperar", "recuperar_multi"):
        monkeypatch.setattr(service, nombre, _recuperar_bloqueante)
    monkeypatch.setattr(service, "_crear_cliente", lambda *_: ClienteInstantaneo())


async def _contar_latidos(tarea: asyncio.Task, intervalo: float = 0.005) -> int:
    """Cuántas veces consigue despertarse el bucle mientras `tarea` está en curso.

    Es la medición directa de «¿puede el servidor atender a alguien más?». Con el trabajo
    bloqueante en el bucle, el contador se queda en ~0.
    """
    latidos = 0
    while not tarea.done():
        await asyncio.sleep(intervalo)
        latidos += 1
    return latidos


async def test_la_recuperacion_no_congela_el_bucle(rag_lento):
    """Durante una interpretación, el bucle sigue despertándose para atender otras cosas."""
    tarea = asyncio.create_task(service.interpretar(PeticionInterpretacion.model_validate(PETICION)))
    latidos = await _contar_latidos(tarea)
    await tarea

    # Con to_thread caben ~80 latidos de 5 ms en 0.4 s; en el bucle serían 0 o 1.
    assert latidos > 10, (
        f"sólo {latidos} latidos durante la recuperación: el bucle estuvo bloqueado, "
        "la recuperación volvió a ejecutarse sin asyncio.to_thread"
    )


async def test_dos_interpretaciones_se_solapan(rag_lento):
    """Dos peticiones concurrentes no se serializan: comparten el tiempo de recuperación."""
    inicio = time.perf_counter()
    await asyncio.gather(
        service.interpretar(PeticionInterpretacion.model_validate(PETICION)),
        service.interpretar(PeticionInterpretacion.model_validate(PETICION)),
    )
    transcurrido = time.perf_counter() - inicio

    # Serializadas costarían >= 2*BLOQUEO_S; solapadas, algo más de BLOQUEO_S.
    assert transcurrido < BLOQUEO_S * 1.7, (
        f"{transcurrido:.2f}s para dos interpretaciones de {BLOQUEO_S}s: se serializaron"
    )


async def test_el_alta_no_congela_el_bucle(alta_abierta):
    """scrypt (n=2**14) es caro A PROPÓSITO; en el bucle, cada alta congela el servicio."""
    import httpx

    from app import db
    from app.main import app

    # ASGITransport no dispara el lifespan (TestClient sí), así que la tabla no existiría.
    db.inicializar_db()

    transporte = httpx.ASGITransport(app=app)
    async with httpx.AsyncClient(transport=transporte, base_url="http://test") as cliente:
        tarea = asyncio.create_task(
            cliente.post(
                "/api/auth/registro",
                json={
                    "nombre": "Hilo",
                    "apellido": "Vet",
                    "email": "hilo@example.com",
                    "password": "clave-segura-1",
                },
            )
        )
        latidos = await _contar_latidos(tarea, intervalo=0.002)
        resp = await tarea

    assert resp.status_code == 200, resp.text
    assert latidos > 3, (
        f"sólo {latidos} latidos durante el alta: scrypt volvió a correr en el bucle de eventos"
    )