Verify shared adaptation and harden PC stream lifecycle

This commit is contained in:
saphid committed 2026-09-28 22:50:58 +10:00
1 parent c5fc552c60
commit cd202433bb
15 files changed
+641 -63

No files matched your search

+12 -10
View File
@@ -56,12 +56,12 @@ BROWSER_FLAGS = []
PANEL_BOX = (1920, 1080)
# Opens one viewer on gamescope's X display and gives its window its own panel
# id. Args: appid url width height tag [browser flags...]. The page puts "[tag]" in its title at
# id. Args: appid url width height tag profile [browser flags...]. The page puts "[tag]" in its title at
# once, which is how its X window is found (Chromium may hand the URL to an
# instance that's already running, so there's no process to follow).
LAUNCH = r"""set -u
appid=$1 url=$2 w=$3 h=$4 tag=$5
shift 5
appid=$1 url=$2 w=$3 h=$4 tag=$5 profile=$6
shift 6
export DISPLAY=:0 LC_ALL=C.UTF-8
unset WAYLAND_DISPLAY
if ! xprop -root GAMESCOPE_FOCUSABLE_WINDOWS >/dev/null 2>&1; then
@@ -73,15 +73,15 @@ common=(--ozone-platform=x11 --force-device-scale-factor=1 --no-first-run --no-d
--disable-features=Translate,MediaRouter --autoplay-policy=no-user-gesture-required
"--window-size=$w,$h" "$@" "--app=$url")
if [ -x "$HOME/chromium-xr/chrome" ]; then
cmd=("$HOME/chromium-xr/chrome" "--user-data-dir=$HOME/.local/share/frame-control/mac-view" "${common[@]}")
cmd=("$HOME/chromium-xr/chrome" "--user-data-dir=$HOME/.local/share/frame-control/$profile" "${common[@]}")
elif flatpak info org.chromium.Chromium >/dev/null 2>&1; then
cmd=(flatpak run org.chromium.Chromium
"--user-data-dir=$HOME/.var/app/org.chromium.Chromium/data/frame-mac-view" "${common[@]}")
"--user-data-dir=$HOME/.var/app/org.chromium.Chromium/data/frame-$profile" "${common[@]}")
else
echo "NO_BROWSER"
exit 3
fi
log=/tmp/frame-mac-view.log
log=/tmp/frame-$profile.log
setsid nohup "${cmd[@]}" >>"$log" 2>&1 </dev/null &
for _ in $(seq 1 60); do
sleep 0.5
@@ -117,6 +117,7 @@ def fit(w, h, box=PANEL_BOX):
class MacView:
host = "mac"
viewer_profile = "mac-view"
def __init__(self, tunnel_ssh, run, frame, track=None):
self.tunnel_ssh = list(tunnel_ssh)
@@ -200,7 +201,7 @@ class MacView:
break
self.agent.kill()
else:
raise MacViewError(f"The Mac streaming helper didn't start: {line.strip() or 'no output'}")
raise MacViewError(f"The streaming helper didn't start: {line.strip() or 'no output'}")
if self.port != int(m.group(1)):
self._drop_tunnel()
self.port = int(m.group(1))
@@ -216,7 +217,7 @@ class MacView:
with urllib.request.urlopen(req, timeout=10) as r:
return json.load(r)
except (urllib.error.URLError, OSError, ValueError) as e:
raise MacViewError(f"The Mac streaming helper didn't answer: {e}")
raise MacViewError(f"The streaming helper didn't answer: {e}")
# ---- the tunnel from the Frame ----
@@ -381,7 +382,7 @@ class MacView:
"bpp": q["bpp"]}
url = f"http://127.0.0.1:{self.remote_port}/view?{urlencode(params)}"
w, h = fit(width or 1280, height or 720)
args = " ".join(_quote(str(a)) for a in (appid, url, w, h, tag, *self.browser_flags))
args = " ".join(_quote(str(a)) for a in (appid, url, w, h, tag, self.viewer_profile, *self.browser_flags))
try:
out = self.run("bash -s -- " + args, stdin=LAUNCH, timeout=60)
except Exception as e: # noqa: BLE001 - the server's Failure carries the Frame's words
@@ -413,7 +414,8 @@ class MacView:
if self.shown or self.shows != shows or self.launching:
return
try:
self.run("pkill -f '[f]rame-control/mac-view|[d]ata/frame-mac-view' || true", timeout=10)
pattern = "[f]rame-control/" + self.viewer_profile + "|[d]ata/frame-" + self.viewer_profile
self.run("pkill -f " + _quote(pattern) + " || true", timeout=10)
except Exception:
pass
+95 -14
View File
@@ -13,6 +13,7 @@ from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import json
import math
import os
import queue
from pathlib import Path
import secrets
import socket
@@ -138,8 +139,12 @@ class Session:
self.stop_event, self.key_event = threading.Event(), threading.Event()
self.lock, self.input_lock = threading.RLock(), threading.RLock()
self.pending, self.last_submit = {}, 0
self.reconfiguring = False
self.test_inputs = queue.Queue(maxsize=128)
self.input = None if self.src == 'test' else WindowsInput(agent.windows, source) if agent.windows else PortalInput(self.native, source)
self.writer = None
self.released = False
self.input_enabled = self.src == 'test' or self.source.get('devices', 3) == 3
def gate(self, stage, pts, capture, arrived):
# Exceptions cannot cross a ctypes callback boundary.
@@ -154,7 +159,7 @@ class Session:
self.stats.captured += 1
self.controller.call('capture', arrived)
state = self.controller.state()
if self.pending or arrived-self.last_submit < 1000000/state['fps'] or not self.controller.call('gate', arrived, 1):
if self.reconfiguring or len(self.pending) >= 3 or arrived-self.last_submit < 1000000/state['fps'] or not self.controller.call('gate', arrived, 1):
self.stats.skipped += 1
return 0
self.pending[pts] = dict(cap=capture, arr=arrived, e0=arrived, tier=state['tier'], br=state['target'])
@@ -164,10 +169,20 @@ class Session:
self.stop_event.set()
return 0
def fresh_pipewire(self):
if 'portal' not in self.source:
return
error = C.create_string_buffer(1024)
fd = self.native.lib.fc_portal_refresh(self.source['portal'], error, len(error))
if fd < 0:
raise RuntimeError(error.value.decode(errors='replace'))
self.source['fd'] = fd
def produce(self):
native, capture = self.native.lib, None
try:
encoder = "x264enc" if self.src == "test" and self.native.has("x264enc") else self.agent.encoder
self.fresh_pipewire()
description = pipeline(self.source, sys.platform, encoder, self.w, self.h, self.fps, self.bitrate, self.codec)
self.callback = GATE(self.gate)
error = C.create_string_buffer(1024)
@@ -175,8 +190,14 @@ class Session:
if not capture:
raise RuntimeError(error.value.decode(errors='replace'))
last_update = last_stats = self.native.now()
last_output = last_update
last_output = last_config = last_update
applied = (self.w, self.h, self.bitrate)
wanted = applied
while not self.stop_event.is_set():
while not self.test_inputs.empty():
event = self.test_inputs.get_nowait()
native.fc_capture_test(capture, int(event.get('i', 0)) & 0xffffffff)
self.stats.input(event)
output = Encoded()
result = native.fc_capture_pull(capture, C.byref(output))
now = self.native.now()
@@ -185,9 +206,9 @@ class Session:
if result:
last_output = now
with self.lock:
# Exactly one raw frame is in flight and B-frames are disabled.
# x264 may offset PTS; associate by that single frame,
# keeping the actual pre-encode capture timestamp.
# At most three raw frames are in flight; no B-frames.
# x264 offsets PTS, so match the FIFO encode order while
# retaining the actual pre-encode capture timestamp.
raw_pts = next(iter(self.pending), None)
record = self.pending.get(raw_pts)
if record is None:
@@ -208,12 +229,36 @@ class Session:
native.fc_capture_key(capture)
if now-last_update >= 100000:
target = self.controller.update(now)
# x264 supports live bitrate changes. Hardware elements
# advertise their mutability; the initial bitrate always
# applies. Frame gating remains active for every encoder.
state = self.controller.state()
w, h = dimensions(self.w, self.h, max(320, int(max(self.w, self.h)*state['scale'])))
if target and self.codec == 'h264' and encoder == 'x264enc':
native.fc_capture_bitrate(capture, target)
self.bitrate = target
# Hardware properties are not uniformly mutable in PLAYING.
# Drain then reopen the pipeline at a keyframe when its
# budget changes materially. The portal fd/session stays
# alive, so this does not bypass or repeat user consent.
desired_bitrate = target or applied[2]
bitrate_change = self.codec == 'h264' and encoder != 'x264enc' and abs(desired_bitrate-applied[2]) > applied[2]*.2
if not self.reconfiguring and now-last_config >= 1000000 and ((w, h) != applied[:2] or bitrate_change):
wanted = (w, h, desired_bitrate)
with self.lock:
self.reconfiguring = True
last_update = now
if self.reconfiguring:
with self.lock:
drained = not self.pending
if drained:
native.fc_capture_close(capture)
capture = None
self.fresh_pipewire()
description = pipeline(self.source, sys.platform, encoder, wanted[0], wanted[1], self.fps, wanted[2], self.codec)
capture = native.fc_capture_open(description.encode(), self.callback, error, len(error))
if not capture:
raise RuntimeError(error.value.decode(errors='replace'))
applied, self.bitrate, last_config = wanted, wanted[2], now
with self.lock:
self.reconfiguring = False
if now-last_stats >= 1000000:
self.ws.send(dict(self.stats.summary(), t='stats', bitrate=self.bitrate,
size='%dx%d' % (self.w, self.h), tier=self.controller.state()['tier']))
@@ -235,7 +280,7 @@ class Session:
return
self.ws.send(dict(t='hello', r=self.key))
self.ws.send(dict(t='info', src=self.src, title=self.source.get('title', self.source.get('name', 'Test pattern')),
app='PC', codec=self.codec, input=self.src == 'test' or self.source.get('devices', 3) == 3,
app='PC', codec=self.codec, input=self.input_enabled,
inputMessage='Allow pointer and keyboard control in the host sharing dialog.',
aspect=self.w/self.h, warm=0))
with self.lock:
@@ -264,8 +309,23 @@ class Session:
if self.stop_event.is_set():
break
if self.input:
self.input.handle(m)
self.stats.input(m)
if not self.input_enabled and t != 'release':
continue
try:
self.input.handle(m)
self.stats.input(m)
except RuntimeError as e:
self.input_enabled = False
try:
self.input.release()
except RuntimeError:
pass
self.ws.send(dict(t='error', message=str(e)))
elif m.get('i'):
try:
self.test_inputs.put_nowait(m)
except queue.Full:
pass
finally:
self.end()
self.writer.join(6)
@@ -277,7 +337,8 @@ class Session:
self.stop_event.set()
self.ws.close()
with self.input_lock:
if self.input:
if self.input and not self.released:
self.released = True
try:
self.input.release()
except (RuntimeError, OSError):
@@ -294,6 +355,7 @@ class Agent:
self.windows = Windows() if sys.platform == 'win32' else None
self.encoder = 'mfh264enc' if self.windows else 'vah264enc' if native.has('vah264enc') else 'x264enc'
self.shutting_down = False
self.selection_generation = 0
def lists(self):
if self.windows:
@@ -316,6 +378,7 @@ class Agent:
if len(self.sources) >= 8:
raise ValueError('Stop a panel before sharing another source')
self.selecting, self.selection_error = True, ''
generation = self.selection_generation
def choose():
error = C.create_string_buffer(1024)
portal = self.native.lib.fc_portal_select(error, len(error))
@@ -323,7 +386,7 @@ class Agent:
self.selecting = False
if not portal:
self.selection_error = error.value.decode(errors='replace')
elif self.shutting_down:
elif self.shutting_down or generation != self.selection_generation:
self.native.lib.fc_portal_close(portal)
else:
fd, node, w, h, devices = [self.native.lib.fc_portal_value(portal, i) for i in range(5)]
@@ -334,6 +397,8 @@ class Agent:
def stop(self, src=None):
with self.lock:
if src is None:
self.selection_generation += 1
self.grants.revoke(src)
sessions = [s for s in self.sessions.values() if src is None or s.src == src]
for session in sessions:
@@ -452,8 +517,19 @@ class Handler(BaseHTTPRequestHandler):
# under the same lock as redemption, before any capture starts.
session = Session(agent, WebSocket(self), key, source, query)
for old in list(agent.sessions.values()):
if old.key == key:
if old.key == key or old.src == source['src']:
if old.key != key:
agent.grants.keys.pop(old.key, None)
try:
old.ws.send(dict(t='close'))
except OSError:
pass
old.end()
if old.writer:
old.writer.join(6)
if old.writer.is_alive():
session.controller.close()
raise ValueError('The previous capture is still stopping; retry shortly')
ident = agent.next_id
agent.next_id += 1
agent.sessions[ident] = session
@@ -468,6 +544,11 @@ class Handler(BaseHTTPRequestHandler):
except (OSError, ValueError, RuntimeError):
session.end()
finally:
session.end()
if session.writer:
session.writer.join(6)
if not session.writer or not session.writer.is_alive():
session.controller.close()
with agent.lock:
agent.sessions.pop(ident, None)
self.close_connection = True
+31 -14
View File
@@ -45,12 +45,14 @@ class Native:
'fc_capture_pull': (C.c_int, [C.c_void_p, C.POINTER(Encoded)]),
'fc_capture_error': (C.c_char_p, [C.c_void_p]),
'fc_capture_bitrate': (None, [C.c_void_p, C.c_int]),
'fc_capture_key': (None, [C.c_void_p]), 'fc_capture_close': (None, [C.c_void_p]),
'fc_capture_key': (None, [C.c_void_p]),
'fc_capture_test': (None, [C.c_void_p, C.c_uint32]), 'fc_capture_close': (None, [C.c_void_p]),
}
if sys.platform.startswith('linux'):
signatures.update({
'fc_portal_select': (C.c_void_p, [C.c_char_p, C.c_int]),
'fc_portal_close': (None, [C.c_void_p]),
'fc_portal_refresh': (C.c_int, [C.c_void_p, C.c_char_p, C.c_int]),
'fc_portal_value': (C.c_int, [C.c_void_p, C.c_int]),
'fc_portal_input': (C.c_int, [C.c_void_p, C.c_int, C.c_double, C.c_double, C.c_int, C.c_int]),
})
@@ -113,7 +115,7 @@ def dimensions(w, h, maximum):
def pipeline(source, platform, encoder, w, h, fps, bitrate, codec='h264'):
"""Only locally constructed numeric source IDs enter the pipeline parser."""
if source['src'] == 'test':
capture = 'videotestsrc is-live=true pattern=ball'
capture = 'videotestsrc name=source is-live=true pattern=ball'
elif platform == 'win32':
kind, ident = source['src'].split(':')
if kind not in ('window', 'display') or not re.fullmatch(r'[0-9]+', ident):
@@ -124,7 +126,7 @@ def pipeline(source, platform, encoder, w, h, fps, bitrate, codec='h264'):
capture = 'pipewiresrc fd=%d path=%d do-timestamp=true' % (source['fd'], source['node'])
# The source gate runs before conversion/encoding. There is no leaky queue
# of H.264 frames; every encoded reference frame reaches the socket.
raw = '%s ! video/x-raw,framerate=%d/1 ! identity name=gate ! videoconvert ! videoscale ! video/x-raw,width=%d,height=%d' % (capture, fps, w, h)
raw = '%s ! video/x-raw,framerate=%d/1 ! identity name=gate ! videoconvert ! videoscale add-borders=false ! video/x-raw,width=%d,height=%d' % (capture, fps, w, h)
if codec == 'jpeg':
enc = 'jpegenc name=enc quality=80'
parse = ''
@@ -158,12 +160,16 @@ class PortalInput:
raise RuntimeError('The desktop portal refused input; check the sharing permission')
def release(self):
for b in list(self.buttons):
self.emit(1, code=b, down=0)
self.buttons.discard(b)
for k in list(self.keys):
self.emit(3, code=k, down=0)
self.keys.discard(k)
error = None
for values, kind in ((self.buttons, 1), (self.keys, 3)):
for code in list(values):
try:
self.emit(kind, code=code, down=0)
values.discard(code)
except RuntimeError as e:
error = e
if error:
raise error
def handle(self, m):
t = m['t']
@@ -282,12 +288,22 @@ class WindowsInput:
self.buttons, self.keys = set(), set()
def release(self):
error = None
for button in list(self.buttons):
self.host.send(mouse=(0, 0, 0, {0: 4, 1: 0x40, 2: 0x10}[button]))
self.buttons.discard(button)
try:
self.host.send(mouse=(0, 0, 0, {0: 4, 1: 0x40, 2: 0x10}[button]))
self.buttons.discard(button)
except RuntimeError as e:
error = e
for vk in list(self.keys):
self.host.send(key=(vk, 0, 2))
self.keys.discard(vk)
try:
extended = 1 if vk in (33, 34, 35, 36, 37, 38, 39, 40, 45, 46, 91, 92, 163, 165) else 0
self.host.send(key=(vk, 0, 2 | extended))
self.keys.discard(vk)
except RuntimeError as e:
error = e
if error:
raise error
def handle(self, m):
t, u = m['t'], self.host.user
@@ -328,7 +344,8 @@ class WindowsInput:
if re.fullmatch(r'Key[A-Z]', code) or re.fullmatch(r'Digit[0-9]', code):
vk = ord(code[-1])
if vk:
self.host.send(key=(vk, 0, 0 if down else 2))
extended = 1 if vk in (33, 34, 35, 36, 37, 38, 39, 40, 45, 46, 91, 92, 163, 165) else 0
self.host.send(key=(vk, 0, extended | (0 if down else 2)))
(self.keys.add if down else self.keys.discard)(vk)
elif down and len(str(m.get('key', ''))) == 1:
self.handle({'t': 'text', 's': m['key']})
+24 -1
View File
@@ -6,11 +6,13 @@ are different. No second viewer, SSH supervisor or benchmark launcher.
"""
import os
import sys
import subprocess
from frame_macview import MacView, MacViewError, ROOT
from frame_pc_capture import LIBRARY, NATIVE
class PCView(MacView):
viewer_profile = 'pc-view'
host = 'windows' if sys.platform == 'win32' else 'linux'
def __init__(self, *args, **kwargs):
@@ -39,7 +41,9 @@ class PCView(MacView):
env = super().agent_environment()
env['GST_PLUGIN_PATH_1_0'] = str(NATIVE / 'lib' / 'gstreamer-1.0')
env['GST_PLUGIN_SYSTEM_PATH_1_0'] = ''
env['GST_REGISTRY_1_0'] = str(NATIVE.parent / 'registry.bin') if os.access(NATIVE.parent, os.W_OK) else os.path.join(os.path.expanduser('~'), '.cache', 'frame-control-gst.bin')
cache = os.path.join(os.path.expanduser('~'), '.cache', 'frame-control')
os.makedirs(cache, exist_ok=True)
env['GST_REGISTRY_1_0'] = os.path.join(cache, 'gstreamer-registry.bin')
env['GST_REGISTRY_FORK'] = 'no'
if sys.platform == 'win32':
env['PATH'] = str(NATIVE / 'bin') + os.pathsep + env.get('PATH', '')
@@ -47,6 +51,25 @@ class PCView(MacView):
env['LD_LIBRARY_PATH'] = str(NATIVE / 'lib') + os.pathsep + env.get('LD_LIBRARY_PATH', '')
return env
def shutdown(self):
self.closing = True
try:
self.stop()
except MacViewError:
pass
if self.tunnel and self.tunnel.poll() is None:
self.tunnel.terminate()
if self.agent and self.agent.poll() is None:
# EOF gives the host a chance to release held input and portal
# sessions, including on Windows where terminate is uncatchable.
if self.agent.stdin:
self.agent.stdin.close()
try:
self.agent.wait(10)
except subprocess.TimeoutExpired:
self.agent.kill()
self.agent.wait(5)
def _show(self, src, quality, width, height):
if src.startswith('separate:'):
raise MacViewError('Separate virtual displays are a Mac-only feature. Choose a window or screen.')
+2 -1
View File
@@ -1234,7 +1234,8 @@ def _sweep_one(prefix, d):
# The tunnel gets its own connection: the shared master's options would win
# over anything added after them.
macview = frame_pcview.host_view(["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8"],
lambda remote, stdin=None, timeout=30: ssh(remote, stdin=stdin, timeout=timeout),
lambda remote, stdin=None, timeout=30: ssh(remote, stdin=stdin.encode("utf-8") if stdin is not None else None,
timeout=timeout, text=False).decode("utf-8", errors="replace"),
FRAME, track=_live_tunnels.add)