File size: 6,389 Bytes
6374131
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
141
142
143
144
145
146
147
148
149
150
import os
import logging
import threading
import time
import requests
from datetime import datetime, timedelta

logger = logging.getLogger(__name__)

def get_hf_spaces_config():
    """
    Obtiene la lista de Hugging Face Spaces o URLs de workers configuradas desde las variables de entorno.
    Soporta repositorios (ej: 'Isaac105/melodix-worker') o URLs completas.
    """
    raw = os.getenv("HF_SPACES") or os.getenv("HF_WORKER_SPACES") or os.getenv("HF_WORKERS") or ""
    spaces = [s.strip() for s in raw.split(",") if s.strip()]
    return spaces

def wake_hf_workers():
    """
    Intenta despertar los Hugging Face Spaces o workers remotos configurados.
    Utiliza la API de Hugging Face (HfApi.restart_space / get_space_runtime)
    y pings HTTP directos a las URLs de los Spaces.
    """
    def _do_wake():
        hf_token = os.getenv("HF_TOKEN")
        spaces_config = get_hf_spaces_config()
        
        # Buscar también en la base de datos workers con URL configurada
        try:
            from app.database import SessionLocal
            from app.models.process import Worker
            db = SessionLocal()
            try:
                db_workers = db.query(Worker).all()
                for w in db_workers:
                    if w.url and w.url not in spaces_config:
                        spaces_config.append(w.url)
            finally:
                db.close()
        except Exception as e:
            logger.warning(f"[HF Wake] Error consultando URLs de workers en DB: {e}")

        if not spaces_config:
            logger.info("[HF Wake] No hay Hugging Face Spaces o URLs configuradas en HF_SPACES ni en DB.")
            return

        logger.info(f"[HF Wake] Iniciando intento de despertar para {len(spaces_config)} space(s) / worker(s)...")

        # Inicializar HfApi si existe token
        hf_api = None
        if hf_token:
            try:
                from huggingface_hub import HfApi
                hf_api = HfApi(token=hf_token)
            except Exception as e:
                logger.warning(f"[HF Wake] No se pudo inicializar HfApi: {e}")

        for target in spaces_config:
            try:
                # Si el target es formato repo_id ("Usuario/nombre-space")
                if "/" in target and not target.startswith("http"):
                    space_repo = target
                    if hf_api:
                        try:
                            runtime = hf_api.get_space_runtime(repo_id=space_repo)
                            stage = getattr(runtime, 'stage', 'UNKNOWN')
                            logger.info(f"[HF Wake] Space '{space_repo}' estado runtime: {stage}")
                            if stage in ["SLEEPING", "PAUSED", "STOPPED", "OFF", "BUILD_ERROR", "RUNTIME_ERROR"]:
                                logger.info(f"[HF Wake] Reiniciando Space dormido '{space_repo}' vía HF API...")
                                hf_api.restart_space(repo_id=space_repo)
                        except Exception as ex_api:
                            logger.warning(f"[HF Wake] Error en restart_space para '{space_repo}': {ex_api}")

                    parts = space_repo.split('/')
                    target_url = f"https://huggingface.co/spaces/{parts[0]}/{parts[1]}"
                else:
                    target_url = target

                # Hacer ping HTTP para activar el proxy de Hugging Face
                if target_url.startswith("http"):
                    logger.info(f"[HF Wake] Enviando ping HTTP para despertar a: {target_url}")
                    try:
                        resp = requests.get(target_url, timeout=12, headers={"User-Agent": "Melodix-Wake-Bot/1.0"})
                        logger.info(f"[HF Wake] Ping HTTP a {target_url} enviado. Status code: {resp.status_code}")
                    except Exception as http_err:
                        logger.warning(f"[HF Wake] Respuesta del ping HTTP a {target_url}: {http_err}")
            except Exception as err:
                logger.error(f"[HF Wake] Error despertando worker '{target}': {err}")

    # Ejecutar en hilo daemon independiente para no bloquear el hilo de FastAPI
    threading.Thread(target=_do_wake, daemon=True).start()

def check_and_wake_workers_if_needed():
    """
    Verifica si hay tareas pendientes en la DB y si NO hay workers activos.
    Si hay tareas sin atender, invoca wake_hf_workers().
    """
    try:
        from app.database import SessionLocal
        from app.models.process import Task, Worker
        db = SessionLocal()
        try:
            # 1. Tareas en espera
            pending_count = db.query(Task).filter(Task.status.in_(["PENDING", "PROCESSING"])).count()
            if pending_count == 0:
                return

            # 2. Workers reportando en los últimos 2 minutos
            two_mins_ago = datetime.utcnow() - timedelta(minutes=2)
            active_workers_count = db.query(Worker).filter(
                Worker.status.in_(["online", "busy"]),
                Worker.last_seen >= two_mins_ago
            ).count()

            logger.info(f"[HF Monitor] Tareas pendientes: {pending_count} | Workers activos: {active_workers_count}")

            if active_workers_count == 0:
                logger.warning(f"[HF Monitor] ⚠️ {pending_count} tarea(s) en espera y 0 workers activos. Despertando Hugging Face Workers...")
                wake_hf_workers()
        finally:
            db.close()
    except Exception as e:
        logger.error(f"[HF Monitor] Error verificando estado de tareas/workers: {e}")


_monitor_thread_started = False

def start_worker_monitor_loop(interval_seconds: int = 45):
    """
    Inicia un hilo en segundo plano que monitorea periódicamente la cola de tareas
    y despierta los workers de Hugging Face si están dormidos.
    """
    global _monitor_thread_started
    if _monitor_thread_started:
        return
    _monitor_thread_started = True

    def _loop():
        logger.info(f"[HF Monitor] Bucle de monitoreo de workers iniciado (intervalo: {interval_seconds}s).")
        while True:
            try:
                check_and_wake_workers_if_needed()
            except Exception as e:
                logger.error(f"[HF Monitor] Error en ciclo de monitoreo: {e}")
            time.sleep(interval_seconds)

    thread = threading.Thread(target=_loop, daemon=True)
    thread.start()