@@ -5,7 +5,7 @@ use async_shutdown::ShutdownManager;
55use serde:: { Deserialize , Serialize } ;
66use tokio:: sync:: mpsc;
77use tokio:: sync:: watch;
8- use tokio_enet:: { Event , Host , HostConfig , Packet , PacketMode , PeerState } ;
8+ use tokio_enet:: { Event , Host , HostConfig , Packet , PacketMode , PeerId , PeerState } ;
99
1010use self :: input:: gamepad:: GamepadConfig ;
1111use self :: { feedback:: FeedbackCommand , input:: InputHandler } ;
@@ -391,32 +391,39 @@ fn build_termination_payload(error_code: u32) -> Vec<u8> {
391391 buf
392392}
393393
394- /// Send an encrypted control packet to the first peer of `host` if it is connected .
394+ /// Send an encrypted control packet to the connected peer.
395395#[ allow( clippy:: too_many_arguments) ]
396396fn send_to_peer (
397397 host : & mut Host ,
398+ peer_id : PeerId ,
398399 cipher : & mut GcmCipher ,
399400 key : & [ u8 ] ,
400401 key_id : i64 ,
401402 sequence_number : u32 ,
402403 payload : & [ u8 ] ,
403404 label : & str ,
404- ) {
405- if let Ok ( packet) = encode_control ( cipher, key, key_id, sequence_number, payload) {
406- if let Some ( peer) = host. peer_mut ( tokio_enet:: PeerId ( 0 ) ) {
407- if peer. state ( ) == PeerState :: Connected {
408- let _ = peer
409- . send ( 0 , Packet :: new ( packet. as_slice ( ) , PacketMode :: ReliableSequenced ) )
410- . map_err ( |e| tracing:: warn!( "Failed to send {label} to peer: {e}" ) ) ;
411- }
412- }
405+ ) -> bool {
406+ let Some ( peer) = host. peer_mut ( peer_id) else {
407+ return false ;
408+ } ;
409+ if peer. state ( ) != PeerState :: Connected {
410+ return false ;
413411 }
412+
413+ let Ok ( packet) = encode_control ( cipher, key, key_id, sequence_number, payload) else {
414+ return false ;
415+ } ;
416+
417+ peer. send ( 0 , Packet :: new ( packet. as_slice ( ) , PacketMode :: ReliableSequenced ) )
418+ . map_err ( |e| tracing:: warn!( "Failed to send {label} to peer: {e}" ) )
419+ . is_ok ( )
414420}
415421
416422/// Build and send an HDR mode control message, then advance `sequence_number`.
417423#[ allow( clippy:: too_many_arguments) ]
418424fn send_hdr_state (
419425 host : & mut Host ,
426+ peer_id : PeerId ,
420427 cipher : & mut GcmCipher ,
421428 state : & HdrModeState ,
422429 key : & [ u8 ] ,
@@ -430,9 +437,10 @@ fn send_hdr_state(
430437 None
431438 } ;
432439 let payload = build_hdr_mode_payload ( state. enabled , metadata. as_ref ( ) ) ;
433- send_to_peer ( host, cipher, key, key_id, * sequence_number, & payload, label) ;
434- * sequence_number += 1 ;
435- tracing:: info!( "Sent HDR mode ({label}): enabled={}" , state. enabled) ;
440+ if send_to_peer ( host, peer_id, cipher, key, key_id, * sequence_number, & payload, label) {
441+ * sequence_number += 1 ;
442+ tracing:: info!( "Sent HDR mode ({label}): enabled={}" , state. enabled) ;
443+ }
436444}
437445
438446#[ allow( clippy:: too_many_arguments) ]
@@ -461,6 +469,7 @@ async fn run_control_loop(
461469 let mut sequence_number = 0u32 ;
462470 let mut send_hdr_mode = false ;
463471 let mut audio_triggered = false ;
472+ let mut connected_peer: Option < PeerId > = None ;
464473
465474 // Cached AES-GCM cipher, shared by the input (decrypt) and feedback (encrypt)
466475 // directions, rebuilt only when the input key rotates.
@@ -479,23 +488,41 @@ async fn run_control_loop(
479488 // Drain all pending feedback messages so a burst doesn't back up the
480489 // channel (and stall the inputtino IO thread that produces them).
481490 while let Ok ( command) = feedback_rx. try_recv ( ) {
482- tracing:: debug!( "Sending control feedback command: {command:?}" ) ;
483- let payload = command. as_packet ( ) ;
484- let ( key, key_id) = {
485- let keys = context. keys_rx . borrow ( ) ;
486- ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
487- } ;
488- send_to_peer ( & mut host, & mut cipher, & key, key_id, sequence_number, & payload, "feedback" ) ;
489- sequence_number += 1 ;
491+ if let Some ( peer_id) = connected_peer {
492+ tracing:: debug!( "Sending control feedback command: {command:?}" ) ;
493+ let payload = command. as_packet ( ) ;
494+ let ( key, key_id) = {
495+ let keys = context. keys_rx . borrow ( ) ;
496+ ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
497+ } ;
498+ if send_to_peer (
499+ & mut host,
500+ peer_id,
501+ & mut cipher,
502+ & key,
503+ key_id,
504+ sequence_number,
505+ & payload,
506+ "feedback" ,
507+ ) {
508+ sequence_number += 1 ;
509+ }
510+ }
490511 }
491512
492513 match host
493514 . service ( Duration :: from_millis ( 10 ) )
494515 . await
495516 . map_err ( |e| tracing:: error!( "Failure in enet host: {e}" ) )
496517 {
497- Ok ( Some ( Event :: Connect { .. } ) ) => { } ,
498- Ok ( Some ( Event :: Disconnect { .. } ) ) => { } ,
518+ Ok ( Some ( Event :: Connect { peer_id, .. } ) ) => {
519+ connected_peer = Some ( peer_id) ;
520+ } ,
521+ Ok ( Some ( Event :: Disconnect { peer_id, .. } ) ) => {
522+ if connected_peer == Some ( peer_id) {
523+ connected_peer = None ;
524+ }
525+ } ,
499526 Ok ( Some ( Event :: Receive { ref packet, .. } ) ) => {
500527 let mut control_message = match ControlMessage :: from_bytes ( packet. data ( ) ) {
501528 Ok ( control_message) => control_message,
@@ -576,22 +603,44 @@ async fn run_control_loop(
576603 // Send HDR mode notification after the host.service() match to avoid double mutable borrow.
577604 if send_hdr_mode {
578605 send_hdr_mode = false ;
579- let state = hdr_metadata_rx. borrow_and_update ( ) . clone ( ) ;
580- let ( key, key_id) = {
581- let keys = context. keys_rx . borrow ( ) ;
582- ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
583- } ;
584- send_hdr_state ( & mut host, & mut cipher, & state, & key, key_id, & mut sequence_number, "initial" ) ;
606+ if let Some ( peer_id) = connected_peer {
607+ let state = hdr_metadata_rx. borrow_and_update ( ) . clone ( ) ;
608+ let ( key, key_id) = {
609+ let keys = context. keys_rx . borrow ( ) ;
610+ ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
611+ } ;
612+ send_hdr_state (
613+ & mut host,
614+ peer_id,
615+ & mut cipher,
616+ & state,
617+ & key,
618+ key_id,
619+ & mut sequence_number,
620+ "initial" ,
621+ ) ;
622+ }
585623 }
586624
587625 // Check for HDR metadata updates from the video pipeline.
588626 if context. hdr && hdr_metadata_rx. has_changed ( ) . unwrap_or ( false ) {
589- let state = hdr_metadata_rx. borrow_and_update ( ) . clone ( ) ;
590- let ( key, key_id) = {
591- let keys = context. keys_rx . borrow ( ) ;
592- ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
593- } ;
594- send_hdr_state ( & mut host, & mut cipher, & state, & key, key_id, & mut sequence_number, "metadata update" ) ;
627+ if let Some ( peer_id) = connected_peer {
628+ let state = hdr_metadata_rx. borrow_and_update ( ) . clone ( ) ;
629+ let ( key, key_id) = {
630+ let keys = context. keys_rx . borrow ( ) ;
631+ ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
632+ } ;
633+ send_hdr_state (
634+ & mut host,
635+ peer_id,
636+ & mut cipher,
637+ & state,
638+ & key,
639+ key_id,
640+ & mut sequence_number,
641+ "metadata update" ,
642+ ) ;
643+ }
595644 }
596645 }
597646
@@ -605,7 +654,18 @@ async fn run_control_loop(
605654 let keys = context. keys_rx . borrow ( ) ;
606655 ( keys. remote_input_key . clone ( ) , keys. remote_input_key_id )
607656 } ;
608- send_to_peer ( & mut host, & mut cipher, & key, key_id, sequence_number, & termination_payload, "termination" ) ;
657+ if let Some ( peer_id) = connected_peer {
658+ send_to_peer (
659+ & mut host,
660+ peer_id,
661+ & mut cipher,
662+ & key,
663+ key_id,
664+ sequence_number,
665+ & termination_payload,
666+ "termination" ,
667+ ) ;
668+ }
609669 let _ = host. flush ( ) . await ;
610670
611671 // Explicitly drop the ENet host before the delay shutdown token
0 commit comments