From ea5b82412f726b36abf98c12723ccd97d8c7c55c Mon Sep 17 00:00:00 2001 From: Joaquin Date: Mon, 13 Jul 2026 17:20:02 +0200 Subject: [PATCH] Proyecto kiradesktop - escritorio virtual CPU/GPU para captura DRM --- README.md | 125 ++++++++ cpu/kiradesktop-cpu.py | 411 ++++++++++++++++++++++++++ cpu/kiradesktop.py | 412 ++++++++++++++++++++++++++ cpu/kiradesktop.service | 23 ++ gpu/kiradesktop-gpu.py | 560 ++++++++++++++++++++++++++++++++++++ gpu/kiradesktop-gpu.service | 25 ++ 6 files changed, 1556 insertions(+) create mode 100644 README.md create mode 100644 cpu/kiradesktop-cpu.py create mode 100644 cpu/kiradesktop.py create mode 100644 cpu/kiradesktop.service create mode 100644 gpu/kiradesktop-gpu.py create mode 100644 gpu/kiradesktop-gpu.service diff --git a/README.md b/README.md new file mode 100644 index 0000000..137950b --- /dev/null +++ b/README.md @@ -0,0 +1,125 @@ +# KiraDesktop + +Sistema de escritorios virtuales remotos (X11 + XFCE) pensado para **captura de +streams con protección DRM (Widevine L3) reproducidos en Firefox**. Cada +instancia levanta un escritorio Linux headless, expone acceso VNC para +control/depuración y ofrece una API REST que produce un stream de vídeo+audio +en formato MPEG-TS (vía `ffmpeg`) listo para consumir con cualquier +reproductor/pipeline aguas abajo. + +La idea general: como el contenido DRM no se puede volcar "en crudo" (el +decodificador protege los frames), se renderiza en una sesión X11 real dentro +de Firefox y se captura la salida de pantalla (X11 grab) + audio (PulseAudio +null-sink), re-codificando con `ffmpeg` a H.264/AAC en MPEG-TS. El resultado es +funcionalmente equivalente a "grabar la pantalla" del escritorio remoto. + +## Arquitectura: variante CPU vs variante GPU + +### Variante CPU (`cpu/`) + +- Contenedores LXC (Proxmox `pct`), clonados: `kiradesktop-cpu-1` ... + `kiradesktop-cpu-12`. Son **12 instancias idénticas** en el cluster, cada + una con **1 escritorio** por contenedor. +- Servidor VNC: `Xtigervnc` (framebuffer virtual, no requiere GPU). +- Captura: `ffmpeg` con `-f x11grab` + encoder software `libx264` (preset + `superfast`, escalado a 1280x720 para aliviar la carga de CPU aunque el VNC + se sirve en 1080p). +- Cada contenedor corre un único proceso (`kiradesktop-cpu.py`) que gestiona + ese único escritorio (`/desktop/0/...`). +- Incluye `kiradesktop.py`, una versión anterior/backup del mismo script. Ver + sección "Diferencias" más abajo. + +### Variante GPU (`gpu/`) + +- Una VM KVM (`kirastream-gpu1`, Proxmox VM id 211) con GPU NVIDIA pasada por + passthrough/vGPU. +- Servidor VNC: `x11vnc` sobre sesiones `Xorg` reales con el driver NVIDIA + (permite aceleración GLX/VDPAU y decodificación por hardware NVDEC/VA-API + dentro de Firefox: `LIBVA_DRIVER_NAME=nvidia`, `MOZ_X11_EGL=1`). +- Captura: `ffmpeg` con `-f x11grab` + encoder por hardware `h264_nvenc`. +- Un único proceso (`kiradesktop-gpu.py`) gestiona **hasta 6 escritorios en + paralelo** (`/desktop/{0..5}/...`), cada uno con su propio display + (`:1`-`:6`), puerto VNC (`5901`-`5906`) y PulseAudio aislado. + +## API REST (puerto 8900 en ambas variantes) + +| Método | Ruta | Descripción | +|--------|-----------------------------------------|---------------------------------------| +| GET | `/health` | Estado de escritorio(s) | +| POST | `/desktop/{id}/start` | Arranca el escritorio `id` | +| POST | `/desktop/{id}/stop` | Detiene el escritorio `id` | +| GET | `/desktop/{id}/stream` | Stream MPEG-TS (multi-cliente) | +| GET | `/sync/firefox/export[/{slot}]` | Exporta el perfil de Firefox (.tar.gz)| +| POST | `/sync/firefox/import` | Importa un perfil de Firefox | + +En la variante CPU `{id}` es siempre `0` (un solo escritorio por contenedor). +En la variante GPU `{id}` va de `0` a `5` (hasta 6 escritorios por VM). + +## Puertos + +- **8900/tcp** — API REST (FastAPI + uvicorn, `0.0.0.0:8900`). +- **5901/tcp** (y consecutivos `5902`-`5906` en GPU) — VNC directo al + framebuffer X11 de cada escritorio. + +## Despliegue (systemd) + +Cada instancia corre como servicio systemd (`kiradesktop.service` en CPU, +`kiradesktop-gpu.service` en GPU) con `Restart=always` y toda la configuración +por variables de entorno — así el mismo `.py` sirve para las 12 instancias CPU +clonadas y para la VM GPU, cambiando solo el entorno: + +- `MY_IP` — IP pública/interna que se anuncia en las respuestas de la API + (`vnc_addr`, `stream_url`). **Debe fijarse por instancia** al clonar el + contenedor/VM (es lo único que realmente cambia entre clones CPU). +- `RESOLUTION`, `FRAMERATE`, `VIDEO_BITRATE`, `AUDIO_BITRATE` — parámetros de + codificación. +- `VNC_PORT` / `BASE_VNC_PORT` — puerto(s) VNC. +- `WM_CMD` — gestor de ventanas a lanzar (`startxfce4` por defecto). +- Variante GPU además: `NUM_DESKTOPS` (6), `BASE_DISPLAY`, `NVENC_PRESET`, + `XORG_CONFIG` (ruta al `xorg.conf` con el driver NVIDIA). + +Instalación típica en cada nodo: + +```bash +cp kiradesktop-cpu.py /opt/kiradesktop-cpu.py # o kiradesktop-gpu.py +cp kiradesktop.service /etc/systemd/system/ # o kiradesktop-gpu.service +systemctl daemon-reload +systemctl enable --now kiradesktop.service # o kiradesktop-gpu.service +``` + +## Diferencias entre `kiradesktop.py` y `kiradesktop-cpu.py` + +`cpu/kiradesktop.py` es una **versión previa/backup** del script CPU actual +(`cpu/kiradesktop-cpu.py`), tal como se encontró en `/opt` del contenedor. Se +incluye igual por trazabilidad histórica, pero **el que está desplegado y en +uso (referenciado por el `.service`) es `kiradesktop-cpu.py`**. + +La diferencia real entre ambos archivos es mínima: + +- Se eliminó un comentario explicativo sobre el escalado a 720p. +- Se cambió el preset del encoder `libx264` de `ultrafast` a `superfast` + (mejor relación calidad/CPU a cambio de un poco más de carga). + +El resto del código (gestión de PulseAudio, Xtigervnc, ffmpeg, API FastAPI, +export/import de perfil de Firefox) es idéntico. + +## Nota de seguridad (importante para un despliegue real) + +Los servidores VNC de ambas variantes se levantan **sin contraseña** +(`Xtigervnc -SecurityTypes None` en CPU, `x11vnc -nopw` en GPU) y la API REST +no implementa autenticación. Esto asume que el servicio corre en una red +interna/confiable (o detrás de un firewall/VPN/reverse-proxy que añada auth). +**No exponer el puerto 8900 ni los puertos VNC directamente a Internet** sin +añadir una capa de autenticación (por ejemplo, un proxy con Basic Auth o un +túnel VPN) delante de ambos servicios. + +No se encontraron contraseñas, tokens ni claves de API hardcodeadas en el +código (se auditó con `grep -iE "password|secret|token|key"` antes de subir +el repo). + +## Requisitos + +- Python 3 con `fastapi`, `uvicorn`. +- `ffmpeg` (con soporte `libx264` en CPU, o `h264_nvenc`/driver NVIDIA en GPU). +- `pulseaudio`, `dbus-launch`, XFCE (`startxfce4`) o el WM que se configure. +- CPU: `Xtigervnc`. GPU: `Xorg` + driver propietario NVIDIA + `x11vnc`. diff --git a/cpu/kiradesktop-cpu.py b/cpu/kiradesktop-cpu.py new file mode 100644 index 0000000..c5bff9c --- /dev/null +++ b/cpu/kiradesktop-cpu.py @@ -0,0 +1,411 @@ +#!/usr/bin/env python3 +""" +kiradesktop-cpu.py — Gestor de escritorio virtual CPU para KiraStream +Gestiona 1 escritorio X11 virtual con VNC y streaming MPEG-TS via libx264. +Diseñado para captura de contenido DRM (Widevine L3 en Firefox). + +API en puerto 8900: + GET /health → estado del escritorio + GET /desktop/0/stream → stream MPEG-TS + POST /desktop/0/start → arrancar escritorio + POST /desktop/0/stop → parar escritorio +""" + +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 + +RESOLUTION = os.getenv("RESOLUTION", "1920x1080") +FRAMERATE = int(os.getenv("FRAMERATE", "30")) +VIDEO_BITRATE = os.getenv("VIDEO_BITRATE", "4M") +AUDIO_BITRATE = os.getenv("AUDIO_BITRATE", "128k") +XVNC_BIN = os.getenv("XVNC_BIN", "Xtigervnc") +FFMPEG_BIN = os.getenv("FFMPEG_BIN", "ffmpeg") +MY_IP = os.getenv("MY_IP", "10.10.10.220") +VNC_PORT = int(os.getenv("VNC_PORT", "5901")) +X_DISPLAY = os.getenv("X_DISPLAY", ":1") +WM_CMD = os.getenv("WM_CMD", "startxfce4") + +logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(levelname)s %(name)s %(message)s") +logger = logging.getLogger("kiradesktop-cpu") + +PA_DIR = "/tmp/pa-desktop" +HOME_DIR = "/tmp/home-desktop" + + +def _double(b: str) -> str: + try: + return f"{int(b[:-1]) * 2}{b[-1]}" + except Exception: + return "8M" + +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 + + +class DesktopSlot: + def __init__(self): + self.id = 0 + self.pa_dir = PA_DIR + self.home_dir = HOME_DIR + + self._xvnc : 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() + + @property + def running(self) -> bool: + return self._xvnc is not None and self._xvnc.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/0/stream" + + def _env(self, extra: dict | None = None) -> dict: + e = {**os.environ, + "DISPLAY": X_DISPLAY, + "PULSE_SERVER": f"unix:{self.pa_dir}/native", + "PULSE_RUNTIME_PATH": self.pa_dir, + "HOME": self.home_dir, + "DBUS_SESSION_BUS_ADDRESS": ""} + if extra: + e.update(extra) + return e + + def _setup_pulseaudio_config(self) -> None: + """Create PulseAudio config with null-sink for LXC (no hardware audio).""" + pa_config_dir = f"{self.home_dir}/.config/pulse" + os.makedirs(pa_config_dir, mode=0o700, exist_ok=True) + default_pa = f"{pa_config_dir}/default.pa" + with open(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" + ) + + async def start(self) -> dict: + async with self._lock: + if self.running: + return {"ok": True, "already": True} + + os.makedirs(self.pa_dir, mode=0o700, exist_ok=True) + os.makedirs(self.home_dir, mode=0o700, exist_ok=True) + + # PulseAudio with null-sink (no ALSA hardware in LXC) + 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) + + # Xtigervnc + self._xvnc = await asyncio.create_subprocess_exec( + XVNC_BIN, X_DISPLAY, + "-rfbport", str(VNC_PORT), + "-SecurityTypes", "None", + "-localhost", "no", + "-geometry", RESOLUTION, + "-depth", "24", "-dpi", "96", + "-AlwaysShared", + stderr=asyncio.subprocess.PIPE, + stdout=asyncio.subprocess.DEVNULL, + ) + self._tasks.append(asyncio.create_task(_drain(self._xvnc, "xvnc"))) + await asyncio.sleep(1.5) + if self._xvnc.returncode is not None: + raise RuntimeError(f"Xtigervnc exited with code {self._xvnc.returncode}") + + # Window manager + 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, "wm"))) + await asyncio.sleep(3.0) + + # ffmpeg captura CPU + await self._start_ffmpeg() + + logger.info("Desktop started — display=%s vnc=%d stream=%s", + X_DISPLAY, VNC_PORT, self.stream_url) + return { + "ok": True, "display": X_DISPLAY, + "vnc_port": VNC_PORT, "vnc_addr": f"{MY_IP}:{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 + "-f", "pulse", "-ac", "2", + "-fragment_size", "7680", + "-thread_queue_size", "512", + "-use_wallclock_as_timestamps", "1", + "-i", "default", + # Video SEGUNDO + "-f", "x11grab", "-draw_mouse", "1", + "-framerate", str(FRAMERATE), + "-video_size", RESOLUTION, + "-thread_queue_size", "512", + "-use_wallclock_as_timestamps", "1", + "-i", f"{X_DISPLAY}+0,0", + "-map", "0:a", + "-map", "1:v", + "-vf", "scale=1280:720", + "-c:v", "libx264", "-preset", "superfast", "-tune", "zerolatency", + "-pix_fmt", "yuv420p", + "-b:v", VIDEO_BITRATE, "-maxrate", VIDEO_BITRATE, + "-bufsize", VIDEO_BITRATE, + "-g", "30", "-keyint_min", "15", + "-c:a", "aac", "-b:a", AUDIO_BITRATE, "-ar", "48000", "-ac", "2", + "-muxdelay", "0", "-muxpreload", "0", + "-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, "ffmpeg"))) + self._fan = asyncio.create_task(self._fan_out()) + logger.info("ffmpeg capture started (pid=%d)", 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("fan_out: %s", e) + finally: + for q in list(self._clients): + try: + q.put_nowait(None) + except Exception: + pass + self._clients.clear() + + async def stop(self) -> dict: + async with self._lock: + if not self.running: + return {"ok": True, "already_stopped": True} + logger.info("Stopping desktop...") + for proc in (self._ffmpeg, self._wm, self._xvnc): + 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._xvnc = self._wm = self._pa = self._ffmpeg = None + logger.info("Desktop stopped") + return {"ok": True} + + def status(self) -> dict: + return { + "id": 0, "display": X_DISPLAY, + "vnc_port": VNC_PORT, "vnc_addr": f"{MY_IP}:{VNC_PORT}", + "running": self.running, "streaming": self.streaming, + "clients": len(self._clients), "stream_url": self.stream_url, + } + + +slot = DesktopSlot() +app = FastAPI(title="KiraDesktop CPU Daemon", version="1.0") + + +@app.get("/health") +async def health(): + return { + "node": f"kiradesktop-cpu ({MY_IP})", + "type": "cpu", + "desktops": 1, + "resolution": RESOLUTION, + "desktop_status": [slot.status()], + } + + +@app.post("/desktop/0/start") +async def start_desktop(): + try: + return await slot.start() + except Exception as e: + logger.error("start: %s", e) + raise HTTPException(500, str(e)) + + +@app.post("/desktop/0/stop") +async def stop_desktop(): + return await slot.stop() + + +@app.get("/desktop/0/stream") +async def stream_desktop(): + 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") + + +@app.get("/sync/firefox/export") +async def export_firefox(): + mozilla_dir = f"{HOME_DIR}/.mozilla" + if not os.path.exists(mozilla_dir): + raise HTTPException(404, "No Firefox profile found") + buf = io.BytesIO() + with tarfile.open(fileobj=buf, mode="w:gz") as tar: + tar.add(mozilla_dir, arcname=".mozilla") + return Response( + content=buf.getvalue(), + media_type="application/octet-stream", + headers={"Content-Disposition": "attachment; filename=firefox-profile.tar.gz"}, + ) + + +@app.post("/sync/firefox/import") +async def import_firefox(request: Request): + data = await request.body() + if not data: + raise HTTPException(400, "No data received") + mozilla_dir = f"{HOME_DIR}/.mozilla" + os.makedirs(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(HOME_DIR) + _fix_firefox_profile(mozilla_dir) + return {"ok": True} + + +def _fix_firefox_profile(mozilla_dir: str) -> None: + ff_dir = f"{mozilla_dir}/firefox" + if not os.path.isdir(ff_dir): + return + + # Convert absolute paths → relative in profiles.ini + 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) + + # Fix installs.ini absolute paths + 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) + + # Remove lock files + 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 + + +if __name__ == "__main__": + uvicorn.run(app, host="0.0.0.0", port=8900, workers=1) diff --git a/cpu/kiradesktop.py b/cpu/kiradesktop.py new file mode 100644 index 0000000..aee48bc --- /dev/null +++ b/cpu/kiradesktop.py @@ -0,0 +1,412 @@ +#!/usr/bin/env python3 +""" +kiradesktop-cpu.py — Gestor de escritorio virtual CPU para KiraStream +Gestiona 1 escritorio X11 virtual con VNC y streaming MPEG-TS via libx264. +Diseñado para captura de contenido DRM (Widevine L3 en Firefox). + +API en puerto 8900: + GET /health → estado del escritorio + GET /desktop/0/stream → stream MPEG-TS + POST /desktop/0/start → arrancar escritorio + POST /desktop/0/stop → parar escritorio +""" + +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 + +RESOLUTION = os.getenv("RESOLUTION", "1920x1080") +FRAMERATE = int(os.getenv("FRAMERATE", "30")) +VIDEO_BITRATE = os.getenv("VIDEO_BITRATE", "4M") +AUDIO_BITRATE = os.getenv("AUDIO_BITRATE", "128k") +XVNC_BIN = os.getenv("XVNC_BIN", "Xtigervnc") +FFMPEG_BIN = os.getenv("FFMPEG_BIN", "ffmpeg") +MY_IP = os.getenv("MY_IP", "10.10.10.220") +VNC_PORT = int(os.getenv("VNC_PORT", "5901")) +X_DISPLAY = os.getenv("X_DISPLAY", ":1") +WM_CMD = os.getenv("WM_CMD", "startxfce4") + +logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(levelname)s %(name)s %(message)s") +logger = logging.getLogger("kiradesktop-cpu") + +PA_DIR = "/tmp/pa-desktop" +HOME_DIR = "/tmp/home-desktop" + + +def _double(b: str) -> str: + try: + return f"{int(b[:-1]) * 2}{b[-1]}" + except Exception: + return "8M" + +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 + + +class DesktopSlot: + def __init__(self): + self.id = 0 + self.pa_dir = PA_DIR + self.home_dir = HOME_DIR + + self._xvnc : 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() + + @property + def running(self) -> bool: + return self._xvnc is not None and self._xvnc.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/0/stream" + + def _env(self, extra: dict | None = None) -> dict: + e = {**os.environ, + "DISPLAY": X_DISPLAY, + "PULSE_SERVER": f"unix:{self.pa_dir}/native", + "PULSE_RUNTIME_PATH": self.pa_dir, + "HOME": self.home_dir, + "DBUS_SESSION_BUS_ADDRESS": ""} + if extra: + e.update(extra) + return e + + def _setup_pulseaudio_config(self) -> None: + """Create PulseAudio config with null-sink for LXC (no hardware audio).""" + pa_config_dir = f"{self.home_dir}/.config/pulse" + os.makedirs(pa_config_dir, mode=0o700, exist_ok=True) + default_pa = f"{pa_config_dir}/default.pa" + with open(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" + ) + + async def start(self) -> dict: + async with self._lock: + if self.running: + return {"ok": True, "already": True} + + os.makedirs(self.pa_dir, mode=0o700, exist_ok=True) + os.makedirs(self.home_dir, mode=0o700, exist_ok=True) + + # PulseAudio with null-sink (no ALSA hardware in LXC) + 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) + + # Xtigervnc + self._xvnc = await asyncio.create_subprocess_exec( + XVNC_BIN, X_DISPLAY, + "-rfbport", str(VNC_PORT), + "-SecurityTypes", "None", + "-localhost", "no", + "-geometry", RESOLUTION, + "-depth", "24", "-dpi", "96", + "-AlwaysShared", + stderr=asyncio.subprocess.PIPE, + stdout=asyncio.subprocess.DEVNULL, + ) + self._tasks.append(asyncio.create_task(_drain(self._xvnc, "xvnc"))) + await asyncio.sleep(1.5) + if self._xvnc.returncode is not None: + raise RuntimeError(f"Xtigervnc exited with code {self._xvnc.returncode}") + + # Window manager + 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, "wm"))) + await asyncio.sleep(3.0) + + # ffmpeg captura CPU + await self._start_ffmpeg() + + logger.info("Desktop started — display=%s vnc=%d stream=%s", + X_DISPLAY, VNC_PORT, self.stream_url) + return { + "ok": True, "display": X_DISPLAY, + "vnc_port": VNC_PORT, "vnc_addr": f"{MY_IP}:{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 + "-f", "pulse", "-ac", "2", + "-fragment_size", "7680", + "-thread_queue_size", "512", + "-use_wallclock_as_timestamps", "1", + "-i", "default", + # Video SEGUNDO + "-f", "x11grab", "-draw_mouse", "1", + "-framerate", str(FRAMERATE), + "-video_size", RESOLUTION, + "-thread_queue_size", "512", + "-use_wallclock_as_timestamps", "1", + "-i", f"{X_DISPLAY}+0,0", + "-map", "0:a", + "-map", "1:v", + # Escalar a 720p para reducir carga CPU del encoder (VNC sigue en 1080p) + "-vf", "scale=1280:720", + "-c:v", "libx264", "-preset", "ultrafast", "-tune", "zerolatency", + "-pix_fmt", "yuv420p", + "-b:v", VIDEO_BITRATE, "-maxrate", VIDEO_BITRATE, + "-bufsize", VIDEO_BITRATE, + "-g", "30", "-keyint_min", "15", + "-c:a", "aac", "-b:a", AUDIO_BITRATE, "-ar", "48000", "-ac", "2", + "-muxdelay", "0", "-muxpreload", "0", + "-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, "ffmpeg"))) + self._fan = asyncio.create_task(self._fan_out()) + logger.info("ffmpeg capture started (pid=%d)", 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("fan_out: %s", e) + finally: + for q in list(self._clients): + try: + q.put_nowait(None) + except Exception: + pass + self._clients.clear() + + async def stop(self) -> dict: + async with self._lock: + if not self.running: + return {"ok": True, "already_stopped": True} + logger.info("Stopping desktop...") + for proc in (self._ffmpeg, self._wm, self._xvnc): + 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._xvnc = self._wm = self._pa = self._ffmpeg = None + logger.info("Desktop stopped") + return {"ok": True} + + def status(self) -> dict: + return { + "id": 0, "display": X_DISPLAY, + "vnc_port": VNC_PORT, "vnc_addr": f"{MY_IP}:{VNC_PORT}", + "running": self.running, "streaming": self.streaming, + "clients": len(self._clients), "stream_url": self.stream_url, + } + + +slot = DesktopSlot() +app = FastAPI(title="KiraDesktop CPU Daemon", version="1.0") + + +@app.get("/health") +async def health(): + return { + "node": f"kiradesktop-cpu ({MY_IP})", + "type": "cpu", + "desktops": 1, + "resolution": RESOLUTION, + "desktop_status": [slot.status()], + } + + +@app.post("/desktop/0/start") +async def start_desktop(): + try: + return await slot.start() + except Exception as e: + logger.error("start: %s", e) + raise HTTPException(500, str(e)) + + +@app.post("/desktop/0/stop") +async def stop_desktop(): + return await slot.stop() + + +@app.get("/desktop/0/stream") +async def stream_desktop(): + 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") + + +@app.get("/sync/firefox/export") +async def export_firefox(): + mozilla_dir = f"{HOME_DIR}/.mozilla" + if not os.path.exists(mozilla_dir): + raise HTTPException(404, "No Firefox profile found") + buf = io.BytesIO() + with tarfile.open(fileobj=buf, mode="w:gz") as tar: + tar.add(mozilla_dir, arcname=".mozilla") + return Response( + content=buf.getvalue(), + media_type="application/octet-stream", + headers={"Content-Disposition": "attachment; filename=firefox-profile.tar.gz"}, + ) + + +@app.post("/sync/firefox/import") +async def import_firefox(request: Request): + data = await request.body() + if not data: + raise HTTPException(400, "No data received") + mozilla_dir = f"{HOME_DIR}/.mozilla" + os.makedirs(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(HOME_DIR) + _fix_firefox_profile(mozilla_dir) + return {"ok": True} + + +def _fix_firefox_profile(mozilla_dir: str) -> None: + ff_dir = f"{mozilla_dir}/firefox" + if not os.path.isdir(ff_dir): + return + + # Convert absolute paths → relative in profiles.ini + 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) + + # Fix installs.ini absolute paths + 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) + + # Remove lock files + 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 + + +if __name__ == "__main__": + uvicorn.run(app, host="0.0.0.0", port=8900, workers=1) diff --git a/cpu/kiradesktop.service b/cpu/kiradesktop.service new file mode 100644 index 0000000..6e607d0 --- /dev/null +++ b/cpu/kiradesktop.service @@ -0,0 +1,23 @@ +[Unit] +Description=KiraDesktop CPU - Virtual Desktop Streamer +After=network.target + +[Service] +Type=simple +WorkingDirectory=/opt +Environment=MY_IP=10.10.10.220 +Environment=VNC_PORT=5901 +Environment=X_DISPLAY=:1 +Environment=RESOLUTION=1920x1080 +Environment=FRAMERATE=30 +Environment=VIDEO_BITRATE=4M +Environment=AUDIO_BITRATE=128k +Environment=WM_CMD=startxfce4 +ExecStart=/usr/bin/python3 /opt/kiradesktop-cpu.py +Restart=always +RestartSec=5 +KillSignal=SIGTERM +TimeoutStopSec=15 + +[Install] +WantedBy=multi-user.target diff --git a/gpu/kiradesktop-gpu.py b/gpu/kiradesktop-gpu.py new file mode 100644 index 0000000..29fe785 --- /dev/null +++ b/gpu/kiradesktop-gpu.py @@ -0,0 +1,560 @@ +#!/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) diff --git a/gpu/kiradesktop-gpu.service b/gpu/kiradesktop-gpu.service new file mode 100644 index 0000000..c3a8206 --- /dev/null +++ b/gpu/kiradesktop-gpu.service @@ -0,0 +1,25 @@ +[Unit] +Description=KiraDesktop GPU - Virtual Desktop Manager (6 escritorios GPU) +After=network.target kiragpu.service + +[Service] +Type=simple +WorkingDirectory=/opt +Environment=NUM_DESKTOPS=6 +Environment=BASE_DISPLAY=1 +Environment=BASE_VNC_PORT=5901 +Environment=RESOLUTION=1920x1080 +Environment=FRAMERATE=60 +Environment=VIDEO_BITRATE=8M +Environment=AUDIO_BITRATE=128k +Environment=NVENC_PRESET=p4 +Environment=MY_IP=10.10.10.239 +Environment=WM_CMD=startxfce4 +ExecStart=/usr/bin/python3 /opt/kiradesktop-gpu.py +Restart=always +RestartSec=5 +KillSignal=SIGTERM +TimeoutStopSec=15 + +[Install] +WantedBy=multi-user.target