""" 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