mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 05:02:50 +02:00
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) <noreply@anthropic.com>
This commit is contained in:
1 parent
d9d538690e
commit
a992f6d896
5 files changed
+135
-14
No files matched your search
+9
-2
@@ -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;
|
||||
|
||||
@@ -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))
|
||||
|
||||
+37
-7
@@ -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:
|
||||
|
||||
+11
-3
@@ -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()
|
||||
|
||||
+8
-2
@@ -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')
|
||||
|
||||
Reference in new issue
Block a user