mirror of
https://github.com/saphid/frame-control.git
synced 2026-10-06 04:04:21 +02:00
Retain the final idle-window update while the stream gate is closed
This commit is contained in:
1 parent
2466c5aa14
commit
da5c954ce4
5 files changed
+81
-17
No files matched your search
+23
-4
@@ -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);
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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')
|
||||
|
||||
+12
-9
@@ -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:
|
||||
|
||||
@@ -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 = ''
|
||||
|
||||
Reference in new issue
Block a user