Files
kiradesktop/gpu/kiradesktop-gpu.py

561 lines
22 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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)