Files
kiradesktop/cpu/kiradesktop-cpu.py

412 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",
"-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)