@@ -32,6 +32,7 @@ use smallvec::SmallVec;
3232use tokio:: net:: UdpSocket ;
3333use tokio:: runtime:: Handle ;
3434use tokio:: sync:: { broadcast, mpsc} ;
35+ use tokio:: time:: Instant ;
3536
3637use vey_io_ext:: { UdpMoveRecv , UdpMoveSend , UdpSocketExt } ;
3738use vey_io_sys:: udp:: { RecvMsgHdr , SendMsgHdr } ;
@@ -96,7 +97,10 @@ impl UdpMoveRecv for AcceptedUdpPacketReceiver {
9697 fn poll_recv_packet ( & mut self , cx : & mut Context < ' _ > ) -> Poll < Result < Bytes , Self :: RecvError > > {
9798 match self . inner . poll_recv ( cx) {
9899 Poll :: Ready ( Some ( packet) ) => Poll :: Ready ( Ok ( packet) ) ,
99- Poll :: Ready ( None ) => Poll :: Ready ( Ok ( Bytes :: new ( ) ) ) ,
100+ Poll :: Ready ( None ) => Poll :: Ready ( Err ( io:: Error :: new (
101+ io:: ErrorKind :: UnexpectedEof ,
102+ "packet receiver closed" ,
103+ ) ) ) ,
100104 Poll :: Pending => Poll :: Pending ,
101105 }
102106 }
@@ -816,6 +820,16 @@ where
816820 mut rt_state : RuntimeState ,
817821 mut ct_table : LruCache < ClientConnectionKey , StreamDispatcher , FixedState > ,
818822 ) {
823+ let mut wait_sleep = Box :: pin ( tokio:: time:: sleep ( self . conn_track . offline_wait_time ( ) ) ) ;
824+ let mut allow_new = true ;
825+
826+ info ! (
827+ "SRT[{}_v{}#{}] enters offline-wait mode" ,
828+ self . server. name( ) ,
829+ self . server_version,
830+ self . instance_id
831+ ) ;
832+
819833 loop {
820834 tokio:: select! {
821835 biased;
@@ -825,14 +839,32 @@ where
825839 self . handle_events( & mut event_recv_buf, & mut ct_table) . await ;
826840 event_recv_buf. clear( ) ;
827841 if ct_table. is_empty( ) {
828- break ;
842+ return ;
829843 }
830844 }
831845 r = self . recv_packets( & rt_state. socket) => {
832846 match r {
833847 Ok ( packets) => {
834848 for ( cc_info, data) in packets {
835- self . handle_packet( cc_info, data, & rt_state, & mut ct_table) ;
849+ if allow_new {
850+ self . handle_packet( cc_info, data, & rt_state, & mut ct_table) ;
851+ } else {
852+ let key = cc_info. connection_key( ) ;
853+ if let Some ( dispatcher) = ct_table. get( & key) {
854+ match dispatcher. sender. try_send( data) {
855+ Ok ( _) => { }
856+ Err ( mpsc:: error:: TrySendError :: Full ( _) ) => {
857+ dispatcher. state. add_recv_dropped( ) ;
858+ }
859+ Err ( mpsc:: error:: TrySendError :: Closed ( _) ) => {
860+ dispatcher. state. add_recv_dropped( ) ;
861+ ct_table. pop( & key) ;
862+ }
863+ }
864+ } else {
865+ self . listen_stats. add_dropped( ) ;
866+ }
867+ }
836868 }
837869 }
838870 Err ( e) => {
@@ -842,6 +874,14 @@ where
842874 }
843875 }
844876 }
877+ _ = & mut wait_sleep => {
878+ if !allow_new {
879+ break ;
880+ }
881+ info!( "SRT[{}_v{}#{}] enters offline-quit mode" , self . server. name( ) , self . server_version, self . instance_id) ;
882+ allow_new = false ;
883+ wait_sleep. as_mut( ) . reset( Instant :: now( ) + self . conn_track. offline_quit_time( ) ) ;
884+ }
845885 }
846886 }
847887 }
0 commit comments