413 lines
14 KiB
Python
413 lines
14 KiB
Python
#!/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)
|