Skip to content

Commit d306e5f

Browse files
committed
fix: stabilize v4l2 loopback viewer output
1 parent 4bde0e7 commit d306e5f

2 files changed

Lines changed: 462 additions & 31 deletions

File tree

publish.py

Lines changed: 219 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,10 @@ def normalize_signaling_server_url(server_url: str) -> str:
142142
"high-444": "f40032",
143143
}
144144

145+
# avdec_h264 is the established viewer decoder and recovers cleanly from RTP
146+
# loss. Keep OpenH264 as a fallback for minimal images without gst-libav.
147+
H264_VIEWER_DECODER_FALLBACKS = ("avdec_h264", "openh264dec")
148+
145149
def sanitize_profile_level_id(value: Optional[str]) -> Optional[str]:
146150
"""Normalize a profile-level-id string to lowercase hex without 0x prefix."""
147151
if not value:
@@ -620,6 +624,59 @@ def request_pad_compat(element, template_name: str):
620624
return element.get_request_pad(template_name)
621625

622626

627+
def configure_live_input_selector(selector) -> Dict[str, Any]:
628+
"""Configure an input-selector for live sources with unrelated timestamps.
629+
630+
Viewer idle sources start with the pipeline, while WebRTC sources are added
631+
later and commonly begin with an earlier running time. The default
632+
``active-segment`` synchronization can therefore leave the newly selected
633+
source behind the idle segment. Clock synchronization keeps inactive live
634+
pads aligned so switching does not stall downstream elements such as
635+
``videorate``.
636+
637+
Properties are detected individually because older GStreamer releases do
638+
not provide all of them (notably ``drop-backwards``).
639+
"""
640+
requested = (
641+
("sync-streams", True),
642+
("sync-mode", "clock"),
643+
("cache-buffers", True),
644+
("drop-backwards", True),
645+
)
646+
applied: Dict[str, Any] = {}
647+
for property_name, value in requested:
648+
try:
649+
if selector.find_property(property_name):
650+
selector.set_property(property_name, value)
651+
applied[property_name] = value
652+
except Exception:
653+
# Keep the selector usable on vendor or legacy plugin variants
654+
# whose advertised properties reject newer values.
655+
continue
656+
return applied
657+
658+
659+
def build_v4l2sink_sink_description(
660+
device: str,
661+
io_mode: int,
662+
caps: str,
663+
) -> str:
664+
"""Build the common sink after already-normalized selector inputs.
665+
666+
Each selector input has identical raw-video caps. Keeping a second
667+
``videorate`` after the selector makes its timestamp history span source
668+
changes, which can suppress every frame after an idle-to-WebRTC switch.
669+
"""
670+
return (
671+
"queue max-size-buffers=2 leaky=downstream ! "
672+
"videoconvert ! "
673+
f"{caps} ! "
674+
"identity name=viewer_v4l2sink_drop_allocation drop-allocation=true ! "
675+
f"v4l2sink name=viewer_v4l2sink device={device} "
676+
f"io-mode={io_mode} sync=false"
677+
)
678+
679+
623680
def _run_v4l2_h264_encoder_probe() -> Tuple[bool, str]:
624681
"""Exercise a few small system-memory frames through v4l2h264enc."""
625682
probe = (
@@ -3329,6 +3386,7 @@ def __init__(self, params):
33293386
self.v4l2sink_io_mode = clamp_int(getattr(params, "v4l2sink_io_mode", 0), 0, 5)
33303387
self.v4l2sink_device = resolve_v4l2sink_device(self.v4l2sink) if self.v4l2sink else None
33313388
self.v4l2sink_selector = None
3389+
self.v4l2sink_selector_kind = None
33323390
self.v4l2sink_sink_bin = None
33333391
self.v4l2sink_sources = {}
33343392
self.v4l2sink_current_pad = None
@@ -4998,6 +5056,11 @@ def _prime_viewer_display(self):
49985056

49995057
def _ensure_display_chain(self):
50005058
"""Ensure viewer display selector and sink are ready."""
5059+
# A V4L2 loopback receiver is an alternate output, not a second local
5060+
# preview. Building both chains wastes scarce Pi resources and can make
5061+
# the physical display sink compete with the loopback sink for frames.
5062+
if getattr(self, "v4l2sink_device", None):
5063+
return None
50015064
if not self.pipe:
50025065
return None
50035066

@@ -5189,17 +5252,43 @@ def _ensure_v4l2sink_chain(self):
51895252
self._v4l2sink_chain_unavailable_reason = reason
51905253
return None
51915254

5192-
selector_factory_names = ("input-selector", "inputselector")
5255+
selector_factory_names = []
5256+
if gst_element_supports_property("compositor", "force-live"):
5257+
selector_factory_names.append("compositor")
5258+
selector_factory_names.extend(("input-selector", "inputselector"))
51935259
for selector_factory in selector_factory_names:
5194-
selector = Gst.ElementFactory.make(selector_factory, selector_name)
5260+
if selector_factory == "compositor":
5261+
selector = None
5262+
try:
5263+
factory = Gst.ElementFactory.find("compositor")
5264+
create_with_properties = getattr(factory, "create_with_properties", None)
5265+
if create_with_properties:
5266+
selector = create_with_properties(
5267+
["name", "force-live"],
5268+
[selector_name, True],
5269+
)
5270+
except Exception:
5271+
selector = None
5272+
else:
5273+
selector = Gst.ElementFactory.make(selector_factory, selector_name)
51955274
if selector:
51965275
self.v4l2sink_selector = selector
5276+
self.v4l2sink_selector_kind = (
5277+
"compositor" if selector_factory == "compositor" else "selector"
5278+
)
51975279
print(f"[v4l2sink] Using selector factory `{selector_factory}`")
5198-
try:
5199-
if selector.find_property("cache-buffers"):
5200-
selector.set_property("cache-buffers", True)
5201-
except Exception:
5202-
pass
5280+
if self.v4l2sink_selector_kind == "compositor":
5281+
for property_name, value in (
5282+
("background", "black"),
5283+
("ignore-inactive-pads", True),
5284+
):
5285+
try:
5286+
if selector.find_property(property_name):
5287+
selector.set_property(property_name, value)
5288+
except Exception:
5289+
pass
5290+
else:
5291+
configure_live_input_selector(selector)
52035292
break
52045293

52055294
if not self.v4l2sink_selector:
@@ -5216,13 +5305,10 @@ def _ensure_v4l2sink_chain(self):
52165305
f"video/x-raw,format={self.v4l2sink_format},width=(int){self.v4l2sink_width},"
52175306
f"height=(int){self.v4l2sink_height},framerate=(fraction){self.v4l2sink_fps}/1"
52185307
)
5219-
sink_desc = (
5220-
"queue max-size-buffers=2 leaky=downstream ! "
5221-
"videorate ! videoscale ! videoconvert ! "
5222-
f"{caps} ! "
5223-
"identity name=viewer_v4l2sink_drop_allocation drop-allocation=true ! "
5224-
f"v4l2sink name=viewer_v4l2sink device={self.v4l2sink_device} "
5225-
f"io-mode={self.v4l2sink_io_mode} sync=false"
5308+
sink_desc = build_v4l2sink_sink_description(
5309+
self.v4l2sink_device,
5310+
self.v4l2sink_io_mode,
5311+
caps,
52265312
)
52275313
self.v4l2sink_sink_bin = Gst.parse_bin_from_description(sink_desc, True)
52285314
self.v4l2sink_sink_bin.set_name("viewer_v4l2sink_sink_bin")
@@ -5249,6 +5335,8 @@ def _ensure_v4l2sink_splash_sources(self):
52495335
"""Create idle/blank sources for V4L2 sink output."""
52505336
if not self.pipe or not self.v4l2sink_selector:
52515337
return
5338+
if self.v4l2sink_selector_kind == "compositor":
5339+
return
52525340

52535341
caps = (
52545342
f"video/x-raw,format={self.v4l2sink_format},width=(int){self.v4l2sink_width},"
@@ -5269,7 +5357,12 @@ def _ensure_v4l2sink_splash_sources(self):
52695357
except Exception as exc:
52705358
printwarn(f"Failed to initialize V4L2 sink blank source: {exc}")
52715359

5272-
def _link_v4l2sink_bin(self, bin_obj: Gst.Bin, label: str):
5360+
def _link_v4l2sink_bin(
5361+
self,
5362+
bin_obj: Gst.Bin,
5363+
label: str,
5364+
wait_until_ready: bool = False,
5365+
):
52735366
"""Connect a bin's output to the V4L2 sink selector."""
52745367
if not self.v4l2sink_selector:
52755368
raise RuntimeError("V4L2 sink selector is unavailable")
@@ -5307,14 +5400,65 @@ def _link_v4l2sink_bin(self, bin_obj: Gst.Bin, label: str):
53075400
self.v4l2sink_selector.release_request_pad(selector_pad)
53085401
raise RuntimeError(f"Failed to link V4L2 sink source '{label}': {link_result}")
53095402

5310-
self.v4l2sink_sources[label] = {
5403+
source_info = {
53115404
"bin": bin_obj,
53125405
"selector_pad": selector_pad,
53135406
"src_pad": src_pad,
5407+
"ready": not wait_until_ready,
53145408
}
5409+
self.v4l2sink_sources[label] = source_info
5410+
5411+
if self.v4l2sink_selector_kind == "compositor":
5412+
try:
5413+
selector_pad.set_property("zorder", 0 if label == "blank" else 1)
5414+
selector_pad.set_property("alpha", 1.0 if label == "blank" else 0.0)
5415+
except Exception as exc:
5416+
printwarn(f"Failed to configure V4L2 compositor pad '{label}': {exc}")
5417+
5418+
if wait_until_ready:
5419+
try:
5420+
source_info["ready_probe_id"] = src_pad.add_probe(
5421+
Gst.PadProbeType.BUFFER,
5422+
self._on_v4l2sink_source_buffer,
5423+
(label, source_info),
5424+
)
5425+
except Exception as exc:
5426+
# A pad probe is available on all supported GStreamer versions,
5427+
# but retain the old immediate-switch behavior on vendor forks.
5428+
source_info["ready"] = True
5429+
printwarn(
5430+
f"Could not wait for V4L2 sink source '{label}' to produce a frame: {exc}"
5431+
)
53155432
bin_obj.sync_state_with_parent()
53165433
return selector_pad
53175434

5435+
def _on_v4l2sink_source_buffer(self, pad, probe_info, user_data):
5436+
"""Schedule a source switch after the remote branch produces a frame."""
5437+
label, source_info = user_data
5438+
source_info["ready_probe_id"] = None
5439+
print(f"[v4l2sink] Source '{label}' produced its first frame")
5440+
event_loop = getattr(self, "event_loop", None)
5441+
if event_loop and event_loop.is_running():
5442+
event_loop.call_soon_threadsafe(
5443+
self._mark_v4l2sink_source_ready,
5444+
label,
5445+
source_info,
5446+
)
5447+
else:
5448+
self._mark_v4l2sink_source_ready(label, source_info)
5449+
return Gst.PadProbeReturn.REMOVE
5450+
5451+
def _mark_v4l2sink_source_ready(self, label: str, source_info: Dict[str, Any]):
5452+
"""Activate a pending source if it is still the registered generation."""
5453+
if self.v4l2sink_sources.get(label) is not source_info:
5454+
return False
5455+
5456+
source_info["ready"] = True
5457+
if source_info.pop("activate_when_ready", False):
5458+
if self._activate_v4l2sink_source(label):
5459+
self.v4l2sink_state = "remote"
5460+
return False
5461+
53185462
def _activate_v4l2sink_source(self, label: str) -> bool:
53195463
"""Activate a registered source on the V4L2 sink selector."""
53205464
source = self.v4l2sink_sources.get(label)
@@ -5327,7 +5471,13 @@ def _activate_v4l2sink_source(self, label: str) -> bool:
53275471
return True
53285472

53295473
try:
5330-
self.v4l2sink_selector.set_property("active-pad", pad)
5474+
if self.v4l2sink_selector_kind == "compositor":
5475+
for source_label, source_info in self.v4l2sink_sources.items():
5476+
source_info["selector_pad"].set_property(
5477+
"alpha", 1.0 if source_label == label else 0.0
5478+
)
5479+
else:
5480+
self.v4l2sink_selector.set_property("active-pad", pad)
53315481
self.v4l2sink_current_pad = pad
53325482
print(f"[v4l2sink] Activated source '{label}'")
53335483
return True
@@ -5338,11 +5488,25 @@ def _activate_v4l2sink_source(self, label: str) -> bool:
53385488
def _set_v4l2sink_mode(self, mode: str, remote_label: Optional[str] = None):
53395489
"""Switch to the appropriate V4L2 sink source."""
53405490
if mode == "remote" and remote_label:
5341-
if remote_label in self.v4l2sink_sources:
5342-
self._activate_v4l2sink_source(remote_label)
5343-
self.v4l2sink_state = "remote"
5491+
source = self.v4l2sink_sources.get(remote_label)
5492+
if source:
5493+
if not source.get("ready", True):
5494+
source["activate_when_ready"] = True
5495+
print(f"[v4l2sink] Waiting for source '{remote_label}' to produce a frame")
5496+
return
5497+
if self._activate_v4l2sink_source(remote_label):
5498+
self.v4l2sink_state = "remote"
53445499
return
53455500
printwarn(f"V4L2 sink remote source '{remote_label}' not registered")
5501+
if self.v4l2sink_selector_kind == "compositor":
5502+
try:
5503+
for source in self.v4l2sink_sources.values():
5504+
source["selector_pad"].set_property("alpha", 0.0)
5505+
self.v4l2sink_current_pad = None
5506+
self.v4l2sink_state = "idle"
5507+
except Exception as exc:
5508+
printwarn(f"Failed to blank V4L2 sink output: {exc}")
5509+
return
53465510
if "blank" in self.v4l2sink_sources:
53475511
self._activate_v4l2sink_source("blank")
53485512
self.v4l2sink_state = "idle"
@@ -5367,6 +5531,14 @@ def _release_v4l2sink_source(self, label: str):
53675531
pass
53685532

53695533
selector_pad = source.get("selector_pad")
5534+
src_pad = source.get("src_pad")
5535+
ready_probe_id = source.get("ready_probe_id")
5536+
if src_pad and ready_probe_id:
5537+
try:
5538+
src_pad.remove_probe(ready_probe_id)
5539+
except Exception:
5540+
pass
5541+
53705542
if selector_pad:
53715543
try:
53725544
self.v4l2sink_selector.release_request_pad(selector_pad)
@@ -5427,6 +5599,7 @@ def _reset_v4l2sink_chain_state(self):
54275599
pass
54285600

54295601
self.v4l2sink_selector = None
5602+
self.v4l2sink_selector_kind = None
54305603
self.v4l2sink_sink_bin = None
54315604
self.v4l2sink_sources = {}
54325605
self.v4l2sink_current_pad = None
@@ -5803,6 +5976,9 @@ def _activate_display_source(self, label: str) -> bool:
58035976

58045977
def _set_display_mode(self, mode: str, remote_label: Optional[str] = None):
58055978
"""Switch to the appropriate splash/remote source."""
5979+
if getattr(self, "v4l2sink_device", None):
5980+
self._set_v4l2sink_mode(mode, remote_label=remote_label)
5981+
return
58065982
print(f"[display] Switching mode -> {mode} (remote={remote_label})")
58075983
if self._display_direct_mode:
58085984
if mode == "remote" and remote_label:
@@ -7063,7 +7239,9 @@ def get_v4l2sink_decoder(codec: str, fallbacks: Tuple[str, ...]) -> Tuple[str, b
70637239
f"{build_v4l2sink_output_chain(using_hw_decoder)}"
70647240
)
70657241
elif codec_type == "H264":
7066-
decoder_desc, using_hw_decoder = get_v4l2sink_decoder("H264", ("openh264dec", "avdec_h264"))
7242+
decoder_desc, using_hw_decoder = get_v4l2sink_decoder(
7243+
"H264", H264_VIEWER_DECODER_FALLBACKS
7244+
)
70677245
pipeline_desc = (
70687246
"queue ! rtph264depay ! h264parse ! "
70697247
f"{decoder_desc} ! "
@@ -7142,7 +7320,7 @@ def get_v4l2sink_decoder(codec: str, fallbacks: Tuple[str, ...]) -> Tuple[str, b
71427320

71437321
remote_label = f"remote_{pad.get_name()}"
71447322
try:
7145-
self._link_v4l2sink_bin(out, remote_label)
7323+
self._link_v4l2sink_bin(out, remote_label, wait_until_ready=True)
71467324
self.v4l2sink_remote_map[pad.get_name()] = remote_label
71477325
except Exception as exc:
71487326
printwarn(f"Failed to attach V4L2 sink video bin: {exc}")
@@ -9716,17 +9894,18 @@ def _classify_loss(loss_percent: Optional[float]) -> Tuple[str, str]:
97169894
self.pipe.add(client['webrtc'])
97179895
if self.view:
97189896
self._install_viewer_rtpbin_overrides(client['webrtc'])
9719-
try:
9720-
self._ensure_display_chain()
9721-
self._set_display_mode("idle")
9722-
except Exception as exc:
9723-
printwarn(f"Display initialization failed: {exc}")
97249897
if self.v4l2sink_device:
97259898
try:
97269899
self._ensure_v4l2sink_chain()
97279900
self._set_v4l2sink_mode("idle")
97289901
except Exception as exc:
97299902
printwarn(f"V4L2 sink initialization failed: {exc}")
9903+
else:
9904+
try:
9905+
self._ensure_display_chain()
9906+
self._set_display_mode("idle")
9907+
except Exception as exc:
9908+
printwarn(f"Display initialization failed: {exc}")
97309909

97319910
if self.vp8 or self.vp9 or self.av1 or self.h264:
97329911
direction = GstWebRTC.WebRTCRTPTransceiverDirection.RECVONLY
@@ -10252,13 +10431,22 @@ async def start_pipeline(self, UUID):
1025210431

1025310432
# Set pipeline to NULL after cleaning up elements
1025410433
try:
10255-
self._reset_display_chain_state()
10256-
self._reset_v4l2sink_chain_state()
10257-
self.pipe.set_state(Gst.State.NULL)
10434+
if self.v4l2sink_device:
10435+
# The primed idle output owns a v4l2loopback buffer pool.
10436+
# Wait for the parent to close that device before a peer
10437+
# pipeline opens it again, or the replacement pool can
10438+
# stall while activating on its first frame.
10439+
self.pipe.set_state(Gst.State.NULL)
10440+
self.pipe.get_state(2 * Gst.SECOND)
10441+
self._reset_display_chain_state()
10442+
self._reset_v4l2sink_chain_state()
10443+
else:
10444+
self._reset_display_chain_state()
10445+
self._reset_v4l2sink_chain_state()
10446+
self.pipe.set_state(Gst.State.NULL)
1025810447
except Exception as e:
1025910448
printwarn(f"Failed to set pipeline to NULL: {e}")
1026010449
self.pipe = None
10261-
1026210450
await self.createPeer(UUID)
1026310451

1026410452
def stop_pipeline(self, UUID, wait=False, expected_client=None):

0 commit comments

Comments
 (0)