fix: 完善 RK3588 HDMI RX 信号检测与自动恢复

- 统一使用 QUERY_DV_TIMINGS 判断 HDMI RX 输入状态
  - 增加 source-following 设备状态 API 与前端只读展示
  - 支持长时间无信号后的 MJPEG/WebRTC 自动恢复
  - 修复模式切换时采集设备尚未释放导致的 EBUSY
  - 取消陈旧 WebRTC 重连并统一采集恢复策略
  - 移除视频输入状态区域的冗余标题
This commit is contained in:
mofeng-git
2026-07-30 19:41:44 +08:00
parent 971c263bf8
commit 9fb23476ac
36 changed files with 1883 additions and 1492 deletions

View File

@@ -15,7 +15,8 @@ use crate::hid::HidController;
use crate::video::capture::DEFAULT_CAPTURE_BUFFER_COUNT;
use crate::video::codec::h264_bitstream;
use crate::video::device::{
enumerate_devices, select_recovery_device, VideoDevice, VideoDeviceRecoveryHint,
enumerate_devices, select_recovery_device, VideoControlMode, VideoDevice, VideoDeviceInfo,
VideoDeviceRecoveryHint,
};
use crate::video::types::{
BitratePreset, EncoderBackend, PipelineStateNotification, PixelFormat, Resolution,
@@ -28,6 +29,18 @@ use super::signaling::{ConnectionState, IceCandidate, SdpAnswer, SdpOffer};
use super::universal_session::{UniversalSession, UniversalSessionConfig};
const H264_PROFILE_DETECT_TIMEOUT: Duration = Duration::from_millis(500);
const PIPELINE_RELEASE_TIMEOUT: Duration = Duration::from_secs(5);
fn update_signal_recovery_edge(pending: &AtomicBool, state: &str) -> bool {
match state {
"no_signal" => {
pending.store(true, Ordering::Release);
false
}
"streaming" => pending.swap(false, Ordering::AcqRel),
_ => false,
}
}
#[derive(Debug, Clone)]
pub struct WebRtcStreamerConfig {
@@ -63,8 +76,7 @@ pub struct CaptureDeviceConfig {
pub jpeg_quality: u8,
pub subdev_path: Option<PathBuf>,
pub bridge_kind: Option<String>,
/// V4L2 driver name (e.g. `uvcvideo`) for UVC-specific recovery hints.
pub v4l2_driver: Option<String>,
pub control_mode: VideoControlMode,
pub recovery_hint: VideoDeviceRecoveryHint,
}
@@ -99,6 +111,7 @@ pub struct WebRtcStreamer {
hid_controller: RwLock<Option<Arc<HidController>>>,
events: RwLock<Option<Arc<EventBus>>>,
recovery_in_progress: AtomicBool,
signal_recovery_pending: Arc<AtomicBool>,
self_weak: StdRwLock<Option<std::sync::Weak<Self>>>,
}
@@ -119,6 +132,7 @@ impl WebRtcStreamer {
hid_controller: RwLock::new(None),
events: RwLock::new(None),
recovery_in_progress: AtomicBool::new(false),
signal_recovery_pending: Arc::new(AtomicBool::new(false)),
self_weak: StdRwLock::new(None),
});
let weak = Arc::downgrade(&streamer);
@@ -154,11 +168,7 @@ impl WebRtcStreamer {
// Close all existing sessions
self.close_all_sessions().await;
// Stop current pipeline
if let Some(ref pipeline) = *self.video_pipeline.read().await {
pipeline.stop();
}
*self.video_pipeline.write().await = None;
self.stop_video_pipeline_and_release().await?;
// Update codec
*self.video_codec.write().await = codec;
@@ -231,18 +241,48 @@ impl WebRtcStreamer {
}
}
/// Serialize pipeline teardown with creation and return only after V4L2
/// STREAMOFF, buffer teardown and FD close have completed.
async fn stop_video_pipeline_and_release(&self) -> Result<()> {
let mut pipeline_guard = self.video_pipeline.write().await;
let Some(pipeline) = pipeline_guard.as_ref().cloned() else {
return Ok(());
};
pipeline.stop_and_wait(PIPELINE_RELEASE_TIMEOUT).await?;
*pipeline_guard = None;
Ok(())
}
fn build_pipeline_state_notifier(
device: String,
events: Option<Arc<EventBus>>,
recovery_pending: Arc<AtomicBool>,
) -> Option<Arc<dyn Fn(PipelineStateNotification) + Send + Sync>> {
events.map(|events| {
Arc::new(move |notification: PipelineStateNotification| {
let recovered = update_signal_recovery_edge(&recovery_pending, notification.state);
events.publish(SystemEvent::StreamStateChanged {
state: notification.state.to_string(),
device: Some(device.clone()),
reason: notification.reason.map(|reason| reason.to_string()),
next_retry_ms: notification.next_retry_ms,
});
if recovered {
events.publish(SystemEvent::StreamRecovered {
device: device.clone(),
});
if let Some(applied) = notification.applied_config {
events.publish(SystemEvent::StreamConfigApplied {
transition_id: None,
device: device.clone(),
resolution: (applied.resolution.width, applied.resolution.height),
format: applied.format.to_string(),
fps: applied.fps,
});
}
events.mark_device_info_dirty();
}
}) as Arc<dyn Fn(PipelineStateNotification) + Send + Sync>
})
}
@@ -330,7 +370,7 @@ impl WebRtcStreamer {
jpeg_quality,
subdev_path: device.subdev_path.clone(),
bridge_kind: device.bridge_kind.clone(),
v4l2_driver: Some(device.driver.clone()),
control_mode: device.control_mode,
recovery_hint: VideoDeviceRecoveryHint::from(&device),
};
@@ -347,6 +387,7 @@ impl WebRtcStreamer {
debug!("WebRTC video recovery already in progress");
return;
}
self.signal_recovery_pending.store(true, Ordering::Release);
let streamer = self.clone();
tokio::spawn(async move {
@@ -412,24 +453,11 @@ impl WebRtcStreamer {
{
Ok(reconnected) => {
info!(
"WebRTC video recovered with {} after {} attempts, reconnected {} sessions",
"WebRTC capture reopened with {} after {} attempts; reconnected {} sessions and waiting for first frame",
device.path.display(),
attempt,
reconnected
);
streamer
.publish_stream_event(SystemEvent::StreamRecovered {
device: device.path.display().to_string(),
})
.await;
streamer
.publish_stream_event(SystemEvent::StreamStateChanged {
state: "streaming".to_string(),
device: Some(device.path.display().to_string()),
reason: None,
next_retry_ms: None,
})
.await;
streamer.recovery_in_progress.store(false, Ordering::SeqCst);
return;
}
@@ -460,6 +488,13 @@ impl WebRtcStreamer {
let pipeline_config = {
let config = self.config.read().await;
SharedVideoPipelineConfig {
control_mode: self
.capture_device
.read()
.await
.as_ref()
.map(|capture| capture.control_mode)
.unwrap_or(VideoControlMode::Configurable),
resolution: config.resolution,
input_format: config.input_format,
output_codec: Self::codec_type_to_encoder_type(codec),
@@ -476,6 +511,7 @@ impl WebRtcStreamer {
pipeline.set_state_notifier(Self::build_pipeline_state_notifier(
device.device_path.display().to_string(),
self.events.read().await.clone(),
self.signal_recovery_pending.clone(),
));
pipeline
.start_with_device(
@@ -484,7 +520,6 @@ impl WebRtcStreamer {
device.jpeg_quality,
device.subdev_path,
device.bridge_kind,
device.v4l2_driver,
)
.await?;
} else {
@@ -532,7 +567,8 @@ impl WebRtcStreamer {
let should_reconnect = pending_geometry.is_some();
if let Some((r, f)) = pending_geometry {
streamer.sync_video_geometry_from_negotiated(r, f).await;
let fps = streamer.config.read().await.fps;
streamer.sync_video_input_from_negotiated(r, f, fps).await;
}
if should_reconnect {
let streamer_for_reconnect = streamer.clone();
@@ -570,9 +606,10 @@ impl WebRtcStreamer {
});
let pipeline_cfg = pipeline.config().await;
self.sync_video_geometry_from_negotiated(
self.sync_video_input_from_negotiated(
pipeline_cfg.resolution,
pipeline_cfg.input_format,
pipeline_cfg.fps,
)
.await;
@@ -725,32 +762,41 @@ impl WebRtcStreamer {
&self,
device_path: PathBuf,
jpeg_quality: u8,
subdev_path: Option<PathBuf>,
bridge_kind: Option<String>,
v4l2_driver: Option<String>,
device_info: Option<VideoDeviceInfo>,
) {
let (subdev_path, bridge_kind, control_mode, recovery_hint) = match device_info {
Some(info) => (
info.subdev_path.clone(),
info.bridge_kind.clone(),
info.control_mode,
VideoDeviceRecoveryHint::from(&info),
),
None => {
let recovery_hint = VideoDevice::open_readonly(&device_path)
.and_then(|device| device.info())
.map(|info| VideoDeviceRecoveryHint::from(&info))
.unwrap_or_else(|_| VideoDeviceRecoveryHint {
path: device_path.clone(),
name: String::new(),
driver: String::new(),
bus_info: String::new(),
card: String::new(),
is_capture_card: true,
});
(None, None, VideoControlMode::Configurable, recovery_hint)
}
};
debug!(
"Setting direct capture device for WebRTC: {:?} (subdev={:?}, kind={:?}, driver={:?})",
device_path, subdev_path, bridge_kind, v4l2_driver
"Setting direct capture device for WebRTC: {:?} (subdev={:?}, kind={:?}, mode={:?})",
device_path, subdev_path, bridge_kind, control_mode
);
let recovery_hint = VideoDevice::open_readonly(&device_path)
.and_then(|device| device.info())
.map(|info| VideoDeviceRecoveryHint::from(&info))
.unwrap_or_else(|_| VideoDeviceRecoveryHint {
path: device_path.clone(),
name: String::new(),
driver: v4l2_driver.clone().unwrap_or_default(),
bus_info: String::new(),
card: String::new(),
is_capture_card: true,
});
*self.capture_device.write().await = Some(CaptureDeviceConfig {
device_path,
buffer_count: DEFAULT_CAPTURE_BUFFER_COUNT,
jpeg_quality,
subdev_path,
bridge_kind,
v4l2_driver,
control_mode,
recovery_hint,
});
}
@@ -764,12 +810,13 @@ impl WebRtcStreamer {
///
/// This stops the encoding pipeline and closes all sessions.
pub async fn prepare_for_config_change(&self) {
// Stop pipeline and close sessions - will be recreated on next session
if let Some(ref pipeline) = *self.video_pipeline.read().await {
pipeline.stop();
}
*self.video_pipeline.write().await = None;
self.close_all_sessions().await;
if let Err(error) = self.stop_video_pipeline_and_release().await {
warn!(
"Failed to release video pipeline for config change: {}",
error
);
}
}
// === Configuration ===
@@ -804,12 +851,6 @@ impl WebRtcStreamer {
resolution.width, resolution.height, format, fps
);
// Stop existing pipeline
if let Some(ref pipeline) = *self.video_pipeline.read().await {
pipeline.stop();
}
*self.video_pipeline.write().await = None;
// Close all existing sessions - they need to reconnect
let session_count = self.close_all_sessions().await;
if session_count > 0 {
@@ -818,6 +859,13 @@ impl WebRtcStreamer {
session_count
);
}
if let Err(error) = self.stop_video_pipeline_and_release().await {
warn!(
"Failed to release video pipeline for config change: {}",
error
);
return;
}
// Update config (preserve user-configured bitrate)
{
@@ -836,31 +884,36 @@ impl WebRtcStreamer {
self.notify_device_info_dirty().await;
}
/// Update resolution/format to match DV-negotiated capture without stopping
/// Update the input mode to match DV-negotiated capture without stopping
/// the pipeline or closing sessions. Used when hardware timing differs from
/// saved settings (e.g. RK628 `S_FMT` follows source while SQLite still has
/// a user-chosen preset).
pub async fn sync_video_geometry_from_negotiated(
pub async fn sync_video_input_from_negotiated(
&self,
resolution: Resolution,
format: PixelFormat,
fps: u32,
) {
{
let mut config = self.config.write().await;
if config.resolution == resolution && config.input_format == format {
if config.resolution == resolution && config.input_format == format && config.fps == fps
{
return;
}
info!(
"WebRTC geometry aligned to negotiated capture: {}x{} {:?} (was {}x{} {:?})",
"WebRTC input aligned to negotiated capture: {}x{} {:?} @ {} fps (was {}x{} {:?} @ {} fps)",
resolution.width,
resolution.height,
format,
fps,
config.resolution.width,
config.resolution.height,
config.input_format
config.input_format,
config.fps,
);
config.resolution = resolution;
config.input_format = format;
config.fps = fps;
}
self.notify_device_info_dirty().await;
@@ -868,12 +921,6 @@ impl WebRtcStreamer {
/// Update encoder backend (software/hardware selection)
pub async fn update_encoder_backend(&self, encoder_backend: Option<EncoderBackend>) {
// Stop existing pipeline
if let Some(ref pipeline) = *self.video_pipeline.read().await {
pipeline.stop();
}
*self.video_pipeline.write().await = None;
// Close all existing sessions - they need to reconnect with new encoder
let session_count = self.close_all_sessions().await;
if session_count > 0 {
@@ -882,6 +929,13 @@ impl WebRtcStreamer {
session_count
);
}
if let Err(error) = self.stop_video_pipeline_and_release().await {
warn!(
"Failed to release video pipeline for encoder backend change: {}",
error
);
return;
}
// Update config
let mut config = self.config.write().await;
@@ -1123,17 +1177,12 @@ impl WebRtcStreamer {
/// Close all sessions and wait for the video pipeline to fully release the
/// capture device. Use this when the caller needs the V4L2 device immediately
/// afterwards (e.g. switching to MJPEG mode).
pub async fn close_all_sessions_and_release_device(&self) -> usize {
pub async fn close_all_sessions_and_release_device(&self) -> Result<usize> {
let count = self.close_all_sessions().await;
if let Some(ref pipeline) = *self.video_pipeline.read().await {
pipeline
.stop_and_wait(std::time::Duration::from_secs(3))
.await;
}
*self.video_pipeline.write().await = None;
self.stop_video_pipeline_and_release().await?;
count
Ok(count)
}
/// Get session count
@@ -1256,16 +1305,7 @@ impl WebRtcStreamer {
if pipeline_running {
info!("Restarting video pipeline to apply new bitrate: {}", preset);
// Stop existing pipeline
if let Some(ref pipeline) = *self.video_pipeline.read().await {
pipeline.stop();
}
// Wait for pipeline to stop
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
// Clear pipeline reference - will be recreated
*self.video_pipeline.write().await = None;
self.stop_video_pipeline_and_release().await?;
let has_source = self.capture_device.read().await.is_some();
if !has_source {
@@ -1306,18 +1346,10 @@ impl crate::video::traits::VideoOutput for WebRtcStreamer {
&self,
device_path: PathBuf,
jpeg_quality: u8,
subdev_path: Option<PathBuf>,
bridge_kind: Option<String>,
v4l2_driver: Option<String>,
device_info: Option<VideoDeviceInfo>,
) {
self.set_capture_device(
device_path,
jpeg_quality,
subdev_path,
bridge_kind,
v4l2_driver,
)
.await;
self.set_capture_device(device_path, jpeg_quality, device_info)
.await;
}
async fn current_video_codec(&self) -> VideoCodecType {
@@ -1332,7 +1364,7 @@ impl crate::video::traits::VideoOutput for WebRtcStreamer {
self.close_all_sessions().await;
}
async fn close_all_sessions_and_release_device(&self) -> usize {
async fn close_all_sessions_and_release_device(&self) -> Result<usize> {
self.close_all_sessions_and_release_device().await
}
@@ -1398,6 +1430,7 @@ impl Default for WebRtcStreamer {
hid_controller: RwLock::new(None),
events: RwLock::new(None),
recovery_in_progress: AtomicBool::new(false),
signal_recovery_pending: Arc::new(AtomicBool::new(false)),
self_weak: StdRwLock::new(None),
}
}
@@ -1431,4 +1464,14 @@ mod tests {
assert!(!WebRtcStreamer::should_stop_pipeline(0, 1));
assert!(!WebRtcStreamer::should_stop_pipeline(2, 3));
}
#[test]
fn recovery_edge_is_emitted_once_after_first_streaming_frame() {
let pending = AtomicBool::new(false);
assert!(!update_signal_recovery_edge(&pending, "streaming"));
assert!(!update_signal_recovery_edge(&pending, "no_signal"));
assert!(!update_signal_recovery_edge(&pending, "no_signal"));
assert!(update_signal_recovery_edge(&pending, "streaming"));
assert!(!update_signal_recovery_edge(&pending, "streaming"));
}
}