Spaces:
Sleeping
Sacar del bucle de eventos la recuperación RAG, SQLite y scrypt
Browse filesARCHITECTURE_REVIEW §3.1. `interpretar` es `async`, pero sus dos operaciones más
caras eran SÍNCRONAS y corrían en el bucle de eventos, así que mientras una
interpretación recuperaba, el proceso entero no contestaba a NADA: 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.
- `ai/service.py`: `recuperar` / `recuperar_multi` a un hilo. Es lo más caro del
request —embedding bge-m3, búsqueda LanceDB y cross-encoder sobre hasta
`rag_candidatos` filas—, del orden de segundos en cpu-basic.
- `routers/auth.py`: SQLite y scrypt de login y alta. scrypt (n=2**14) es caro a
propósito: es justo el trabajo que no puede vivir en el bucle.
- `routers/lab.py`: la escritura de `lab_persistir`, que corre por cada muestra
con el analizador enviando en ráfaga (120/minute).
`security/authz.py` no toca la BD (la sesión es una cookie firmada), así que las
peticiones autenticadas ya no pasaban por SQLite. `main.py` se deja como está: su
carga corre en el lifespan, antes de que haya tráfico que bloquear.
Las pruebas miden la PROPIEDAD, no la implementación: cuentan cuántas veces
consigue despertarse el bucle mientras hay trabajo en curso, y comprueban que dos
interpretaciones concurrentes se solapan en vez de serializarse. Verificado que
fallan sin el arreglo: 1 latido en vez de ~80, y 0.81s para dos recuperaciones de
0.4s. Márgenes holgados a propósito —miden concurrencia, no latencia— para que no
parpadeen en CI.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
- backend/app/ai/service.py +10 -2
- backend/app/routers/auth.py +14 -7
- backend/app/routers/lab.py +6 -1
- backend/tests/test_bucle_no_bloqueante.py +139 -0
|
@@ -7,6 +7,7 @@ Un reintento ante fallo de validación; si persiste, error tipado (nunca texto c
|
|
| 7 |
|
| 8 |
from __future__ import annotations
|
| 9 |
|
|
|
|
| 10 |
import logging
|
| 11 |
|
| 12 |
from ..config import obtener_config
|
|
@@ -208,13 +209,20 @@ async def interpretar(pet: PeticionInterpretacion) -> RespuestaInterpretacion:
|
|
| 208 |
# consulta concatenada).
|
| 209 |
nombres_patrones = [p.nombre for p in pet.patrones]
|
| 210 |
nombres_hallazgos = [h.nombre for h in pet.hallazgos]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 211 |
if cfg.rag_multiconsulta:
|
| 212 |
-
fragmentos =
|
|
|
|
| 213 |
construir_consultas(nombres_patrones, nombres_hallazgos),
|
| 214 |
especie=pet.paciente.especie,
|
| 215 |
)
|
| 216 |
else:
|
| 217 |
-
fragmentos =
|
|
|
|
| 218 |
construir_consulta(nombres_patrones, nombres_hallazgos),
|
| 219 |
especie=pet.paciente.especie,
|
| 220 |
)
|
|
|
|
| 7 |
|
| 8 |
from __future__ import annotations
|
| 9 |
|
| 10 |
+
import asyncio
|
| 11 |
import logging
|
| 12 |
|
| 13 |
from ..config import obtener_config
|
|
|
|
| 209 |
# consulta concatenada).
|
| 210 |
nombres_patrones = [p.nombre for p in pet.patrones]
|
| 211 |
nombres_hallazgos = [h.nombre for h in pet.hallazgos]
|
| 212 |
+
# `asyncio.to_thread` y no llamada directa: la recuperación es SÍNCRONA y cara —embedding
|
| 213 |
+
# bge-m3, búsqueda en LanceDB y cross-encoder bge-reranker-v2-m3 sobre hasta
|
| 214 |
+
# `rag_candidatos` filas—, del orden de segundos en cpu-basic. Ejecutada en el bucle de
|
| 215 |
+
# eventos dejaba el proceso entero sordo mientras durase: ni health check, ni el login de
|
| 216 |
+
# otro veterinario, ni una ingesta del puente. La concurrencia efectiva era 1.
|
| 217 |
if cfg.rag_multiconsulta:
|
| 218 |
+
fragmentos = await asyncio.to_thread(
|
| 219 |
+
recuperar_multi,
|
| 220 |
construir_consultas(nombres_patrones, nombres_hallazgos),
|
| 221 |
especie=pet.paciente.especie,
|
| 222 |
)
|
| 223 |
else:
|
| 224 |
+
fragmentos = await asyncio.to_thread(
|
| 225 |
+
recuperar,
|
| 226 |
construir_consulta(nombres_patrones, nombres_hallazgos),
|
| 227 |
especie=pet.paciente.especie,
|
| 228 |
)
|
|
@@ -9,6 +9,8 @@ Mejoras de seguridad frente a auth.php:
|
|
| 9 |
|
| 10 |
from __future__ import annotations
|
| 11 |
|
|
|
|
|
|
|
| 12 |
from fastapi import APIRouter, Depends, HTTPException, Request, Response, status
|
| 13 |
from pydantic import BaseModel, EmailStr, Field
|
| 14 |
|
|
@@ -75,17 +77,21 @@ async def estado(request: Request) -> dict:
|
|
| 75 |
@router.post("/auth/login")
|
| 76 |
@limiter.limit(obtener_config().limite_login)
|
| 77 |
async def login(request: Request, body: LoginBody, response: Response) -> dict:
|
|
|
|
|
|
|
|
|
|
|
|
|
| 78 |
ip = request.client.host if request.client else "?"
|
| 79 |
-
if
|
| 80 |
raise HTTPException(status.HTTP_429_TOO_MANY_REQUESTS, "Demasiados intentos. Espera unos minutos.")
|
| 81 |
|
| 82 |
-
usuario =
|
| 83 |
-
if not usuario or not
|
| 84 |
-
|
| 85 |
# Mensaje genérico: no revela si el email existe.
|
| 86 |
raise HTTPException(status.HTTP_401_UNAUTHORIZED, "Email o contraseña incorrectos.")
|
| 87 |
|
| 88 |
-
|
| 89 |
csrf = _emitir_sesion(response, usuario["email"], usuario["nombre"])
|
| 90 |
return {"ok": True, "nombre": usuario["nombre"], "csrf": csrf}
|
| 91 |
|
|
@@ -101,9 +107,10 @@ async def registro(request: Request, body: RegistroBody, response: Response) ->
|
|
| 101 |
status.HTTP_403_FORBIDDEN,
|
| 102 |
"El alta de cuentas está restringida. Solicita acceso al administrador.",
|
| 103 |
)
|
| 104 |
-
if
|
| 105 |
raise HTTPException(status.HTTP_409_CONFLICT, "Ya existe una cuenta con ese email.")
|
| 106 |
-
crear_usuario
|
|
|
|
| 107 |
csrf = _emitir_sesion(response, body.email, body.nombre)
|
| 108 |
return {"ok": True, "nombre": body.nombre, "csrf": csrf}
|
| 109 |
|
|
|
|
| 9 |
|
| 10 |
from __future__ import annotations
|
| 11 |
|
| 12 |
+
import asyncio
|
| 13 |
+
|
| 14 |
from fastapi import APIRouter, Depends, HTTPException, Request, Response, status
|
| 15 |
from pydantic import BaseModel, EmailStr, Field
|
| 16 |
|
|
|
|
| 77 |
@router.post("/auth/login")
|
| 78 |
@limiter.limit(obtener_config().limite_login)
|
| 79 |
async def login(request: Request, body: LoginBody, response: Response) -> dict:
|
| 80 |
+
# Todo lo que toca SQLite o scrypt va a un hilo: son llamadas SÍNCRONAS dentro de un
|
| 81 |
+
# endpoint `async`, así que en el bucle de eventos bloqueaban el proceso entero. scrypt es
|
| 82 |
+
# además caro A PROPÓSITO (n=2**14, decenas de ms): es justo el trabajo que no puede vivir
|
| 83 |
+
# en el bucle, y el login es el endpoint que más veces lo ejecuta.
|
| 84 |
ip = request.client.host if request.client else "?"
|
| 85 |
+
if await asyncio.to_thread(intentos_recientes, body.email, ip, _VENTANA_THROTTLE_S) >= _MAX_INTENTOS:
|
| 86 |
raise HTTPException(status.HTTP_429_TOO_MANY_REQUESTS, "Demasiados intentos. Espera unos minutos.")
|
| 87 |
|
| 88 |
+
usuario = await asyncio.to_thread(buscar_usuario, body.email)
|
| 89 |
+
if not usuario or not await asyncio.to_thread(verificar_password, body.password, usuario["password"]):
|
| 90 |
+
await asyncio.to_thread(registrar_intento, body.email, ip)
|
| 91 |
# Mensaje genérico: no revela si el email existe.
|
| 92 |
raise HTTPException(status.HTTP_401_UNAUTHORIZED, "Email o contraseña incorrectos.")
|
| 93 |
|
| 94 |
+
await asyncio.to_thread(limpiar_intentos, body.email)
|
| 95 |
csrf = _emitir_sesion(response, usuario["email"], usuario["nombre"])
|
| 96 |
return {"ok": True, "nombre": usuario["nombre"], "csrf": csrf}
|
| 97 |
|
|
|
|
| 107 |
status.HTTP_403_FORBIDDEN,
|
| 108 |
"El alta de cuentas está restringida. Solicita acceso al administrador.",
|
| 109 |
)
|
| 110 |
+
if await asyncio.to_thread(buscar_usuario, body.email):
|
| 111 |
raise HTTPException(status.HTTP_409_CONFLICT, "Ya existe una cuenta con ese email.")
|
| 112 |
+
# `crear_usuario` hashea con scrypt además de escribir: doble motivo para salir del bucle.
|
| 113 |
+
await asyncio.to_thread(crear_usuario, body.nombre, body.apellido, body.email, body.password)
|
| 114 |
csrf = _emitir_sesion(response, body.email, body.nombre)
|
| 115 |
return {"ok": True, "nombre": body.nombre, "csrf": csrf}
|
| 116 |
|
|
@@ -10,6 +10,8 @@ la auth); la consulta usa la sesión existente. El mapeo código→analito ocurr
|
|
| 10 |
|
| 11 |
from __future__ import annotations
|
| 12 |
|
|
|
|
|
|
|
| 13 |
from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
|
| 14 |
|
| 15 |
from .. import db
|
|
@@ -39,7 +41,10 @@ async def post_ingesta(
|
|
| 39 |
mapeado = mapear_resultado(cuerpo)
|
| 40 |
almacen.guardar(mapeado)
|
| 41 |
if obtener_config().lab_persistir:
|
| 42 |
-
|
|
|
|
|
|
|
|
|
|
| 43 |
mapeado.muestra_id.strip().lower(),
|
| 44 |
mapeado.momento.isoformat(),
|
| 45 |
mapeado.model_dump_json(),
|
|
|
|
| 10 |
|
| 11 |
from __future__ import annotations
|
| 12 |
|
| 13 |
+
import asyncio
|
| 14 |
+
|
| 15 |
from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
|
| 16 |
|
| 17 |
from .. import db
|
|
|
|
| 41 |
mapeado = mapear_resultado(cuerpo)
|
| 42 |
almacen.guardar(mapeado)
|
| 43 |
if obtener_config().lab_persistir:
|
| 44 |
+
# A un hilo como el resto de SQLite: el analizador manda en ráfaga
|
| 45 |
+
# (`limite_lab_ingesta` = 120/minute) y esto corre por cada muestra.
|
| 46 |
+
await asyncio.to_thread(
|
| 47 |
+
db.guardar_resultado_lab,
|
| 48 |
mapeado.muestra_id.strip().lower(),
|
| 49 |
mapeado.momento.isoformat(),
|
| 50 |
mapeado.model_dump_json(),
|
|
@@ -0,0 +1,139 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
"""El trabajo síncrono y caro no puede correr en el bucle de eventos.
|
| 2 |
+
|
| 3 |
+
`interpretar` es `async`, pero sus dos operaciones más caras son SÍNCRONAS: la recuperación RAG
|
| 4 |
+
(embedding bge-m3 + LanceDB + cross-encoder sobre hasta `rag_candidatos` filas, segundos en
|
| 5 |
+
cpu-basic) y todo SQLite/scrypt en los endpoints de auth. Llamadas directamente desde una
|
| 6 |
+
corrutina, bloquean el proceso ENTERO mientras duran: ni health check, ni el login de otro
|
| 7 |
+
veterinario, ni una ingesta del puente. La concurrencia efectiva era 1 en una app que se
|
| 8 |
+
anuncia multiusuario.
|
| 9 |
+
|
| 10 |
+
Estas pruebas miden la propiedad, no la implementación: si alguien sustituye un
|
| 11 |
+
`await asyncio.to_thread(...)` por la llamada directa, vuelven a fallar. Los márgenes son
|
| 12 |
+
holgados a propósito —miden concurrencia, no latencia— para que no parpadeen en CI.
|
| 13 |
+
"""
|
| 14 |
+
|
| 15 |
+
from __future__ import annotations
|
| 16 |
+
|
| 17 |
+
import asyncio
|
| 18 |
+
import time
|
| 19 |
+
|
| 20 |
+
import pytest
|
| 21 |
+
|
| 22 |
+
from app.ai import service
|
| 23 |
+
from app.ai.base import ErrorModelo # noqa: F401 (documenta la superficie que se falsea)
|
| 24 |
+
from app.schemas import InterpretacionClinica, PeticionInterpretacion
|
| 25 |
+
|
| 26 |
+
BLOQUEO_S = 0.4
|
| 27 |
+
|
| 28 |
+
PETICION = {
|
| 29 |
+
"paciente": {"especie": "canino"},
|
| 30 |
+
"hallazgos": [
|
| 31 |
+
{
|
| 32 |
+
"clave": "hct",
|
| 33 |
+
"nombre": "Hematocrito",
|
| 34 |
+
"valor": 22.0,
|
| 35 |
+
"unidad": "%",
|
| 36 |
+
"direccion": "bajo",
|
| 37 |
+
"gravedad": "grave",
|
| 38 |
+
}
|
| 39 |
+
],
|
| 40 |
+
"patrones": [{"nombre": "Anemia", "descripcion": "…", "gravedad": "grave"}],
|
| 41 |
+
"imagenes": [],
|
| 42 |
+
}
|
| 43 |
+
|
| 44 |
+
|
| 45 |
+
class ClienteInstantaneo:
|
| 46 |
+
"""Cliente de modelo que responde sin coste: aquí se mide la recuperación, no la generación."""
|
| 47 |
+
|
| 48 |
+
nombre = "medgemma-hf"
|
| 49 |
+
prosa = True
|
| 50 |
+
modelo = "hf-space"
|
| 51 |
+
|
| 52 |
+
async def interpretar(self, *_a, **_k):
|
| 53 |
+
return InterpretacionClinica(interpretacion="ok " * 20, requiere_derivacion=True)
|
| 54 |
+
|
| 55 |
+
|
| 56 |
+
@pytest.fixture
|
| 57 |
+
def rag_lento(monkeypatch):
|
| 58 |
+
"""Retriever síncrono y lento, como el real: `time.sleep` bloquea el hilo que lo ejecute."""
|
| 59 |
+
|
| 60 |
+
def _recuperar_bloqueante(*_a, **_k):
|
| 61 |
+
time.sleep(BLOQUEO_S)
|
| 62 |
+
return []
|
| 63 |
+
|
| 64 |
+
for nombre in ("recuperar", "recuperar_multi"):
|
| 65 |
+
monkeypatch.setattr(service, nombre, _recuperar_bloqueante)
|
| 66 |
+
monkeypatch.setattr(service, "_crear_cliente", lambda *_: ClienteInstantaneo())
|
| 67 |
+
|
| 68 |
+
|
| 69 |
+
async def _contar_latidos(tarea: asyncio.Task, intervalo: float = 0.005) -> int:
|
| 70 |
+
"""Cuántas veces consigue despertarse el bucle mientras `tarea` está en curso.
|
| 71 |
+
|
| 72 |
+
Es la medición directa de «¿puede el servidor atender a alguien más?». Con el trabajo
|
| 73 |
+
bloqueante en el bucle, el contador se queda en ~0.
|
| 74 |
+
"""
|
| 75 |
+
latidos = 0
|
| 76 |
+
while not tarea.done():
|
| 77 |
+
await asyncio.sleep(intervalo)
|
| 78 |
+
latidos += 1
|
| 79 |
+
return latidos
|
| 80 |
+
|
| 81 |
+
|
| 82 |
+
async def test_la_recuperacion_no_congela_el_bucle(rag_lento):
|
| 83 |
+
"""Durante una interpretación, el bucle sigue despertándose para atender otras cosas."""
|
| 84 |
+
tarea = asyncio.create_task(service.interpretar(PeticionInterpretacion.model_validate(PETICION)))
|
| 85 |
+
latidos = await _contar_latidos(tarea)
|
| 86 |
+
await tarea
|
| 87 |
+
|
| 88 |
+
# Con to_thread caben ~80 latidos de 5 ms en 0.4 s; en el bucle serían 0 o 1.
|
| 89 |
+
assert latidos > 10, (
|
| 90 |
+
f"sólo {latidos} latidos durante la recuperación: el bucle estuvo bloqueado, "
|
| 91 |
+
"la recuperación volvió a ejecutarse sin asyncio.to_thread"
|
| 92 |
+
)
|
| 93 |
+
|
| 94 |
+
|
| 95 |
+
async def test_dos_interpretaciones_se_solapan(rag_lento):
|
| 96 |
+
"""Dos peticiones concurrentes no se serializan: comparten el tiempo de recuperación."""
|
| 97 |
+
inicio = time.perf_counter()
|
| 98 |
+
await asyncio.gather(
|
| 99 |
+
service.interpretar(PeticionInterpretacion.model_validate(PETICION)),
|
| 100 |
+
service.interpretar(PeticionInterpretacion.model_validate(PETICION)),
|
| 101 |
+
)
|
| 102 |
+
transcurrido = time.perf_counter() - inicio
|
| 103 |
+
|
| 104 |
+
# Serializadas costarían >= 2*BLOQUEO_S; solapadas, algo más de BLOQUEO_S.
|
| 105 |
+
assert transcurrido < BLOQUEO_S * 1.7, (
|
| 106 |
+
f"{transcurrido:.2f}s para dos interpretaciones de {BLOQUEO_S}s: se serializaron"
|
| 107 |
+
)
|
| 108 |
+
|
| 109 |
+
|
| 110 |
+
async def test_el_alta_no_congela_el_bucle(alta_abierta):
|
| 111 |
+
"""scrypt (n=2**14) es caro A PROPÓSITO; en el bucle, cada alta congela el servicio."""
|
| 112 |
+
import httpx
|
| 113 |
+
|
| 114 |
+
from app import db
|
| 115 |
+
from app.main import app
|
| 116 |
+
|
| 117 |
+
# ASGITransport no dispara el lifespan (TestClient sí), así que la tabla no existiría.
|
| 118 |
+
db.inicializar_db()
|
| 119 |
+
|
| 120 |
+
transporte = httpx.ASGITransport(app=app)
|
| 121 |
+
async with httpx.AsyncClient(transport=transporte, base_url="http://test") as cliente:
|
| 122 |
+
tarea = asyncio.create_task(
|
| 123 |
+
cliente.post(
|
| 124 |
+
"/api/auth/registro",
|
| 125 |
+
json={
|
| 126 |
+
"nombre": "Hilo",
|
| 127 |
+
"apellido": "Vet",
|
| 128 |
+
"email": "hilo@example.com",
|
| 129 |
+
"password": "clave-segura-1",
|
| 130 |
+
},
|
| 131 |
+
)
|
| 132 |
+
)
|
| 133 |
+
latidos = await _contar_latidos(tarea, intervalo=0.002)
|
| 134 |
+
resp = await tarea
|
| 135 |
+
|
| 136 |
+
assert resp.status_code == 200, resp.text
|
| 137 |
+
assert latidos > 3, (
|
| 138 |
+
f"s��lo {latidos} latidos durante el alta: scrypt volvió a correr en el bucle de eventos"
|
| 139 |
+
)
|