Proyecto kiradesktop - escritorio virtual CPU/GPU para captura DRM
This commit is contained in:
@@ -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)
|
||||
Reference in New Issue
Block a user