@@ -180,6 +180,7 @@ pub async fn create_answer(
180180 first_frame,
181181 peer_connection,
182182 video_track,
183+ media_codec,
183184 cancellation_token,
184185 cancellation,
185186 }
@@ -360,12 +361,6 @@ impl WebRtcMediaCodec {
360361 }
361362 }
362363
363- fn from_frame ( frame : & crate :: transport:: packet:: FramePacket ) -> anyhow:: Result < Self > {
364- let codec = frame. codec . as_deref ( ) . unwrap_or_default ( ) ;
365- Self :: from_codec_string ( codec)
366- . ok_or_else ( || anyhow:: anyhow!( "unsupported WebRTC video codec `{codec}`" ) )
367- }
368-
369364 fn mime_type ( self ) -> & ' static str {
370365 match self {
371366 Self :: H264 => MIME_TYPE_H264 ,
@@ -439,6 +434,7 @@ struct WebRtcMediaStream {
439434 first_frame : crate :: transport:: packet:: SharedFrame ,
440435 peer_connection : Arc < webrtc:: peer_connection:: RTCPeerConnection > ,
441436 video_track : Arc < TrackLocalStaticSample > ,
437+ media_codec : WebRtcMediaCodec ,
442438 cancellation_token : broadcast:: Sender < ( ) > ,
443439 cancellation : broadcast:: Receiver < ( ) > ,
444440}
@@ -452,12 +448,12 @@ impl WebRtcMediaStream {
452448 first_frame,
453449 peer_connection,
454450 video_track,
451+ media_codec,
455452 cancellation_token,
456453 mut cancellation,
457454 } = self ;
458455 let mut rx = session. subscribe ( ) ;
459456 let mut latest_keyframe = first_frame. clone ( ) ;
460- let mut last_sequence = 0u64 ;
461457 let mut send_timing = WebRtcSendTiming :: new ( ) ;
462458 let mut peer_state_interval = time:: interval ( Duration :: from_millis ( 250 ) ) ;
463459 let mut bootstrap_sleep = Box :: pin ( time:: sleep ( WEBRTC_BOOTSTRAP_KEYFRAME_INTERVAL ) ) ;
@@ -467,44 +463,70 @@ impl WebRtcMediaStream {
467463 let mut waiting_for_keyframe = false ;
468464 let _guard = WebRtcMetricsGuard :: new ( state. metrics . clone ( ) ) ;
469465
470- if let Err ( error) =
471- write_frame_sample ( & video_track, & first_frame, WEBRTC_MIN_REFRESH_INTERVAL ) . await
466+ match write_frame_sample_with_timeout (
467+ & video_track,
468+ & first_frame,
469+ media_codec,
470+ WEBRTC_MIN_REFRESH_INTERVAL ,
471+ )
472+ . await
472473 {
473- warn ! ( "WebRTC initial keyframe write failed for {udid}: {error}" ) ;
474- let _ = peer_connection. close ( ) . await ;
475- return ;
474+ Ok ( true ) => {
475+ state. metrics . frames_sent . fetch_add ( 1 , Ordering :: Relaxed ) ;
476+ }
477+ Ok ( false ) => {
478+ state
479+ . metrics
480+ . frames_dropped_server
481+ . fetch_add ( 1 , Ordering :: Relaxed ) ;
482+ session. request_keyframe ( ) ;
483+ }
484+ Err ( error) => {
485+ warn ! ( "WebRTC initial keyframe write failed for {udid}: {error}" ) ;
486+ let _ = peer_connection. close ( ) . await ;
487+ return ;
488+ }
476489 }
477- state. metrics . frames_sent . fetch_add ( 1 , Ordering :: Relaxed ) ;
478490
479491 loop {
480492 tokio:: select! {
481493 _ = cancellation. recv( ) => {
494+ warn!( "WebRTC media stream replaced for {udid}" ) ;
482495 break ;
483496 }
484497 _ = peer_state_interval. tick( ) => {
485- if matches!(
486- peer_connection. connection_state( ) ,
487- RTCPeerConnectionState :: Closed
488- | RTCPeerConnectionState :: Disconnected
489- | RTCPeerConnectionState :: Failed
490- ) {
498+ let peer_state = peer_connection. connection_state( ) ;
499+ if matches!( peer_state, RTCPeerConnectionState :: Closed | RTCPeerConnectionState :: Failed ) {
500+ warn!( "WebRTC media stream closing for {udid}: peer state {peer_state}" ) ;
491501 break ;
492502 }
493503 }
494504 _ = & mut bootstrap_sleep, if bootstrap_frames_remaining > 0 => {
495- if let Err ( error ) = write_frame_sample (
505+ match write_frame_sample_with_timeout (
496506 & video_track,
497507 & latest_keyframe,
508+ media_codec,
498509 WEBRTC_BOOTSTRAP_KEYFRAME_INTERVAL ,
499510 ) . await {
500- warn!( "WebRTC bootstrap keyframe write failed for {udid}: {error}" ) ;
501- break ;
511+ Ok ( true ) => {
512+ state. metrics. frames_sent. fetch_add( 1 , Ordering :: Relaxed ) ;
513+ }
514+ Ok ( false ) => {
515+ state
516+ . metrics
517+ . frames_dropped_server
518+ . fetch_add( 1 , Ordering :: Relaxed ) ;
519+ session. request_keyframe( ) ;
520+ }
521+ Err ( error) => {
522+ warn!( "WebRTC bootstrap keyframe write failed for {udid}: {error}" ) ;
523+ break ;
524+ }
502525 }
503526 bootstrap_frames_remaining = bootstrap_frames_remaining. saturating_sub( 1 ) ;
504527 bootstrap_sleep
505528 . as_mut( )
506529 . reset( time:: Instant :: now( ) + WEBRTC_BOOTSTRAP_KEYFRAME_INTERVAL ) ;
507- state. metrics. frames_sent. fetch_add( 1 , Ordering :: Relaxed ) ;
508530 }
509531 _ = & mut refresh_sleep => {
510532 session. request_refresh( ) ;
@@ -524,17 +546,11 @@ impl WebRtcMediaStream {
524546 session. request_keyframe( ) ;
525547 continue ;
526548 }
527- Err ( broadcast:: error:: RecvError :: Closed ) => break ,
549+ Err ( broadcast:: error:: RecvError :: Closed ) => {
550+ warn!( "WebRTC media stream closing for {udid}: frame channel closed" ) ;
551+ break ;
552+ }
528553 } ;
529- if last_sequence != 0 && frame. frame_sequence > last_sequence + 1 && !frame. is_keyframe {
530- state
531- . metrics
532- . frames_dropped_server
533- . fetch_add( frame. frame_sequence - last_sequence - 1 , Ordering :: Relaxed ) ;
534- waiting_for_keyframe = true ;
535- session. request_keyframe( ) ;
536- continue ;
537- }
538554 if waiting_for_keyframe && !frame. is_keyframe {
539555 state. metrics. frames_dropped_server. fetch_add( 1 , Ordering :: Relaxed ) ;
540556 continue ;
@@ -545,27 +561,31 @@ impl WebRtcMediaStream {
545561 }
546562 let duration = send_timing. duration_for( & frame) ;
547563 let started_at = time:: Instant :: now( ) ;
548- let write_result = time:: timeout(
549- WEBRTC_WRITE_TIMEOUT ,
550- write_frame_sample( & video_track, & frame, duration) ,
551- ) . await ;
564+ let write_result = write_frame_sample_with_timeout( & video_track, & frame, media_codec, duration) . await ;
552565 adaptive_refresh_interval = adaptive_interval_for_write( started_at. elapsed( ) ) ;
553- if let Err ( error) = write_result
554- . map_err( |_| anyhow:: anyhow!(
555- "timed out writing WebRTC frame after {}ms" ,
556- WEBRTC_WRITE_TIMEOUT . as_millis( )
557- ) )
558- . and_then( |result| result)
559- {
560- warn!( "WebRTC frame write failed for {udid}: {error}" ) ;
561- break ;
566+ match write_result {
567+ Ok ( true ) => {
568+ state. metrics. frames_sent. fetch_add( 1 , Ordering :: Relaxed ) ;
569+ }
570+ Ok ( false ) => {
571+ state
572+ . metrics
573+ . frames_dropped_server
574+ . fetch_add( 1 , Ordering :: Relaxed ) ;
575+ waiting_for_keyframe = true ;
576+ adaptive_refresh_interval = WEBRTC_MAX_REFRESH_INTERVAL ;
577+ session. request_keyframe( ) ;
578+ }
579+ Err ( error) => {
580+ warn!( "WebRTC frame write failed for {udid}: {error}" ) ;
581+ break ;
582+ }
562583 }
563- last_sequence = frame. frame_sequence;
564- state. metrics. frames_sent. fetch_add( 1 , Ordering :: Relaxed ) ;
565584 }
566585 }
567586 }
568587
588+ warn ! ( "WebRTC media stream ended for {udid}" ) ;
569589 clear_webrtc_media_stream ( & udid, & cancellation_token) ;
570590 let _ = peer_connection. close ( ) . await ;
571591 }
@@ -582,9 +602,10 @@ fn adaptive_interval_for_write(write_elapsed: Duration) -> Duration {
582602async fn write_frame_sample (
583603 video_track : & TrackLocalStaticSample ,
584604 frame : & crate :: transport:: packet:: SharedFrame ,
605+ media_codec : WebRtcMediaCodec ,
585606 duration : Duration ,
586607) -> anyhow:: Result < ( ) > {
587- let data = annex_b_sample ( frame) ?;
608+ let data = annex_b_sample ( frame, media_codec ) ?;
588609 video_track
589610 . write_sample ( & Sample {
590611 data : Bytes :: from ( data) ,
@@ -595,8 +616,28 @@ async fn write_frame_sample(
595616 Ok ( ( ) )
596617}
597618
598- fn annex_b_sample ( frame : & crate :: transport:: packet:: FramePacket ) -> anyhow:: Result < Vec < u8 > > {
599- match WebRtcMediaCodec :: from_frame ( frame) ? {
619+ async fn write_frame_sample_with_timeout (
620+ video_track : & TrackLocalStaticSample ,
621+ frame : & crate :: transport:: packet:: SharedFrame ,
622+ media_codec : WebRtcMediaCodec ,
623+ duration : Duration ,
624+ ) -> anyhow:: Result < bool > {
625+ match time:: timeout (
626+ WEBRTC_WRITE_TIMEOUT ,
627+ write_frame_sample ( video_track, frame, media_codec, duration) ,
628+ )
629+ . await
630+ {
631+ Ok ( result) => result. map ( |( ) | true ) ,
632+ Err ( _) => Ok ( false ) ,
633+ }
634+ }
635+
636+ fn annex_b_sample (
637+ frame : & crate :: transport:: packet:: FramePacket ,
638+ media_codec : WebRtcMediaCodec ,
639+ ) -> anyhow:: Result < Vec < u8 > > {
640+ match media_codec {
600641 WebRtcMediaCodec :: H264 => h264_annex_b_sample ( frame) ,
601642 WebRtcMediaCodec :: Hevc => hevc_annex_b_sample ( frame) ,
602643 }
0 commit comments