From a992f6d8965531270c7d15c5bbf646ef682ae7f0 Mon Sep 17 00:00:00 2001 From: saphid <4596216+saphid@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:44:12 +1000 Subject: [PATCH] PC host: review fixes (idle watchdog, controller close race, malformed input, loader path) - The 10 s watchdog no longer ends idle sessions. PipeWire and WGC deliver frames only on damage; fail only when a frame entered the encoder and never came out, or the pipeline never produced its initial frame. - Controller.state/call/update are inert after close(), so a /status poll racing Stop cannot dereference the freed native controller. - Malformed viewer fields drop the event instead of ending the session. - PATH/LD_LIBRARY_PATH prepends add no empty (current directory) entry. - portal.c: document fd ownership. pipewiresrc dups the fd it is given (gstpipewirecore.c, F_DUPFD_CLOEXEC), so the portal closing its own fd is correct; clear it after close. Co-Authored-By: Claude Opus 5.5 (1M context) --- desktop/portal.c | 11 +++++-- tests/test_pcview.py | 70 ++++++++++++++++++++++++++++++++++++++++++ ui/frame_pc_agent.py | 44 +++++++++++++++++++++----- ui/frame_pc_capture.py | 14 +++++++-- ui/frame_pcview.py | 10 ++++-- 5 files changed, 135 insertions(+), 14 deletions(-) diff --git a/desktop/portal.c b/desktop/portal.c index cdf7cc7..fbd11d2 100644 --- a/desktop/portal.c +++ b/desktop/portal.c @@ -61,12 +61,19 @@ FC_API void fc_portal_close(Portal *p) { NULL,NULL,G_DBUS_CALL_FLAGS_NONE,2000,NULL,NULL); if(r)g_variant_unref(r); } - if(p->fd>=0)close(p->fd); + if(p->fd>=0){close(p->fd);p->fd=-1;} g_free(p->session);if(p->bus)g_object_unref(p->bus); if(p->context)g_main_context_unref(p->context);g_free(p); } /* Each pipeline gets a fresh restricted PipeWire connection. A dup of a - * previously consumed protocol socket is not a new connection. */ + * previously consumed protocol socket is not a new connection. + * Ownership: the Portal owns p->fd for its whole life and is its only closer. + * pipewiresrc never takes the fd it is given: its core connects with + * pw_context_connect_fd(ctx, fcntl(fd, F_DUPFD_CLOEXEC, 3), ...) and that + * duplicate is what PipeWire closes on teardown (src/gst/gstpipewirecore.c, + * unchanged from 0.3.19 through 1.x). So closing p->fd here is not a double + * close; not closing it would leak one socket per pipeline reopen. The + * caller closes the previous pipeline before asking for a new fd. */ FC_API int fc_portal_refresh(Portal *p,char *error,int capacity) { GVariantBuilder b;g_variant_builder_init(&b,G_VARIANT_TYPE_VARDICT); GUnixFDList *fds=NULL;GError *e=NULL; diff --git a/tests/test_pcview.py b/tests/test_pcview.py index dbe6c0c..07342ea 100644 --- a/tests/test_pcview.py +++ b/tests/test_pcview.py @@ -126,6 +126,76 @@ class AdapterTests(unittest.TestCase): session.input.release.assert_called_once() self.assertTrue(session.stop_event.is_set()) + def test_idle_source_is_not_a_stall(self): + # PipeWire/WGC deliver frames only on damage: a static desktop sends + # its first frame and then nothing. That must not end the session. + session = object.__new__(agent.Session) + session.lock, session.pending = threading.RLock(), {} + session.stats = Stats(lambda: 0) + session.stats.captured = 1 + self.assertIsNone(session.stalled(60000000, 0, 0)) + # A frame that entered the encoder and never came out is a fault. + session.pending[5] = dict(cap=0, arr=0, e0=1000000, tier=0, br=0) + self.assertIn('encoder', session.stalled(12000000, 0, 0)) + self.assertIsNone(session.stalled(10000000, 0, 0)) + # No initial frame at all after the pipeline opened is a fault. + session.pending.clear() + self.assertIn('capture', session.stalled(12000000, 0, 1)) + self.assertIsNone(session.stalled(9000000, 0, 1)) + + def test_controller_is_inert_after_close(self): + lib = mock.Mock() + lib.fc_new.return_value = 1234 + lib.fc_value.return_value = 100 + controller = capture.Controller(lib, 60, 5000000) + controller.close() + lib.fc_free.assert_called_once_with(1234) + lib.reset_mock() + state = controller.state() + self.assertEqual((state['fps'], state['target'], state['scale']), (60, 0, 1)) + self.assertEqual(controller.call('gate', 1, 1), 0) + self.assertEqual(controller.update(1), 0) + controller.close() + self.assertEqual(lib.method_calls, []) # nothing touched the freed pointer + + def test_library_path_prepend_adds_no_empty_entry(self): + env = {} + frame_pcview.prepend(env, 'LD_LIBRARY_PATH', '/opt/fc/lib') + self.assertEqual(env['LD_LIBRARY_PATH'], '/opt/fc/lib') + frame_pcview.prepend(env, 'LD_LIBRARY_PATH', '/x') + self.assertEqual(env['LD_LIBRARY_PATH'], '/x' + os.pathsep + '/opt/fc/lib') + with mock.patch.dict(os.environ, clear=True): + env = frame_pcview.PCView([], lambda *a: None, 'frame').agent_environment() + for name in ('PATH', 'LD_LIBRARY_PATH'): + parts = env.get(name, 'x').split(os.pathsep) + self.assertNotIn('', parts, name) + + def test_malformed_viewer_input_is_dropped_not_fatal(self): + session = object.__new__(agent.Session) + session.lock, session.input_lock = threading.RLock(), threading.RLock() + session.stop_event, session.key_event = threading.Event(), threading.Event() + session.src, session.codec, session.key, session.w, session.h = 'display:1', 'h264', 'r', 100, 100 + session.source = dict(src='display:1', w=100, h=100, portal=1) + session.input_enabled, session.released = True, False + session.native = mock.Mock() + session.native.now.return_value = 0 + session.controller, session.agent = mock.Mock(), mock.Mock() + session.stats = Stats(lambda: 0) + session.produce = lambda: None + lib = mock.Mock() + lib.fc_portal_input.return_value = 1 + session.input = capture.PortalInput(mock.Mock(lib=lib), session.source) + session.ws = mock.Mock() + session.ws.receive.side_effect = [ + dict(t='m', x='left', y=0), dict(t='wheel', dx=[1]), dict(t='k', e='down', key=None, i='x'), + dict(t='fd', drop='many'), dict(t='rx', s=[1]), dict(t='m', x=.5, y=.5), ConnectionError('closed')] + with self.assertRaises(ConnectionError): + session.start() + # Every malformed message was skipped; the good move after them still ran. + moves = [c.args for c in lib.fc_portal_input.call_args_list if c.args[1] == 0] + self.assertAlmostEqual(moves[-1][2], .5*99) + self.assertEqual(session.ws.receive.call_count, 7) + def test_stats_keep_capture_time_and_bound_records(self): stats = Stats(lambda: 10000000) stats.input(dict(t='m', i=1, tv=9999990)) diff --git a/ui/frame_pc_agent.py b/ui/frame_pc_agent.py index 720d752..271c49e 100644 --- a/ui/frame_pc_agent.py +++ b/ui/frame_pc_agent.py @@ -27,6 +27,10 @@ from frame_pc_capture import Native, Controller, Windows, WindowsInput, PortalIn from frame_stream_stats import Stats +# Viewer fields are untrusted JSON: float('x'), int(None), m['t'] missing. +MALFORMED = (ValueError, TypeError, KeyError, OverflowError) + + class Grants: """Caller holds the agent lock, including replacement and Stop.""" def __init__(self, token): @@ -172,6 +176,22 @@ class Session: self.stop_event.set() return -1 + STALL = 10000000 + + def stalled(self, now, opened, captured_at_open): + """PipeWire and WGC deliver frames only on screen damage, so an idle + desktop legitimately produces no output. Fail only on evidence of a + fault: a frame entered the encoder and never came out, or the capture + never produced its initial frame after the pipeline opened.""" + with self.lock: + oldest = min((r['e0'] for r in self.pending.values()), default=None) + captured = self.stats.captured + if oldest is not None and now-oldest > self.STALL: + return 'The encoder stopped producing frames; check the video encoder' + if captured == captured_at_open and now-opened > self.STALL: + return 'No frames captured for 10 seconds; check capture permissions' + return None + def fresh_pipewire(self): if 'portal' not in self.source: return @@ -193,21 +213,24 @@ class Session: if not capture: raise RuntimeError(error.value.decode(errors='replace')) last_update = last_stats = self.native.now() - last_output = last_config = last_update + last_config = opened = last_update + captured_at_open = self.stats.captured 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) + try: + native.fc_capture_test(capture, int(event.get('i', 0)) & 0xffffffff) + self.stats.input(event) + except MALFORMED: + pass # a bad probe field drops the event, not the stream output = Encoded() result = native.fc_capture_pull(capture, C.byref(output)) now = self.native.now() if result < 0: raise RuntimeError(native.fc_capture_error(capture).decode(errors='replace')) if result: - last_output = now with self.lock: # At most three raw frames are in flight; no B-frames. # x264 offsets PTS, so match the FIFO encode order while @@ -224,8 +247,9 @@ class Session: with self.lock: f['wire'] = self.native.now() self.pending.pop(raw_pts, None) - if now-last_output > 10000000: - raise RuntimeError('No encoded frames for 10 seconds; check capture permissions and the encoder') + stall = self.stalled(now, opened, captured_at_open) + if stall: + raise RuntimeError(stall) if self.key_event.is_set(): self.key_event.clear() if self.codec == 'h264': @@ -260,6 +284,7 @@ class Session: if not capture: raise RuntimeError(error.value.decode(errors='replace')) applied, self.bitrate, last_config = wanted, wanted[2], now + opened, captured_at_open = now, self.stats.captured with self.lock: self.reconfiguring = False if now-last_stats >= 1000000: @@ -304,7 +329,10 @@ class Session: elif t == 'key-frame': self.key_event.set() elif t in ('rx', 'fd', 'clock'): - self.stats.report(m) + try: + self.stats.report(m) + except MALFORMED: + pass if t == 'rx' and isinstance(m.get('s'), int): self.controller.call('ack', m['s'] & 0xffffffff, self.native.now()) elif t in ('m', 'wheel', 'k', 'text', 'release'): @@ -317,6 +345,8 @@ class Session: try: self.input.handle(m) self.stats.input(m) + except MALFORMED: + pass # drop a malformed viewer event; keep the session except RuntimeError as e: self.input_enabled = False try: diff --git a/ui/frame_pc_capture.py b/ui/frame_pc_capture.py index 0e3366c..cb4f5e5 100644 --- a/ui/frame_pc_capture.py +++ b/ui/frame_pc_capture.py @@ -75,7 +75,9 @@ class Native: class Controller: def __init__(self, lib, fps, ceiling): - self.lib, self.lock = lib, threading.RLock() + # close() may run on the stream thread while /status reads state(); + # after close every entry point is a no-op, never a NULL dereference. + self.lib, self.lock, self.fps = lib, threading.RLock(), fps self.enabled = os.environ.get('FRAME_MAC_VIEW_ADAPT') != '0' self.ptr = lib.fc_new(fps, self.enabled) if not self.ptr: @@ -86,17 +88,23 @@ class Controller: def state(self): with self.lock: names = ['target', 'ceiling', 'tier', 'fps', 'scale', 'baseRtt', 'inFlight', 'slack'] - out = {k: self.lib.fc_value(self.ptr, i) for i, k in enumerate(names)} + if not self.ptr: + out = dict.fromkeys(names, 0) + out.update(fps=self.fps, scale=100) + else: + out = {k: self.lib.fc_value(self.ptr, i) for i, k in enumerate(names)} out.update(scale=out['scale'] / 100, baseRtt=out['baseRtt'] / 1000, slack=out['slack'] / 1000, adapt=self.enabled) return out def call(self, name, *args): with self.lock: - return getattr(self.lib, 'fc_' + name)(self.ptr, *args) + return getattr(self.lib, 'fc_' + name)(self.ptr, *args) if self.ptr else 0 def update(self, now): with self.lock: + if not self.ptr: + return 0 old = self.state() result = self.lib.fc_update(self.ptr, now) state = self.state() diff --git a/ui/frame_pcview.py b/ui/frame_pcview.py index 7994a46..782aa3e 100644 --- a/ui/frame_pcview.py +++ b/ui/frame_pcview.py @@ -11,6 +11,12 @@ from frame_macview import MacView, MacViewError, ROOT from frame_pc_capture import LIBRARY, NATIVE +def prepend(env, name, path): + """An empty entry means the current directory to the loader; never add one.""" + rest = env.get(name, '') + env[name] = str(path) + (os.pathsep + rest if rest else '') + + class PCView(MacView): viewer_profile = 'pc-view' host = 'windows' if sys.platform == 'win32' else 'linux' @@ -46,9 +52,9 @@ class PCView(MacView): 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', '') + prepend(env, 'PATH', NATIVE / 'bin') else: - env['LD_LIBRARY_PATH'] = str(NATIVE / 'lib') + os.pathsep + env.get('LD_LIBRARY_PATH', '') + prepend(env, 'LD_LIBRARY_PATH', NATIVE / 'lib') env['PIPEWIRE_MODULE_DIR'] = str(NATIVE / 'lib' / 'pipewire-0.3') env['SPA_PLUGIN_DIR'] = str(NATIVE / 'lib' / 'spa-0.2') env['PIPEWIRE_CONFIG_DIR'] = str(NATIVE / 'share' / 'pipewire')