#!/usr/bin/env python3 """ kiradesktop-gpu.py — Gestor de escritorios virtuales GPU para kirastream-gpu1 Gestiona hasta 6 escritorios X11 con Xorg+NVIDIA real, VNC via x11vnc y streaming MPEG-TS via h264_nvenc. Firefox usa GPU para render y NVDEC (VA-API). API en puerto 8900: GET /health → estado de todos los escritorios GET /desktop/{id}/stream → stream MPEG-TS (soporta múltiples clientes) POST /desktop/{id}/start → arrancar escritorio N POST /desktop/{id}/stop → parar escritorio N """ import asyncio import io import logging import os import shutil import tarfile import time from typing import Optional, Set import uvicorn from fastapi import FastAPI, HTTPException, Request from fastapi.responses import Response, StreamingResponse # ── Configuración ───────────────────────────────────────────────────────────── NUM_DESKTOPS = int(os.getenv("NUM_DESKTOPS", "6")) BASE_DISPLAY = int(os.getenv("BASE_DISPLAY", "1")) # :1 – :6 BASE_VNC_PORT = int(os.getenv("BASE_VNC_PORT", "5901")) # 5901 – 5906 RESOLUTION = os.getenv("RESOLUTION", "1920x1080") FRAMERATE = int(os.getenv("FRAMERATE", "30")) VIDEO_BITRATE = os.getenv("VIDEO_BITRATE", "8M") AUDIO_BITRATE = os.getenv("AUDIO_BITRATE", "128k") NVENC_PRESET = os.getenv("NVENC_PRESET", "p4") FFMPEG_BIN = os.getenv("FFMPEG_BIN", "ffmpeg") MY_IP = os.getenv("MY_IP", "10.10.10.239") WM_CMD = os.getenv("WM_CMD", "startxfce4") XORG_CONFIG = os.getenv("XORG_CONFIG", "/etc/X11/xorg-virtual.conf") logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s %(message)s") logger = logging.getLogger("kiradesktop-gpu") # ── Helpers ─────────────────────────────────────────────────────────────────── def _double(b: str) -> str: try: return f"{int(b[:-1]) * 2}{b[-1]}" except Exception: return "12M" async def _drain(proc: asyncio.subprocess.Process, label: str) -> None: try: async for line in proc.stderr: txt = line.decode(errors="replace").rstrip() if txt: logger.debug("[%s] %s", label, txt) except Exception: pass async def _wait_path(path: str, timeout: float = 10.0) -> bool: deadline = time.time() + timeout while time.time() < deadline: if os.path.exists(path): return True await asyncio.sleep(0.3) return False async def _wait_display(display: str, timeout: float = 12.0) -> bool: """Espera hasta que el display X11 acepte conexiones.""" deadline = time.time() + timeout while time.time() < deadline: try: proc = await asyncio.create_subprocess_exec( "xdpyinfo", "-display", display, stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, ) await asyncio.wait_for(proc.wait(), timeout=2.0) if proc.returncode == 0: return True except Exception: pass await asyncio.sleep(0.5) return False # ── DesktopSlot ─────────────────────────────────────────────────────────────── class DesktopSlot: def __init__(self, id: int): self.id = id self.disp_num = BASE_DISPLAY + id self.vnc_port = BASE_VNC_PORT + id self.pa_dir = f"/tmp/pa-{self.disp_num}" self.home_dir = f"/tmp/home-{self.disp_num}" self._xorg : Optional[asyncio.subprocess.Process] = None self._x11vnc: Optional[asyncio.subprocess.Process] = None self._wm : Optional[asyncio.subprocess.Process] = None self._pa : Optional[asyncio.subprocess.Process] = None self._ffmpeg: Optional[asyncio.subprocess.Process] = None self._tasks : list[asyncio.Task] = [] self._fan : Optional[asyncio.Task] = None self._clients: Set[asyncio.Queue] = set() self._lock = asyncio.Lock() # ── Properties ────────────────────────────────────────────────────────── @property def display(self) -> str: return f":{self.disp_num}" @property def pulse_socket(self) -> str: return f"unix:{self.pa_dir}/native" @property def running(self) -> bool: return self._xorg is not None and self._xorg.returncode is None @property def streaming(self) -> bool: return self._ffmpeg is not None and self._ffmpeg.returncode is None @property def stream_url(self) -> str: return f"http://{MY_IP}:8900/desktop/{self.id}/stream" def _setup_pulseaudio_config(self) -> None: """PulseAudio con null-sink virtual (sin dependencia de hardware ALSA).""" pa_config_dir = f"{self.home_dir}/.config/pulse" os.makedirs(pa_config_dir, mode=0o700, exist_ok=True) with open(f"{pa_config_dir}/default.pa", "w") as f: f.write( ".nofail\n" "load-module module-null-sink sink_name=virtual " "sink_properties=device.description=Virtual\n" "set-default-sink virtual\n" "set-default-source virtual.monitor\n" "load-module module-native-protocol-unix\n" ) def _env(self, extra: dict | None = None) -> dict: e = { **os.environ, "DISPLAY": self.display, "PULSE_SERVER": self.pulse_socket, "PULSE_RUNTIME_PATH": self.pa_dir, "HOME": self.home_dir, "DBUS_SESSION_BUS_ADDRESS": "", # Firefox: decode por hardware (NVDEC via VA-API) y render GPU "LIBVA_DRIVER_NAME": "nvidia", "MOZ_DISABLE_RDD_SANDBOX": "1", "MOZ_X11_EGL": "1", } if extra: e.update(extra) return e def _setup_firefox_profile(self) -> None: """Inyecta user.js con VA-API en todos los perfiles Firefox existentes.""" user_js = ( 'user_pref("media.ffmpeg.vaapi.enabled", true);\n' 'user_pref("media.hardware-video-decoding.force-enabled", true);\n' 'user_pref("gfx.webrender.all", true);\n' 'user_pref("media.av1.enabled", false);\n' ) ff_base = f"{self.home_dir}/.mozilla/firefox" os.makedirs(ff_base, mode=0o700, exist_ok=True) # Inyectar en cualquier perfil que ya exista try: for entry in os.listdir(ff_base): pdir = os.path.join(ff_base, entry) if os.path.isdir(pdir) and entry not in ("Crash Reports", "Pending Pings"): with open(os.path.join(pdir, "user.js"), "w") as f: f.write(user_js) except Exception: pass # Crear perfil propio como fallback profile_dir = f"{ff_base}/kiradesktop" os.makedirs(profile_dir, mode=0o700, exist_ok=True) with open(f"{profile_dir}/user.js", "w") as f: f.write(user_js) # ── Start ──────────────────────────────────────────────────────────────── async def start(self) -> dict: async with self._lock: if self.running: return {"ok": True, "already": True, "display": self.display} os.makedirs(self.pa_dir, mode=0o700, exist_ok=True) os.makedirs(self.home_dir, mode=0o700, exist_ok=True) # 1. PulseAudio con null-sink virtual logger.info("[desktop-%d] Starting PulseAudio", self.id) self._setup_pulseaudio_config() self._pa = await asyncio.create_subprocess_exec( "pulseaudio", "--start", "--exit-idle-time=-1", env={**os.environ, "PULSE_RUNTIME_PATH": self.pa_dir, "HOME": self.home_dir}, stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, ) await _wait_path(f"{self.pa_dir}/native", timeout=8.0) # 2. Xorg con driver NVIDIA real (GPU rendering, GLX, VDPAU) logger.info("[desktop-%d] Starting Xorg on %s (NVIDIA)", self.id, self.display) log_file = f"/tmp/xorg{self.disp_num}.log" self._xorg = await asyncio.create_subprocess_exec( "Xorg", self.display, "-config", XORG_CONFIG, "-noreset", "-logfile", log_file, "-logverbose", "1", stderr=asyncio.subprocess.DEVNULL, stdout=asyncio.subprocess.DEVNULL, ) # Esperar a que el display esté listo ready = await _wait_display(self.display, timeout=12.0) if not ready or self._xorg.returncode is not None: code = self._xorg.returncode # Mostrar últimas líneas del log para diagnóstico try: with open(log_file) as lf: tail = lf.read()[-2000:] except Exception: tail = "(sin log)" raise RuntimeError( f"Xorg no arrancó en {self.display} (exit={code})\n{tail}" ) # 3. x11vnc — acceso VNC al display Xorg logger.info("[desktop-%d] Starting x11vnc on port %d", self.id, self.vnc_port) self._x11vnc = await asyncio.create_subprocess_exec( "x11vnc", "-display", self.display, "-rfbport", str(self.vnc_port), "-nopw", "-forever", "-shared", "-quiet", "-noxdamage", stderr=asyncio.subprocess.DEVNULL, stdout=asyncio.subprocess.DEVNULL, ) await asyncio.sleep(0.5) # 4. Perfil Firefox con VA-API pre-configurado self._setup_firefox_profile() # 5. Window manager logger.info("[desktop-%d] Starting window manager: %s", self.id, WM_CMD) self._wm = await asyncio.create_subprocess_exec( "dbus-launch", "--exit-with-session", *WM_CMD.split(), env=self._env(), stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.PIPE, ) self._tasks.append(asyncio.create_task( _drain(self._wm, f"wm-{self.id}"))) await asyncio.sleep(3.0) # 6. ffmpeg x11grab → h264_nvenc await self._start_ffmpeg() logger.info("[desktop-%d] Ready — display=%s vnc=%d stream=%s", self.id, self.display, self.vnc_port, self.stream_url) return { "ok": True, "display": self.display, "vnc_port": self.vnc_port, "vnc_addr": f"{MY_IP}:{self.vnc_port}", "stream_url": self.stream_url, } async def _start_ffmpeg(self): cmd = [ FFMPEG_BIN, "-hide_banner", "-loglevel", "warning", # Audio PRIMERO: PulseAudio debe estar listo antes de abrir x11grab. # Si x11grab arranca primero, bufferiza ~30 frames mientras espera # audio, causando que el video salga 1s atrasado respecto al audio. "-f", "pulse", "-ac", "2", "-fragment_size", "7680", # chunks de 20ms → listo rápido "-thread_queue_size", "512", "-use_wallclock_as_timestamps", "1", "-i", "default", # Video SEGUNDO — x11grab arranca cuando audio ya está activo "-f", "x11grab", "-draw_mouse", "1", "-framerate", str(FRAMERATE), "-video_size", RESOLUTION, "-thread_queue_size", "512", "-use_wallclock_as_timestamps", "1", "-i", f"{self.display}+0,0", "-map", "0:a", "-map", "1:v", "-vf", "format=nv12,hwupload_cuda", "-c:v", "h264_nvenc", "-preset", NVENC_PRESET, "-rc", "cbr", "-profile:v", "high", "-level:v", "4.2", "-b:v", VIDEO_BITRATE, "-maxrate", VIDEO_BITRATE, "-bufsize", VIDEO_BITRATE, "-g", "60", "-keyint_min", "30", # Audio: AAC "-c:a", "aac", "-b:a", AUDIO_BITRATE, "-ar", "48000", "-ac", "2", "-muxdelay", "0", "-muxpreload", "0", # Output MPEG-TS "-f", "mpegts", "pipe:1", ] self._ffmpeg = await asyncio.create_subprocess_exec( *cmd, env=self._env(), stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, ) self._tasks.append(asyncio.create_task( _drain(self._ffmpeg, f"ffmpeg-{self.id}"))) self._fan = asyncio.create_task(self._fan_out()) logger.info("[desktop-%d] ffmpeg capture started (pid=%d)", self.id, self._ffmpeg.pid) async def _fan_out(self): try: while True: chunk = await self._ffmpeg.stdout.read(65536) if not chunk: break dead: Set[asyncio.Queue] = set() for q in list(self._clients): try: q.put_nowait(chunk) except asyncio.QueueFull: dead.add(q) for q in dead: self._clients.discard(q) try: q.put_nowait(None) except Exception: pass except Exception as e: logger.warning("[desktop-%d] fan_out: %s", self.id, e) finally: for q in list(self._clients): try: q.put_nowait(None) except Exception: pass self._clients.clear() # ── Stop ───────────────────────────────────────────────────────────────── async def stop(self) -> dict: async with self._lock: if not self.running: return {"ok": True, "already_stopped": True} logger.info("[desktop-%d] Stopping...", self.id) for proc in (self._ffmpeg, self._wm, self._x11vnc, self._xorg): if proc and proc.returncode is None: try: proc.kill() await asyncio.wait_for(proc.wait(), timeout=3.0) except Exception: pass try: import subprocess subprocess.run( ["pulseaudio", "--kill"], env={**os.environ, "PULSE_RUNTIME_PATH": self.pa_dir, "HOME": self.home_dir}, timeout=3, capture_output=True, ) except Exception: pass for t in self._tasks: t.cancel() if self._fan: self._fan.cancel() self._tasks.clear() self._fan = None self._xorg = self._x11vnc = self._wm = self._pa = self._ffmpeg = None logger.info("[desktop-%d] Stopped", self.id) return {"ok": True} # ── Status ─────────────────────────────────────────────────────────────── def status(self) -> dict: return { "id": self.id, "display": self.display, "vnc_port": self.vnc_port, "vnc_addr": f"{MY_IP}:{self.vnc_port}", "running": self.running, "streaming": self.streaming, "clients": len(self._clients), "stream_url": self.stream_url, } # ── Estado global ───────────────────────────────────────────────────────────── slots = [DesktopSlot(i) for i in range(NUM_DESKTOPS)] # ── FastAPI ─────────────────────────────────────────────────────────────────── app = FastAPI(title="KiraDesktop GPU Daemon", version="2.0") @app.get("/health") async def health(): return { "node": "kiradesktop-gpu", "type": "gpu", "desktops": NUM_DESKTOPS, "resolution": RESOLUTION, "desktop_status": [s.status() for s in slots], } @app.post("/desktop/{id}/start") async def start_desktop(id: int): if id < 0 or id >= NUM_DESKTOPS: raise HTTPException(404, f"Desktop {id} no existe") try: return await slots[id].start() except Exception as e: logger.error("start desktop %d: %s", id, e) raise HTTPException(500, str(e)) @app.post("/desktop/{id}/stop") async def stop_desktop(id: int): if id < 0 or id >= NUM_DESKTOPS: raise HTTPException(404, f"Desktop {id} no existe") return await slots[id].stop() @app.get("/desktop/{id}/stream") async def stream_desktop(id: int): if id < 0 or id >= NUM_DESKTOPS: raise HTTPException(404, f"Desktop {id} no existe") slot = slots[id] if not slot.streaming: raise HTTPException(503, "Desktop no está en streaming") q: asyncio.Queue = asyncio.Queue(maxsize=128) slot._clients.add(q) async def generator(): try: while True: chunk = await q.get() if chunk is None: break yield chunk except asyncio.CancelledError: pass finally: slot._clients.discard(q) return StreamingResponse(generator(), media_type="video/mp2t") def _export_slot(slot_idx: int | None) -> Response: if slot_idx is not None: if slot_idx < 0 or slot_idx >= len(slots): raise HTTPException(404, f"Slot {slot_idx} no existe") src = slots[slot_idx] if not os.path.exists(f"{src.home_dir}/.mozilla"): raise HTTPException(404, f"Slot {slot_idx} no tiene perfil de Firefox") else: src = next((s for s in slots if os.path.exists(f"{s.home_dir}/.mozilla")), None) if src is None: raise HTTPException(404, "Ningún slot tiene perfil de Firefox") buf = io.BytesIO() with tarfile.open(fileobj=buf, mode="w:gz") as tar: tar.add(f"{src.home_dir}/.mozilla", arcname=".mozilla") return Response( content=buf.getvalue(), media_type="application/octet-stream", headers={"Content-Disposition": "attachment; filename=firefox-profile.tar.gz"}, ) @app.get("/sync/firefox/export") async def export_firefox_gpu(): return _export_slot(None) @app.get("/sync/firefox/export/{slot_id}") async def export_firefox_gpu_slot(slot_id: int): return _export_slot(slot_id) def _fix_firefox_profile(mozilla_dir: str, home_dir: str) -> None: ff_dir = f"{mozilla_dir}/firefox" if not os.path.isdir(ff_dir): return # Fix profiles.ini: convert absolute paths → relative, and patch wrong home paths profiles_ini = f"{ff_dir}/profiles.ini" if os.path.exists(profiles_ini): lines = [] in_profile = False with open(profiles_ini) as f: for line in f: s = line.strip() if s.startswith("[") and s.endswith("]"): in_profile = s.startswith("[Profile") if in_profile and s.startswith("IsRelative=0"): lines.append("IsRelative=1\n") elif in_profile and s.startswith("Path=") and len(s) > 5 and s[5] == "/": lines.append(f"Path={os.path.basename(s[5:].rstrip('/'))}\n") else: lines.append(line) with open(profiles_ini, "w") as f: f.writelines(lines) installs_ini = f"{ff_dir}/installs.ini" if os.path.exists(installs_ini): lines = [] with open(installs_ini) as f: for line in f: s = line.strip() if s.startswith("Default=") and len(s) > 8 and s[8] == "/": lines.append(f"Default={os.path.basename(s[8:].rstrip('/'))}\n") else: lines.append(line) with open(installs_ini, "w") as f: f.writelines(lines) import glob for profile_dir in glob.glob(f"{ff_dir}/*/"): for name in ("lock", ".parentlock", ".parentlock.lck"): try: os.remove(os.path.join(profile_dir, name)) except OSError: pass @app.post("/sync/firefox/import") async def import_firefox_gpu(request: Request): data = await request.body() if not data: raise HTTPException(400, "No data received") for s in slots: mozilla_dir = f"{s.home_dir}/.mozilla" os.makedirs(s.home_dir, exist_ok=True) if os.path.exists(mozilla_dir): shutil.rmtree(mozilla_dir) with tarfile.open(fileobj=io.BytesIO(data), mode="r:gz") as tar: tar.extractall(s.home_dir) _fix_firefox_profile(mozilla_dir, s.home_dir) return {"ok": True, "slots": len(slots)} if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8900, workers=1)