@@ -374,7 +374,9 @@ impl SessionInner {
374374 track_name : String ,
375375 ) -> ( Option < u64 > , bool ) {
376376 let mut guard = self . track_lifecycle . lock ( ) . unwrap ( ) ;
377- let entry = guard. entry ( key. to_string ( ) ) . or_insert_with ( TrackLifecycle :: pending) ;
377+ let entry = guard
378+ . entry ( key. to_string ( ) )
379+ . or_insert_with ( TrackLifecycle :: pending) ;
378380 entry. created = true ;
379381 entry. closed = false ;
380382
@@ -401,7 +403,9 @@ impl SessionInner {
401403
402404 fn activate_track ( & self , key : & str , track_alias : u64 ) {
403405 let mut guard = self . track_lifecycle . lock ( ) . unwrap ( ) ;
404- let entry = guard. entry ( key. to_string ( ) ) . or_insert_with ( TrackLifecycle :: pending) ;
406+ let entry = guard
407+ . entry ( key. to_string ( ) )
408+ . or_insert_with ( TrackLifecycle :: pending) ;
405409 if entry. closed {
406410 return ;
407411 }
@@ -435,7 +439,9 @@ impl SessionInner {
435439
436440 fn close_track ( & self , key : & str ) -> Vec < TrackNotifier > {
437441 let mut guard = self . track_lifecycle . lock ( ) . unwrap ( ) ;
438- let entry = guard. entry ( key. to_string ( ) ) . or_insert_with ( TrackLifecycle :: created) ;
442+ let entry = guard
443+ . entry ( key. to_string ( ) )
444+ . or_insert_with ( TrackLifecycle :: created) ;
439445
440446 if entry. closed {
441447 return Vec :: new ( ) ;
@@ -919,26 +925,128 @@ async fn control_loop(
919925 inner : Arc < SessionInner > ,
920926 shutdown_rx : & mut watch:: Receiver < bool > ,
921927) {
922- loop {
928+ let exit_reason = loop {
923929 tokio:: select! {
924- _ = shutdown_rx. changed( ) => break ,
930+ _ = shutdown_rx. changed( ) => break None ,
925931 msg = control_rx. recv( ) => {
926932 match msg {
927933 Some ( msg) => {
928- if handler. send( & msg) . await . is_err ( ) {
929- break ;
934+ if let Err ( err ) = handler. send( & msg) . await {
935+ break Some ( format! ( "control stream send failed: {err:?}" ) ) ;
930936 }
931937 }
932- None => break ,
938+ None => break None ,
933939 }
934940 }
935941 result = handler. next_message( ) => {
936942 match result {
937943 Ok ( msg) => dispatch_control_response( msg, & inner) ,
938- Err ( _ ) => break ,
944+ Err ( err ) => break Some ( format! ( "control stream receive failed: {err:?}" ) ) ,
939945 }
940946 }
941947 }
948+ } ;
949+
950+ if let Some ( reason) = exit_reason {
951+ fail_pending_control_ops ( & inner, & reason) ;
952+ }
953+ }
954+
955+ fn fail_pending_control_ops ( inner : & Arc < SessionInner > , reason : & str ) {
956+ let pending_publishes = {
957+ let mut guard = inner. pending_publishes . lock ( ) . unwrap ( ) ;
958+ std:: mem:: take ( & mut * guard)
959+ } ;
960+
961+ for ( _, pending) in pending_publishes {
962+ let mut msg_env = OwnedEnv :: new ( ) ;
963+ let pid = pending. caller_pid ;
964+ pending. ref_env . run ( |ref_env| {
965+ let publish_ref = pending. publish_ref_term . load ( ref_env) ;
966+ let _ = msg_env. send_and_clear ( & pid, |env| {
967+ let nil = rustler:: types:: atom:: nil ( ) . to_term ( env) ;
968+ let payload = TransportErrorOut {
969+ op : atoms:: publish ( ) ,
970+ message : reason. to_string ( ) ,
971+ kind : Some ( atoms:: runtime ( ) ) ,
972+ r#ref : publish_ref. in_env ( env) ,
973+ handle : nil,
974+ } ;
975+ ( atoms:: moqx_transport_error ( ) , payload) . encode ( env)
976+ } ) ;
977+ } ) ;
978+ }
979+
980+ let pending_subscribes = {
981+ let mut guard = inner. pending_subscribes . lock ( ) . unwrap ( ) ;
982+ std:: mem:: take ( & mut * guard)
983+ } ;
984+
985+ for ( _, pending) in pending_subscribes {
986+ let mut msg_env = OwnedEnv :: new ( ) ;
987+ let pid = pending. caller_pid ;
988+ pending. term_env . run ( |term_env| {
989+ let sub_ref = pending. subscription_ref_term . load ( term_env) ;
990+ let _ = msg_env. send_and_clear ( & pid, |env| {
991+ let nil = rustler:: types:: atom:: nil ( ) . to_term ( env) ;
992+ let payload = TransportErrorOut {
993+ op : atoms:: subscribe ( ) ,
994+ message : reason. to_string ( ) ,
995+ kind : Some ( atoms:: runtime ( ) ) ,
996+ r#ref : nil,
997+ handle : sub_ref. in_env ( env) ,
998+ } ;
999+ ( atoms:: moqx_transport_error ( ) , payload) . encode ( env)
1000+ } ) ;
1001+ } ) ;
1002+ }
1003+
1004+ let pending_fetches = {
1005+ let mut guard = inner. pending_fetches . lock ( ) . unwrap ( ) ;
1006+ std:: mem:: take ( & mut * guard)
1007+ } ;
1008+
1009+ for ( _, pending) in pending_fetches {
1010+ let mut msg_env = OwnedEnv :: new ( ) ;
1011+ let pid = pending. caller_pid ;
1012+ pending. ref_env . run ( |ref_env| {
1013+ let fetch_ref = pending. ref_term . load ( ref_env) ;
1014+ let _ = msg_env. send_and_clear ( & pid, |env| {
1015+ let nil = rustler:: types:: atom:: nil ( ) . to_term ( env) ;
1016+ let payload = TransportErrorOut {
1017+ op : atoms:: fetch ( ) ,
1018+ message : reason. to_string ( ) ,
1019+ kind : Some ( atoms:: runtime ( ) ) ,
1020+ r#ref : fetch_ref. in_env ( env) ,
1021+ handle : nil,
1022+ } ;
1023+ ( atoms:: moqx_transport_error ( ) , payload) . encode ( env)
1024+ } ) ;
1025+ } ) ;
1026+ }
1027+
1028+ let active_subscriptions = {
1029+ let mut guard = inner. active_subscriptions . lock ( ) . unwrap ( ) ;
1030+ std:: mem:: take ( & mut * guard)
1031+ } ;
1032+
1033+ for ( _, active) in active_subscriptions {
1034+ let mut msg_env = OwnedEnv :: new ( ) ;
1035+ let pid = active. caller_pid ;
1036+ active. ref_env . run ( |ref_env| {
1037+ let sub_ref = active. subscription_ref_term . load ( ref_env) ;
1038+ let _ = msg_env. send_and_clear ( & pid, |env| {
1039+ let nil = rustler:: types:: atom:: nil ( ) . to_term ( env) ;
1040+ let payload = TransportErrorOut {
1041+ op : atoms:: subscribe ( ) ,
1042+ message : reason. to_string ( ) ,
1043+ kind : Some ( atoms:: runtime ( ) ) ,
1044+ r#ref : nil,
1045+ handle : sub_ref. in_env ( env) ,
1046+ } ;
1047+ ( atoms:: moqx_transport_error ( ) , payload) . encode ( env)
1048+ } ) ;
1049+ } ) ;
9421050 }
9431051}
9441052
@@ -1424,8 +1532,13 @@ async fn handle_subgroup_stream(
14241532 // Control and data travel on independent streams. A subgroup stream may arrive
14251533 // just before the control loop processes the matching SubscribeOk(track_alias).
14261534 // Wait briefly for local activation before deciding this alias is unknown.
1427- let has_subscription =
1428- wait_for_active_subscription ( inner, track_alias, Duration :: from_millis ( 2_000 ) , Duration :: from_millis ( 5 ) ) . await ;
1535+ let has_subscription = wait_for_active_subscription (
1536+ inner,
1537+ track_alias,
1538+ Duration :: from_millis ( 2_000 ) ,
1539+ Duration :: from_millis ( 5 ) ,
1540+ )
1541+ . await ;
14291542
14301543 if !has_subscription {
14311544 return Ok ( ( ) ) ;
0 commit comments