@@ -933,7 +933,7 @@ impl SupervisorSessionRegistry {
933933 payload : Some ( gateway_message:: Payload :: RelayOpen ( relay_open) ) ,
934934 } ;
935935 if tx. send ( msg) . await . is_err ( ) {
936- warn ! ( sandbox_id = %sandbox_id, channel_id = %channel_id, "supervisor session: failed to replay pending relay to superseding session" ) ;
936+ warn ! ( sandbox_id = %sandbox_id, channel_id = %channel_id, "supervisor session: failed to replay pending relay to new session" ) ;
937937 break ;
938938 }
939939 }
@@ -1960,12 +1960,12 @@ async fn establish_supervisor_session(
19601960 }
19611961 state. telemetry . sandbox_session_connected ( & sandbox_id) ;
19621962
1963- if superseded {
1964- state
1965- . supervisor_sessions
1966- . replay_pending_relays ( & sandbox_id , & tx )
1967- . await ;
1968- }
1963+ // The previous session may already have been removed before this one
1964+ // registers. Pending relay opens still need to reach the new session.
1965+ state
1966+ . supervisor_sessions
1967+ . replay_pending_relays ( & sandbox_id , & tx )
1968+ . await ;
19691969
19701970 // Step 4: Spawn the session loop that reads inbound messages.
19711971 let state_clone = Arc :: clone ( & state) ;
@@ -2971,6 +2971,52 @@ mod tests {
29712971 }
29722972 }
29732973
2974+ #[ tokio:: test]
2975+ async fn replay_pending_relays_reissues_open_after_disconnected_session ( ) {
2976+ let registry = SupervisorSessionRegistry :: new ( ) ;
2977+ let ( tx_old, mut rx_old) = mpsc:: channel :: < GatewayMessage > ( 4 ) ;
2978+ let ( tx_new, mut rx_new) = mpsc:: channel :: < GatewayMessage > ( 4 ) ;
2979+
2980+ registry. register (
2981+ "sbx" . to_string ( ) ,
2982+ "s-old" . to_string ( ) ,
2983+ tx_old,
2984+ make_shutdown ( ) ,
2985+ ) ;
2986+ let ( channel_id, _relay_rx) = registry
2987+ . open_relay ( "sbx" , Duration :: from_secs ( 1 ) )
2988+ . await
2989+ . expect ( "open_relay should succeed" ) ;
2990+ rx_old
2991+ . recv ( )
2992+ . await
2993+ . expect ( "old session should receive RelayOpen" ) ;
2994+
2995+ assert_eq ! ( registry. remove_if_current( "sbx" , "s-old" ) , Some ( false ) ) ;
2996+ let superseded = registry. register (
2997+ "sbx" . to_string ( ) ,
2998+ "s-new" . to_string ( ) ,
2999+ tx_new,
3000+ make_shutdown ( ) ,
3001+ ) ;
3002+ assert ! ( !superseded, "old session was removed before reconnect" ) ;
3003+
3004+ registry
3005+ . replay_pending_relays ( "sbx" , & registry. lookup_session ( "sbx" ) . unwrap ( ) )
3006+ . await ;
3007+
3008+ let replayed = rx_new
3009+ . recv ( )
3010+ . await
3011+ . expect ( "new session should receive RelayOpen" ) ;
3012+ match replayed. payload {
3013+ Some ( gateway_message:: Payload :: RelayOpen ( open) ) => {
3014+ assert_eq ! ( open. channel_id, channel_id) ;
3015+ }
3016+ other => panic ! ( "expected RelayOpen, got {other:?}" ) ,
3017+ }
3018+ }
3019+
29743020 #[ tokio:: test]
29753021 async fn require_persisted_sandbox_rejects_missing_sandbox ( ) {
29763022 let store = test_store ( ) . await ;
0 commit comments