From da5c954ce4216c418932ddccbb7d5bd53cd64a95 Mon Sep 17 00:00:00 2001 From: saphid <4596216+saphid@users.noreply.github.com> Date: Mon, 28 Sep 2026 23:05:23 +1000 Subject: [PATCH] Retain the final idle-window update while the stream gate is closed --- desktop/capture.c | 27 +++++++++++++++++++++++---- docs/pc-in-headset.md | 8 +++++--- tests/test_pcview.py | 40 ++++++++++++++++++++++++++++++++++++++++ ui/frame_pc_agent.py | 21 ++++++++++++--------- ui/frame_pc_capture.py | 2 +- 5 files changed, 81 insertions(+), 17 deletions(-) diff --git a/desktop/capture.c b/desktop/capture.c index 06e146d..adb28cb 100644 --- a/desktop/capture.c +++ b/desktop/capture.c @@ -8,7 +8,7 @@ typedef int (*Gate)(int stage, int64_t pts, int64_t capture, int64_t arrived); typedef struct { - GstElement *pipeline, *encoder, *sink; + GstElement *pipeline, *encoder, *sink, *raw_queue; Gate gate; GstSample *sample; GstMapInfo map; @@ -34,7 +34,19 @@ static GstPadProbeReturn probe(GstPad *pad,GstPadProbeInfo *info,gpointer data) } if(clock)gst_object_unref(clock); int stage=GPOINTER_TO_INT(g_object_get_data(G_OBJECT(pad),"stage")); - return c->gate(stage,(int64_t)GST_BUFFER_PTS(b),cap,now) ? GST_PAD_PROBE_OK : GST_PAD_PROBE_DROP; + if(stage==1) return c->gate(stage,(int64_t)GST_BUFFER_PTS(b),cap,now)>0 ? GST_PAD_PROBE_OK : GST_PAD_PROBE_DROP; + int phase=0; + for(;;) { + /* While congested, prefer a newer raw picture queued upstream. If + * nothing changed, retain this last picture until the gate opens; + * an idle window must not stay stale after a dropped final update. */ + guint queued=0; + if(phase && c->raw_queue)g_object_get(c->raw_queue,"current-level-buffers",&queued,NULL); + if(queued)return GST_PAD_PROBE_DROP; + int decision=c->gate(phase,(int64_t)GST_BUFFER_PTS(b),cap,now); + if(decision)return decision>0 ? GST_PAD_PROBE_OK : GST_PAD_PROBE_DROP; + phase=2;g_usleep(2000); + } } FC_API Capture *fc_capture_open(const char *pipeline,Gate gate,char *error,int capacity) { GError *e=NULL;Capture *c=g_new0(Capture,1);c->gate=gate; @@ -51,6 +63,7 @@ FC_API Capture *fc_capture_open(const char *pipeline,Gate gate,char *error,int c if(raw)gst_object_unref(raw);if(c->encoder)gst_object_unref(c->encoder); if(c->sink)gst_object_unref(c->sink);gst_object_unref(c->pipeline);g_free(c);return NULL; } + c->raw_queue=gst_bin_get_by_name(GST_BIN(c->pipeline),"raw_queue"); GstPad *p=gst_element_get_static_pad(raw,"src"); g_object_set_data(G_OBJECT(p),"stage",GINT_TO_POINTER(0)); gst_pad_add_probe(p,GST_PAD_PROBE_TYPE_BUFFER,probe,c,NULL);gst_object_unref(p);gst_object_unref(raw); @@ -64,7 +77,7 @@ FC_API int fc_capture_pull(Capture *c,Encoded *out) { if(c->mapped) {gst_buffer_unmap(gst_sample_get_buffer(c->sample),&c->map);c->mapped=0;} if(c->sample){gst_sample_unref(c->sample);c->sample=NULL;} GstBus *bus=gst_element_get_bus(c->pipeline); - GstMessage *m=gst_bus_pop_filtered(bus,GST_MESSAGE_ERROR|GST_MESSAGE_EOS);gst_object_unref(bus); + GstMessage *m=gst_bus_pop_filtered(bus,GST_MESSAGE_ERROR);gst_object_unref(bus); if(m) { if(GST_MESSAGE_TYPE(m)==GST_MESSAGE_ERROR) { GError *e=NULL;char *debug=NULL;gst_message_parse_error(m,&e,&debug); @@ -73,7 +86,12 @@ FC_API int fc_capture_pull(Capture *c,Encoded *out) { gst_message_unref(m);return -1; } c->sample=gst_app_sink_try_pull_sample(GST_APP_SINK(c->sink),100*GST_MSECOND); - if(!c->sample)return 0; + if(!c->sample) { + if(gst_app_sink_is_eos(GST_APP_SINK(c->sink))) { + g_strlcpy(c->error,"The capture source closed",sizeof(c->error));return -1; + } + return 0; + } GstBuffer *b=gst_sample_get_buffer(c->sample); if(!gst_buffer_map(b,&c->map,GST_MAP_READ))return 0; c->mapped=1;out->data=c->map.data;out->size=(int)c->map.size; @@ -102,5 +120,6 @@ FC_API void fc_capture_close(Capture *c) { gst_element_set_state(c->pipeline,GST_STATE_NULL); if(c->mapped)gst_buffer_unmap(gst_sample_get_buffer(c->sample),&c->map); if(c->sample)gst_sample_unref(c->sample); + if(c->raw_queue)gst_object_unref(c->raw_queue); gst_object_unref(c->sink);gst_object_unref(c->encoder);gst_object_unref(c->pipeline);g_free(c); } diff --git a/docs/pc-in-headset.md b/docs/pc-in-headset.md index 3b1fe73..66a5989 100644 --- a/docs/pc-in-headset.md +++ b/docs/pc-in-headset.md @@ -70,7 +70,9 @@ window resize/minimize, and non-US keyboard layouts. `rx`/`fd` timing reports and input messages are unchanged. - `desktop/controller.c` is the rate controller shared by the Mac Swift binding and the PC Python binding. Capture is gated **before** encoding; - encoded reference frames are never discarded. It keeps the Mac's bitrate + encoded reference frames are never discarded. A bounded raw-frame queue + keeps the newest picture, including the last update of an idle window, + until the gate opens. A native one-frame-source test covers that case. It keeps the Mac's bitrate demand protection and tier hysteresis. - PC records use the existing `Stats.swift` JSON schema, with bounded 4096-frame/512-input storage in `ui/frame_stream_stats.py`. The benchmark's @@ -146,13 +148,13 @@ same bounded shaping relay, without administrator privileges. - **Verified, Mac:** all 600 states in a 60-second congestion/recovery trace matched the original Swift controller. `tests/test_pc_controller.py` retains the original trace digest as a regression check. -- **Verified, real Frame, current agent at `cd20243`:** the repeated synthetic +- **Verified, real Frame, agent at `cd20243`:** the repeated synthetic probe drew 400 frames at 39.3 fps, content p50/p95 15.1/28.6 ms, and synthetic input-to-drawn p50 76.6 ms. Test clicks now change the pattern color before injection is timestamped. The frame-rate/late-frame targets still failed; this remains a Frame-hosted x264 test through a Mac relay, not a desktop or physical-laser measurement. Helper exit 0 and cleanup succeeded. - [Current probe result](../bench/results/2026-09-28-cd20243-pc-agent-frame-arm64-final.json). + [Latest device probe result](../bench/results/2026-09-28-cd20243-pc-agent-frame-arm64-final.json). - **Untested:** real Windows WGC → Media Foundation → Frame; real Linux portal → PipeWire → VA-API/x264 → Frame; physical laser input on either. No benchmark numbers for those desktop paths are claimed. diff --git a/tests/test_pcview.py b/tests/test_pcview.py index f18db17..826c3ab 100644 --- a/tests/test_pcview.py +++ b/tests/test_pcview.py @@ -210,6 +210,46 @@ lib.pw_main_loop_destroy(loop) result = subprocess.run([sys.executable, '-c', code], env=env, capture_output=True, text=True, timeout=15) self.assertEqual(result.returncode, 0, result.stderr) + def test_last_raw_picture_survives_a_closed_gate(self): + # A source sends ONE changed frame, then becomes idle. It must reach + # the encoder after congestion clears, without needing another update. + code = """ +import ctypes as C, sys, time +from frame_pc_capture import Native, GATE, Encoded, pipeline +native = Native() +start = None +closing = False +phases = [] +def gate(stage, pts, capture, arrived): + global start + phases.append(stage) + if closing: return -1 + if start is None: start = time.monotonic() + return int(stage == 1 or time.monotonic()-start >= .15) +callback = GATE(gate) +p = pipeline(dict(src='test'), sys.platform, 'x264enc', 320, 180, 30, 300000) +p = p.replace('is-live=true', 'is-live=true num-buffers=1') +error = C.create_string_buffer(1024) +handle = native.lib.fc_capture_open(p.encode(), callback, error, len(error)) +assert handle, error.value +try: + frame = Encoded() + result = 0 + deadline = time.monotonic()+5 + while not result and time.monotonic() < deadline: + result = native.lib.fc_capture_pull(handle, C.byref(frame)) + assert result == 1, native.lib.fc_capture_error(handle) + assert frame.key and frame.size > 17 + assert phases.count(0) == 1 and 2 in phases, phases +finally: + closing = True + native.lib.fc_capture_close(handle) +""" + env = frame_pcview.PCView([], lambda *a: None, 'frame').agent_environment() + env['PYTHONPATH'] = str(ROOT / 'ui') + result = subprocess.run([sys.executable, '-c', code], env=env, capture_output=True, text=True, timeout=10) + self.assertEqual(result.returncode, 0, result.stderr) + def test_auth_and_real_h264_stats(self): self.assertEqual(self.request('/status')[0], 403) status, body = self.request('/ticket?src=test&k='+self.token, 'POST') diff --git a/ui/frame_pc_agent.py b/ui/frame_pc_agent.py index a99c0c5..720d752 100644 --- a/ui/frame_pc_agent.py +++ b/ui/frame_pc_agent.py @@ -150,24 +150,27 @@ class Session: # Exceptions cannot cross a ctypes callback boundary. try: with self.lock: - if self.stop_event.is_set(): - return 0 + if self.stop_event.is_set() or self.reconfiguring: + return -1 if stage == 1: if pts in self.pending: self.pending[pts]['e0'] = self.native.now() return 1 - self.stats.captured += 1 - self.controller.call('capture', arrived) + if stage == 0: + self.stats.captured += 1 + self.controller.call('capture', arrived) + now = self.native.now() state = self.controller.state() - 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 + if len(self.pending) >= 3 or now-self.last_submit < 1000000/state['fps'] or not self.controller.call('gate', now, int(stage == 0)): + if stage == 0: + self.stats.skipped += 1 return 0 - self.pending[pts] = dict(cap=capture, arr=arrived, e0=arrived, tier=state['tier'], br=state['target']) - self.last_submit = arrived + self.pending[pts] = dict(cap=capture, arr=arrived, e0=now, tier=state['tier'], br=state['target']) + self.last_submit = now return 1 except Exception: self.stop_event.set() - return 0 + return -1 def fresh_pipewire(self): if 'portal' not in self.source: diff --git a/ui/frame_pc_capture.py b/ui/frame_pc_capture.py index da002d5..671c1cd 100644 --- a/ui/frame_pc_capture.py +++ b/ui/frame_pc_capture.py @@ -126,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 add-borders=false ! video/x-raw,width=%d,height=%d' % (capture, fps, w, h) + raw = '%s ! video/x-raw,framerate=%d/1 ! queue name=raw_queue leaky=downstream max-size-buffers=1 max-size-bytes=0 max-size-time=0 ! 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 = ''