#!/usr/bin/env python3 """ GStreamer WebRTC Agent for CSS Streaming Two tracks published to Cloudflare Calls: - Video: ximagesrc (Xwayland :99) → VP8 → webrtcbin - Game audio: pipewiresrc → Opus → webrtcbin CSS renders via Wine Wayland driver → Weston (Vulkan/DXVK) → Xwayland → X11 ximagesrc captures from Xwayland X11 display :99 Dependencies: apt: python3-gst-1.0 gstreamer1.0-plugins-bad gstreamer1.0-plugins-good gstreamer1.0-plugins-ugly gstreamer1.0-libav gstreamer1.0-pipewire xdotool weston xwayland wireplumber pipewire pipewire-pulse pip: requests websockets """ import asyncio import json import logging import os import queue as _queue import subprocess import sys import threading import time import gi gi.require_version('Gst', '1.0') gi.require_version('GstWebRTC', '1.0') gi.require_version('GstSdp', '1.0') from gi.repository import Gst, GstWebRTC, GstSdp, GLib import requests import websockets logging.basicConfig(level=logging.INFO, format='%(asctime)s [gst_agent] %(message)s', stream=sys.stderr) log = logging.getLogger(__name__) # ── Configuration ───────────────────────────────────────────────────────────── CF_APP_ID = os.environ.get('CF_APP_ID', '') CF_APP_TOKEN = os.environ.get('CF_APP_TOKEN', '') CF_BASE = f'https://rtc.live.cloudflare.com/v1/apps/{CF_APP_ID}' DISPLAY = os.environ.get('DISPLAY', ':99') CAP_W = int(os.environ.get('CAP_W', '800')) CAP_H = int(os.environ.get('CAP_H', '600')) CAP_FPS = int(os.environ.get('CAP_FPS', '60')) SIG_URL = os.environ.get('SIG_URL', 'ws://localhost:3000') # ── Cloudflare Calls REST helpers ───────────────────────────────────────────── def cf_post(path, body): r = requests.post(f'{CF_BASE}{path}', json=body, headers={ 'Authorization': f'Bearer {CF_APP_TOKEN}', 'Content-Type': 'application/json', }) data = r.json() if 'errorCode' in data: raise RuntimeError(f"CF API error: {data['errorDescription']}") return data def cf_put(path, body): r = requests.put(f'{CF_BASE}{path}', json=body, headers={ 'Authorization': f'Bearer {CF_APP_TOKEN}', 'Content-Type': 'application/json', }) data = r.json() if 'errorCode' in data: raise RuntimeError(f"CF API error: {data['errorDescription']}") return data # ── Input forwarding via a serialized xdotool worker ────────────────────────── # Previously every input event spawned a fresh `xdotool` process via Popen with # NO ordering guarantee. Under load (rapid mouse-move + key events) a keyup # could execute before its keydown, leaving a key logically stuck down in X11 — # that's the "keyboard stops responding until handoff/restart" bug. # # Fix: push all input commands onto a single queue processed by ONE worker # thread. Commands execute strictly in the order they were received, so a # keydown always runs before its matching keyup. Each command still runs as a # short-lived xdotool invocation, but serialized, so ordering is guaranteed. _xdotool_queue = _queue.Queue() def _xdotool_worker(): env = {**os.environ, 'DISPLAY': DISPLAY} while True: args = _xdotool_queue.get() # Coalesce a run of consecutive mouse moves into one, but only by # peeking non-destructively: we drain moves, and if we hit a non-move # command we execute the accumulated move, then execute that command. if args and args[0] == 'mousemove_relative': sum_dx = int(args[-2]) sum_dy = int(args[-1]) pending_cmd = None while True: try: nxt = _xdotool_queue.get_nowait() except _queue.Empty: break if nxt and nxt[0] == 'mousemove_relative': sum_dx += int(nxt[-2]) sum_dy += int(nxt[-1]) else: pending_cmd = nxt break # Execute the accumulated movement as a single command _run_xdotool(env, ['mousemove_relative', '--', str(sum_dx), str(sum_dy)]) # If we pulled a non-move command out of the queue, execute it too if pending_cmd is not None: _run_xdotool(env, [str(a) for a in pending_cmd]) else: _run_xdotool(env, [str(a) for a in args]) def _run_xdotool(env, argv): try: subprocess.run( ['xdotool'] + argv, env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=2, ) except Exception as e: log.error(f'xdotool error: {e}') def xdotool(*args): _xdotool_queue.put(args) # Track held keys so we can release them all on controller change _held_keys = set() def release_all_keys(): # Release all currently held keys and reset X11 modifier state for key in list(_held_keys): try: xdotool('keyup', key) except Exception: pass _held_keys.clear() # Reset X11 modifier state using a neutral key that has no in-game effect # Escape was previously used here but it toggles the in-game menu xdotool('key', '--clearmodifiers', 'ctrl') def handle_input(event): t = event.get('type') if t == 'mouse_move': dx = round(event.get('dx', 0)) dy = round(event.get('dy', 0)) if dx != 0 or dy != 0: xdotool('mousemove_relative', '--', str(dx), str(dy)) elif t == 'mouse_down': xdotool('mousedown', str(event.get('button', 1))) elif t == 'mouse_up': xdotool('mouseup', str(event.get('button', 1))) elif t == 'key_down': key = event.get('key', '') # Block characters that can be used for console command injection # via chat. ";" splits commands and '"' opens/closes command strings. if key in (';', 'semicolon', '"', 'quotedbl', 'quote'): return _held_keys.add(key) xdotool('keydown', key) elif t == 'key_up': key = event.get('key', '') if key in (';', 'semicolon', '"', 'quotedbl', 'quote'): return _held_keys.discard(key) xdotool('keyup', key) elif t == 'mouse_wheel': button = '4' if event.get('direction') == 'up' else '5' xdotool('click', button) elif t == 'release_keys': release_all_keys() elif t == 'connect_to_server': threading.Thread(target=_connect_to_server_sequence, daemon=True).start() # ── Automated menu-connect sequence ────────────────────────────────────────── # Triggered when the signaling server writes a new password (real visitor IP) # into GameMenu.res for a new controller. Opens the in-game escape menu, # moves the mouse to the "UNLOZE | ZombieEscape" button at its fixed absolute # position (measured empirically at 800x600: x=110, y=348), and clicks it. # This re-connects the client with the freshly-written password, which the # server-side connect extension reads to assign the correct name/SteamID — # all without restarting the CSS process itself. CONNECT_BUTTON_X = 110 CONNECT_BUTTON_Y = 348 def _connect_to_server_sequence(): log.info('connect_to_server: opening menu') xdotool('key', 'Escape') time.sleep(0.5) # let the menu render xdotool('mousemove', '--sync', str(CONNECT_BUTTON_X), str(CONNECT_BUTTON_Y)) time.sleep(0.2) # let the button highlight/register hover xdotool('click', '1') log.info('connect_to_server: clicked ZombieEscape button') # ── GStreamer Pipeline ──────────────────────────────────────────────────────── class CSSStreamAgent: def __init__(self): self.pipeline = None self.webrtcbin = None self.local_sdp = None self.loop = GLib.MainLoop() self._offer_ready = threading.Event() self._ice_ready = threading.Event() def build_pipeline(self): pipeline_str = ( # Video: ximagesrc → convert → scale → rate-enforce → VP8 # videoconvert must come before caps negotiation; # videoscale enforces width/height; videorate enforces framerate. f'ximagesrc display-name={DISPLAY} use-damage=0 ! ' f'videoconvert ! ' f'videoscale ! ' f'videorate ! ' f'video/x-raw,framerate={CAP_FPS}/1,width={CAP_W},height={CAP_H} ! ' f'vp8enc deadline=1 target-bitrate=1500000 keyframe-max-dist=60 ' f'error-resilient=partitions lag-in-frames=0 ! ' f'rtpvp8pay pt=96 ! ' f'application/x-rtp,media=video,encoding-name=VP8,payload=96 ! ' f'webrtcbin. ' # Game audio: pipewiresrc → Opus # target-object=game_output pins capture to the game audio sink. # Without this, pipewiresrc defaults to the system capture device # (virtual_mic / microphone) instead of the game audio monitor. # stream.capture.sink=true tells PipeWire to use game_output's # monitor ports rather than treating it as a regular source node. f'pipewiresrc target-object=game_output ' f'stream-properties="props,stream.capture.sink=(boolean)true" ! ' f'audio/x-raw,format=F32LE,rate=48000,channels=2 ! ' f'audioconvert ! ' f'opusenc ! ' f'rtpopuspay pt=97 ! ' f'application/x-rtp,media=audio,encoding-name=OPUS,payload=97 ! ' f'webrtcbin. ' # WebRTC f'webrtcbin name=webrtcbin bundle-policy=max-bundle ' f'stun-server=stun://stun.cloudflare.com:3478' ) log.info(f'Building pipeline {CAP_W}x{CAP_H}@{CAP_FPS}fps via Xwayland') self.pipeline = Gst.parse_launch(pipeline_str) self.webrtcbin = self.pipeline.get_by_name('webrtcbin') self.webrtcbin.connect('on-negotiation-needed', self._on_negotiation_needed) self.webrtcbin.connect('on-ice-candidate', self._on_ice_candidate) self.webrtcbin.connect('notify::ice-connection-state', self._on_ice_state) self.pipeline.set_state(Gst.State.READY) for i in range(2): try: t = self.webrtcbin.get_transceiver(i) if t: t.set_property('direction', GstWebRTC.WebRTCRTPTransceiverDirection.SENDONLY) except Exception: pass def _on_negotiation_needed(self, element): log.info('Creating offer...') promise = Gst.Promise.new_with_change_func(self._on_offer_created, element) element.emit('create-offer', None, promise) def _on_offer_created(self, promise, element): promise.wait() reply = promise.get_reply() offer = reply.get_value('offer') self.local_sdp = offer.sdp.as_text() set_promise = Gst.Promise.new() element.emit('set-local-description', offer, set_promise) set_promise.interrupt() log.info('Offer created and local description set') self._offer_ready.set() def _on_ice_candidate(self, element, mlineindex, candidate): log.debug(f'ICE candidate [{mlineindex}]: {candidate}') def _on_ice_state(self, element, prop): state = element.get_property('ice-connection-state') state_map = { GstWebRTC.WebRTCICEConnectionState.NEW: 'new', GstWebRTC.WebRTCICEConnectionState.CHECKING: 'checking', GstWebRTC.WebRTCICEConnectionState.CONNECTED: 'connected', GstWebRTC.WebRTCICEConnectionState.COMPLETED: 'completed', GstWebRTC.WebRTCICEConnectionState.DISCONNECTED: 'disconnected', GstWebRTC.WebRTCICEConnectionState.FAILED: 'failed', GstWebRTC.WebRTCICEConnectionState.CLOSED: 'closed', } log.info(f'ICE state: {state_map.get(state, state)}') if state in (GstWebRTC.WebRTCICEConnectionState.CONNECTED, GstWebRTC.WebRTCICEConnectionState.COMPLETED): self._ice_ready.set() elif state == GstWebRTC.WebRTCICEConnectionState.FAILED: log.error('ICE failed — exiting process so start.sh can restart it') # os._exit forces immediate termination of the entire process, # including the main thread blocked in asyncio. self.loop.quit() # alone only stops the GLib thread, leaving the asyncio signaling # loop running forever — which is why it got stuck instead of # restarting. os._exit(1) elif state == GstWebRTC.WebRTCICEConnectionState.DISCONNECTED: log.warning('ICE disconnected — monitoring for recovery') def set_remote_description(self, sdp_text, sdp_type='answer'): _, sdp = GstSdp.SDPMessage.new_from_text(sdp_text) wtype = GstWebRTC.WebRTCSDPType.ANSWER if sdp_type == 'answer' else GstWebRTC.WebRTCSDPType.OFFER desc = GstWebRTC.WebRTCSessionDescription.new(wtype, sdp) promise = Gst.Promise.new() self.webrtcbin.emit('set-remote-description', desc, promise) promise.interrupt() log.info('Remote description set') def get_mids(self): mids = [] for line in self.local_sdp.split('\n'): line = line.strip() if line.startswith('a=mid:'): mids.append(line[6:].strip()) return mids def start(self): self.pipeline.set_state(Gst.State.PLAYING) log.info('Pipeline playing') def stop(self): if self.pipeline: self.pipeline.set_state(Gst.State.NULL) def run_glib_loop(self): try: self.loop.run() except Exception as e: log.error(f'GLib loop error: {e}') finally: self.stop() # ── CSS process watchdog ────────────────────────────────────────────────────── def css_watchdog(): # If cstrike_win64.exe crashes (e.g. Wine c0000409 stack overrun), the # stream keeps running but shows a black screen forever with no recovery. # This watchdog detects the crash and triggers a restart. # # IMPORTANT: this must NOT relaunch CSS directly itself. Doing so # (a bare `wine revLoader.exe` Popen call) created a second, uncoordinated # restart authority alongside the signaling server's own kill+relaunch # logic (manual button / scheduled 5h restart). If both fired close # together, two separate Wine processes could end up racing to launch # CSS, producing duplicate/orphaned windows — xdotool's window search # would then target an arbitrary one, causing keyboard input to land on # a stale/backgrounded window while video (which captures the whole X # display, not a specific window) kept working fine. That mismatch is # exactly the "mouse/video fine, keyboard dead" symptom that was # reported after a watchdog relaunch overlapped with a manual restart. # # So instead we just notify the signaling server, which has the single # authoritative kill+relaunch implementation (performCssRestart) shared # by the manual restart button, the scheduled 5-hour restart, and this # watchdog — never more than one restart in flight at a time. # # We only start watching once CSS has been seen alive at least once, to # avoid racing its initial startup. Require several consecutive misses to # avoid false positives during a normal controller-triggered handoff. restart_url = SIG_URL.replace('ws://', 'http://').replace('wss://', 'https://') if not restart_url.endswith('/'): restart_url += '/' restart_url += 'webclient-force-restart' seen_alive = False missing_count = 0 while True: time.sleep(5) try: result = subprocess.run( ['pgrep', '-f', 'cstrike_win64.exe'], capture_output=True, text=True ) alive = result.returncode == 0 except Exception: alive = True # on error, assume alive (don't false-trigger) if alive: seen_alive = True missing_count = 0 continue if not seen_alive: continue missing_count += 1 log.warning(f'CSS process not found ({missing_count}/4)') # 4 consecutive misses (~20s) — long enough to not collide with the # signaling server's own kill+relaunch cycle on controller change. if missing_count >= 4: log.error('CSS appears crashed — requesting restart from signaling server') try: resp = requests.post(restart_url, timeout=5) log.info(f'restart request sent to {restart_url}, status={resp.status_code}') except Exception as e: log.error(f'watchdog restart request failed: {e}') # Reset so we give the restart time to complete before counting # misses again. seen_alive = False missing_count = 0 # ── Signaling ───────────────────────────────────────────────────────────────── async def run_signaling(session_id, video_track_name, audio_track_name): log.info(f'Connecting to signaling server {SIG_URL}') async with websockets.connect(SIG_URL) as ws: await ws.send(json.dumps({ 'type': 'register', 'sessionId': session_id, 'videoTrackName': video_track_name, 'audioTrackName': audio_track_name, 'containerName': os.environ.get('CONTAINER_NAME', 'css-test'), })) log.info('Registered with signaling server') async for raw in ws: try: msg = json.loads(raw) handle_input(msg) except Exception as e: log.error(f'Input error: {e}') # ── Main ────────────────────────────────────────────────────────────────────── def main(): if not CF_APP_ID or not CF_APP_TOKEN: log.error('CF_APP_ID and CF_APP_TOKEN must be set') sys.exit(1) Gst.init(None) log.info(f'Starting — {CAP_W}x{CAP_H}@{CAP_FPS}fps display={DISPLAY}') # Start the serialized xdotool input worker (guarantees keydown/keyup order) xw = threading.Thread(target=_xdotool_worker, daemon=True) xw.start() agent = CSSStreamAgent() agent.build_pipeline() agent.start() t = threading.Thread(target=agent.run_glib_loop, daemon=True) t.start() log.info('Waiting for SDP offer...') if not agent._offer_ready.wait(timeout=15): log.error('Timeout waiting for offer') sys.exit(1) log.info('Creating Cloudflare Calls session...') session_result = cf_post('/sessions/new', { 'sessionDescription': {'type': 'offer', 'sdp': agent.local_sdp} }) session_id = session_result['sessionId'] log.info(f'Session created: {session_id}') agent.set_remote_description(session_result['sessionDescription']['sdp']) log.info('Waiting for ICE...') if not agent._ice_ready.wait(timeout=15): log.error('Timeout waiting for ICE') sys.exit(1) log.info('ICE connected') mids = agent.get_mids() log.info(f'MIDs: {mids}') if len(mids) < 2: log.error(f'Expected 2 MIDs, got {len(mids)}: {mids}') sys.exit(1) ts = int(time.time() * 1000) video_track_name = f'css-video-{ts}' audio_track_name = f'css-audio-{ts}' track_result = cf_post(f'/sessions/{session_id}/tracks/new', { 'sessionDescription': {'type': 'offer', 'sdp': agent.local_sdp}, 'tracks': [ {'location': 'local', 'mid': mids[0], 'trackName': video_track_name}, {'location': 'local', 'mid': mids[1], 'trackName': audio_track_name}, ], }) if track_result.get('sessionDescription'): agent.set_remote_description(track_result['sessionDescription']['sdp']) log.info(f'Video track: {video_track_name}') log.info(f'Audio track: {audio_track_name}') log.info(f'Streaming — sessionId={session_id}') # Start CSS crash watchdog (relaunches the game if it dies unexpectedly) wd = threading.Thread(target=css_watchdog, daemon=True) wd.start() asyncio.run(run_signaling(session_id, video_track_name, audio_track_name)) if __name__ == '__main__': main()