crowdata / app /utils /telemetry.py
YOSOYYONOSOYOTRO's picture
Upload folder using huggingface_hub (part 2)
4223796 verified
Raw
History Blame Contribute Delete
13.5 kB
"""
TelemetrΓ­a de Scrapers β€” CrowData.
Registra el estado de salud de cada scraper tras su ejecuciΓ³n:
- Último estado (ok / error / bloqueado)
- Timestamp del ΓΊltimo run
- Tasa de Γ©xito acumulada (rolling window de 24h)
- Latencia promedio
Los resultados se almacenan en memoria (TTL de 48h) y se exponen
via el endpoint /reports/scrapers/health.
"""
import time
import logging
from collections import deque, defaultdict
from datetime import datetime, timezone
from threading import Lock
from typing import Optional
logger = logging.getLogger(__name__)
# ──────────────────────────────────────────────────────────────────────────────
# Constantes
# ──────────────────────────────────────────────────────────────────────────────
MAX_HISTORY = 50 # MΓ‘ximo de registros por scraper
WINDOW_SECS = 86400 # Ventana de cΓ‘lculo de tasa de Γ©xito: 24h
# ──────────────────────────────────────────────────────────────────────────────
# Estado global en memoria (thread-safe vΓ­a Lock)
# ──────────────────────────────────────────────────────────────────────────────
_lock = Lock()
# scraper_name -> deque of (timestamp_float, status: "ok"|"error"|"blocked", latency_ms: float)
_history: dict[str, deque] = defaultdict(lambda: deque(maxlen=MAX_HISTORY))
# scraper_name -> dict con datos del ΓΊltimo run
_last_run: dict[str, dict] = {}
# Metadatos estΓ‘ticos de cada scraper (descripciΓ³n y fuente)
SCRAPER_META = {
"ARCA/AFIP": {"fuente": "api.afip.gov.ar", "descripcion": "SituaciΓ³n fiscal, IVA, Monotributo"},
"BCRA": {"fuente": "api.bcra.gob.ar", "descripcion": "SituaciΓ³n crediticia y cheques rechazados"},
"BoletΓ­n Oficial": {"fuente": "www.boletinoficial.gob.ar", "descripcion": "Publicaciones oficiales"},
"Boletines Provinciales": {"fuente": "varios", "descripcion": "Boletines oficiales de 24 provincias"},
"IGJ": {"fuente": "sistemas.jus.gob.ar", "descripcion": "Sociedades y directivos registrados"},
"Poder Judicial": {"fuente": "scw.pjn.gov.ar", "descripcion": "Causas judiciales federales"},
"JUBA (Poder Judicial Bs As)":{"fuente": "juba.scba.gov.ar", "descripcion": "Causas en Justicia Bonaerense"},
"DNRPA / Automotores": {"fuente": "dnrpa.gov.ar", "descripcion": "VehΓ­culos registrados"},
"SINAI / Infracciones": {"fuente": "infraccionesba.gba.gob.ar", "descripcion": "Infracciones de trΓ‘nsito ANSV"},
"INPI (Marcas y Patentes)": {"fuente": "markaronline.inpi.gob.ar", "descripcion": "Marcas y patentes comerciales"},
"ANSES": {"fuente": "api.anses.gob.ar", "descripcion": "Aportes, jubilaciones, obra social"},
"Redes Sociales OSINT": {"fuente": "osint", "descripcion": "Perfiles pΓΊblicos en redes sociales"},
"TelefonΓ­a": {"fuente": "osint", "descripcion": "NΓΊmeros telefΓ³nicos vinculados"},
"RENAPER": {"fuente": "argentina.gob.ar", "descripcion": "ValidaciΓ³n de identidad RENAPER"},
"RENAPER Facial": {"fuente": "argentina.gob.ar", "descripcion": "BiometrΓ­a y estado de DNI"},
"Colegios Profesionales": {"fuente": "datos.gob.ar", "descripcion": "MatrΓ­culas profesionales activas"},
"RUIDO / SSSalud": {"fuente": "sssalud.gob.ar", "descripcion": "Cobertura de salud y obra social"},
"COMPR.AR / Contrataciones": {"fuente": "comprear.gob.ar", "descripcion": "Contratos con el Estado"},
"PadrΓ³n Electoral": {"fuente": "padron.gob.ar", "descripcion": "Datos del padrΓ³n electoral"},
"SGARHU / AcadΓ©mico": {"fuente": "siu.edu.ar", "descripcion": "TΓ­tulos universitarios oficiales"},
"SISCOP / Registro Civil": {"fuente": "siscop.gov.ar", "descripcion": "DetecciΓ³n de defunciΓ³n"},
"ARBA Automotores": {"fuente": "arba.gob.ar", "descripcion": "Deuda de patentes vehiculares ARBA"},
"ARBA Catastro": {"fuente": "arba.gob.ar", "descripcion": "Inmuebles y deuda inmobiliaria"},
"CNV (ComisiΓ³n Nacional de Valores)": {"fuente": "cnv.gob.ar", "descripcion": "Registro de agentes financieros"},
"UIF / PEPs": {"fuente": "uif.gob.ar", "descripcion": "Personas Expuestas PolΓ­ticamente"},
"ARBA / AGIP": {"fuente": "arba.gob.ar", "descripcion": "Deudas impositivas provinciales"},
"Name Search OSINT": {"fuente": "osint", "descripcion": "BΓΊsqueda de nombre en fuentes abiertas"},
"PadrΓ³n RUIDO": {"fuente": "sssalud.gob.ar", "descripcion": "Obra social por CUIL"},
"CONTRATAR": {"fuente": "datos.gob.ar", "descripcion": "Contrataciones pΓΊblicas de obra pΓΊblica"},
"AFIP_CONSTANCIA": {"fuente": "soa.afip.gob.ar", "descripcion": "Constancia de inscripciΓ³n fiscal AFIP"},
"DEUDORES_ALIMENTARIOS": {"fuente": "rdam.mjus.gba.gob.ar", "descripcion": "Registro de deudores alimentarios morosos PBA"},
}
# ──────────────────────────────────────────────────────────────────────────────
# API de registro de resultados
# ──────────────────────────────────────────────────────────────────────────────
def record_scraper_result(
scraper_name: str,
status: str, # "ok" | "error" | "blocked" | "empty"
latency_ms: float,
detail: Optional[str] = None,
records_found: int = 0
):
"""
Registra el resultado de una ejecuciΓ³n de un scraper.
Llamar desde el orquestador tras cada scraper en service.py.
Args:
scraper_name: Nombre canΓ³nico del scraper (source_name)
status: "ok" si obtuvo datos, "empty" si ejecutΓ³ pero sin hallazgos,
"error" si lanzΓ³ excepciΓ³n, "blocked" si fue bloqueado antibot
latency_ms: Tiempo de ejecuciΓ³n en milisegundos
detail: Mensaje de error o detalle adicional opcional
records_found: Cantidad de registros obtenidos
"""
ts = time.time()
entry = {
"ts": ts,
"status": status,
"latency_ms": round(latency_ms, 1),
"detail": detail,
"records_found": records_found,
}
with _lock:
_history[scraper_name].append(entry)
_last_run[scraper_name] = {
**entry,
"last_seen": datetime.fromtimestamp(ts, tz=timezone.utc).isoformat(),
}
# ──────────────────────────────────────────────────────────────────────────────
# API de consulta de estado
# ──────────────────────────────────────────────────────────────────────────────
def get_scraper_health() -> list[dict]:
"""
Retorna la lista de scrapers con su estado de salud actual.
Returns:
Lista de dicts con: name, status, success_rate_24h, avg_latency_ms,
last_seen, records_found, fuente, descripcion
"""
now = time.time()
cutoff = now - WINDOW_SECS
results = []
# Incluir todos los scrapers conocidos, aunque no hayan corrido aΓΊn
all_names = set(SCRAPER_META.keys()) | set(_last_run.keys())
with _lock:
for name in sorted(all_names):
meta = SCRAPER_META.get(name, {"fuente": "desconocida", "descripcion": ""})
last = _last_run.get(name)
history = list(_history.get(name, []))
# Calcular tasa de Γ©xito en las ΓΊltimas 24h
window_entries = [e for e in history if e["ts"] >= cutoff]
if window_entries:
ok_count = sum(1 for e in window_entries if e["status"] in ("ok", "empty"))
success_rate = round(ok_count / len(window_entries) * 100, 1)
avg_latency = round(
sum(e["latency_ms"] for e in window_entries) / len(window_entries), 1
)
else:
success_rate = None
avg_latency = None
results.append({
"name": name,
"fuente": meta["fuente"],
"descripcion": meta["descripcion"],
"status": last["status"] if last else "sin_datos",
"last_seen": last["last_seen"] if last else None,
"latency_ms": last["latency_ms"] if last else None,
"avg_latency_ms_24h": avg_latency,
"success_rate_24h": success_rate,
"records_found_last": last["records_found"] if last else 0,
"detail": last.get("detail") if last else None,
"runs_24h": len(window_entries),
})
return results
def get_system_health_summary() -> dict:
"""
Resumen ejecutivo del estado del sistema de scrapers.
"""
health = get_scraper_health()
total = len(health)
statuses = [s["status"] for s in health]
ok_count = statuses.count("ok") + statuses.count("empty")
error_count = statuses.count("error")
blocked_count = statuses.count("blocked")
sin_datos = statuses.count("sin_datos")
return {
"total_scrapers": total,
"operativos": ok_count,
"con_error": error_count,
"bloqueados": blocked_count,
"sin_ejecutar": sin_datos,
"health_pct": round(ok_count / max(total - sin_datos, 1) * 100, 1),
"scrapers": health,
}
# ──────────────────────────────────────────────────────────────────────────────
# Alertas de monitoreo
# ──────────────────────────────────────────────────────────────────────────────
ALERT_CONSECUTIVE_ERRORS = 3 # Errores consecutivos para activar alerta
ALERT_MIN_RUNS = 5 # MΓ­nimo de runs antes de evaluar
ALERT_LOW_SUCCESS_PCT = 30 # % mΓ­nimo de Γ©xito en 24h para no alertar
def get_scrapers_alerts() -> list[dict]:
"""
Detecta scrapers con problemas consistentes y genera alertas.
Returns:
Lista de alertas con: scraper, tipo, severidad, mensaje, detalle
"""
alerts = []
health = get_scraper_health()
for s in health:
name = s["name"]
status = s["status"]
runs = s.get("runs_24h", 0)
success_rate = s.get("success_rate_24h")
last_detail = s.get("detail")
# Alert 1: Scraper con error actual
if status == "error" and runs >= ALERT_CONSECUTIVE_ERRORS:
alerts.append({
"scraper": name,
"tipo": "error_actual",
"severidad": "alta",
"mensaje": f"{name} tiene error activo con {runs} ejecuciones recientes",
"detalle": last_detail,
})
# Alert 2: Scraper bloqueado por antibot
if status == "blocked":
alerts.append({
"scraper": name,
"tipo": "bloqueado_antibot",
"severidad": "alta",
"mensaje": f"{name} bloqueado por sistema antibot",
"detalle": last_detail,
})
# Alert 3: Tasa de Γ©xito baja en 24h
if success_rate is not None and success_rate < ALERT_LOW_SUCCESS_PCT and runs >= ALERT_MIN_RUNS:
alerts.append({
"scraper": name,
"tipo": "baja_tasa_exito",
"severidad": "media",
"mensaje": f"{name} solo tiene {success_rate}% de Γ©xito en las ΓΊltimas 24h ({runs} runs)",
"detalle": None,
})
# Alert 4: Scraper que nunca corriΓ³
if status == "sin_datos":
alerts.append({
"scraper": name,
"tipo": "nunca_ejecutado",
"severidad": "baja",
"mensaje": f"{name} nunca ha sido ejecutado",
"detalle": None,
})
# Ordenar por severidad (alta > media > baja)
severity_order = {"alta": 0, "media": 1, "baja": 2}
alerts.sort(key=lambda a: severity_order.get(a["severidad"], 3))
return alerts