11use std:: { fmt:: Debug , io, net:: SocketAddr , sync:: Arc , time:: Duration } ;
22
3+ use base:: Text ;
34use flume:: { Receiver , Sender } ;
45use futures_lite:: FutureExt ;
56use io:: ErrorKind ;
67use protocol:: {
7- codec:: CryptKey , ClientPlayPacket , MinecraftCodec , Readable , ServerPlayPacket , Writeable ,
8+ codec:: CryptKey , packets:: server:: Disconnect , ClientPlayPacket , MinecraftCodec , Readable ,
9+ ServerPlayPacket , Writeable ,
810} ;
911use tokio:: {
1012 io:: { AsyncReadExt , AsyncWriteExt } ,
@@ -18,6 +20,7 @@ use tokio::{
1820use crate :: {
1921 initial_handler:: { InitialHandling , NewPlayer } ,
2022 options:: Options ,
23+ player_count:: PlayerCount ,
2124} ;
2225
2326/// Tokio task which handles a connection and processes
@@ -32,6 +35,7 @@ pub struct Worker {
3235 reader : Reader ,
3336 writer : Writer ,
3437 options : Arc < Options > ,
38+ player_count : PlayerCount ,
3539 packets_to_send_tx : Sender < ServerPlayPacket > ,
3640 received_packets_rx : Receiver < ClientPlayPacket > ,
3741 new_players : Sender < NewPlayer > ,
@@ -42,6 +46,7 @@ impl Worker {
4246 stream : TcpStream ,
4347 _addr : SocketAddr ,
4448 options : Arc < Options > ,
49+ player_count : PlayerCount ,
4550 new_players : Sender < NewPlayer > ,
4651 ) -> Self {
4752 let ( reader, writer) = stream. into_split ( ) ;
@@ -55,6 +60,7 @@ impl Worker {
5560 reader,
5661 writer,
5762 options,
63+ player_count,
5864 packets_to_send_tx,
5965 received_packets_rx,
6066 new_players,
@@ -75,12 +81,22 @@ impl Worker {
7581 }
7682 }
7783
78- async fn proceed ( self , result : InitialHandling ) {
84+ async fn proceed ( mut self , result : InitialHandling ) {
7985 match result {
8086 InitialHandling :: Disconnect => ( ) ,
8187 InitialHandling :: Join ( new_player) => {
88+ if self . player_count . try_add_player ( ) . is_err ( ) {
89+ self . write ( ServerPlayPacket :: Disconnect ( Disconnect {
90+ reason : Text :: from ( "The server is full!" ) . to_string ( ) ,
91+ } ) )
92+ . await
93+ . ok ( ) ;
94+ return ;
95+ }
96+
97+ let username = new_player. username . clone ( ) ;
8298 let _ = self . new_players . send_async ( new_player) . await ;
83- self . split ( ) ;
99+ self . split ( username ) ;
84100 }
85101 }
86102 }
@@ -89,6 +105,10 @@ impl Worker {
89105 & self . options
90106 }
91107
108+ pub fn player_count ( & self ) -> u32 {
109+ self . player_count . get ( )
110+ }
111+
92112 #[ allow( unused) ]
93113 pub fn enable_compression ( & mut self , threshold : usize ) {
94114 self . reader . codec . enable_compression ( threshold) ;
@@ -108,16 +128,23 @@ impl Worker {
108128 self . writer . write ( packet) . await
109129 }
110130
111- pub fn split ( self ) {
112- let Self { reader, writer, .. } = self ;
131+ pub fn split ( self , username : String ) {
132+ let Self {
133+ reader,
134+ writer,
135+ player_count,
136+ ..
137+ } = self ;
113138 let reader = tokio:: task:: spawn ( async move { reader. run ( ) . await } ) ;
114139 let writer = tokio:: task:: spawn ( async move { writer. run ( ) . await } ) ;
115140
116141 tokio:: task:: spawn ( async move {
117142 let result = reader. race ( writer) . await . expect ( "task panicked" ) ;
118143 if let Err ( e) = result {
119- log:: error!( "Connection lost: {:?}" , e) ;
144+ let message = disconnected_message ( e) ;
145+ log:: debug!( "{} lost connection: {}" , username, message) ;
120146 }
147+ player_count. remove_player ( ) ;
121148 } ) ;
122149 }
123150
@@ -208,3 +235,12 @@ impl Writer {
208235 Ok ( ( ) )
209236 }
210237}
238+
239+ fn disconnected_message ( e : anyhow:: Error ) -> String {
240+ if let Some ( io_error) = e. downcast_ref :: < io:: Error > ( ) {
241+ if io_error. kind ( ) == ErrorKind :: UnexpectedEof {
242+ return "disconnected" . to_owned ( ) ;
243+ }
244+ }
245+ format ! ( "{:?}" , e)
246+ }
0 commit comments