Skip to main content

oxygengine_network_backend_native/
server.rs

1use crate::{client::NativeClient, utils::DoOnDrop};
2use network::{
3    client::{Client, ClientId, ClientState, MessageId},
4    server::{Server, ServerId, ServerState},
5};
6use std::{
7    collections::{HashMap, VecDeque},
8    io::ErrorKind,
9    net::TcpListener,
10    ops::Range,
11    sync::{Arc, RwLock},
12    thread::{sleep, Builder as ThreadBuilder, JoinHandle},
13    time::Duration,
14};
15
16const LISTENER_SLEEP_MS: u64 = 10;
17
18type MsgData = (ClientId, MessageId, Vec<u8>);
19
20pub struct NativeServer {
21    id: ServerId,
22    state: Arc<RwLock<ServerState>>,
23    clients: Arc<RwLock<HashMap<ClientId, NativeClient>>>,
24    clients_ids_cached: Vec<ClientId>,
25    messages: VecDeque<MsgData>,
26    thread: Option<JoinHandle<()>>,
27}
28
29impl Drop for NativeServer {
30    fn drop(&mut self) {
31        self.cleanup();
32    }
33}
34
35impl NativeServer {
36    fn cleanup(&mut self) {
37        if let Ok(mut state) = self.state.write() {
38            *state = ServerState::Closed;
39        }
40        if let Some(thread) = self.thread.take() {
41            thread.join().unwrap();
42        }
43    }
44}
45
46impl Server for NativeServer {
47    fn open(url: &str) -> Option<Self> {
48        let sid = ServerId::default();
49        let url = url.to_owned();
50        let state = Arc::new(RwLock::new(ServerState::Starting));
51        let state2 = state.clone();
52        let clients = Arc::new(RwLock::new(HashMap::default()));
53        let clients2 = clients.clone();
54        let thread = Some(
55            ThreadBuilder::new()
56                .name(format!("Server: {:?}", sid))
57                .spawn(move || {
58                    let state3 = state2.clone();
59                    let _ = DoOnDrop::new(move || {
60                        if let Ok(mut state) = state3.write() {
61                            *state = ServerState::Closed;
62                        }
63                    });
64                    let listener = TcpListener::bind(&url).unwrap();
65                    listener.set_nonblocking(true).unwrap_or_else(|_| {
66                        panic!(
67                            "Server {:?} cannot set non-blocking listening on: {}",
68                            sid, &url
69                        )
70                    });
71                    if let Ok(mut state) = state2.write() {
72                        *state = ServerState::Open;
73                    }
74                    for stream in listener.incoming() {
75                        if let Ok(state) = state2.read() {
76                            if *state == ServerState::Closed {
77                                break;
78                            }
79                        }
80                        match stream {
81                            Ok(stream) => {
82                                let client = NativeClient::from(stream);
83                                let id = client.id();
84                                if let Ok(mut clients) = clients2.write() {
85                                    clients.insert(id, client);
86                                }
87                            }
88                            Err(ref e) if e.kind() == ErrorKind::WouldBlock => {
89                                sleep(Duration::from_millis(LISTENER_SLEEP_MS));
90                            }
91                            Err(ref e) if e.kind() == ErrorKind::UnexpectedEof => {
92                                break;
93                            }
94                            Err(e) => {
95                                panic!("Server {:?} listener {} got IO error: {}", sid, &url, e)
96                            }
97                        }
98                    }
99                    if let Ok(mut state) = state2.write() {
100                        *state = ServerState::Closed;
101                    }
102                })
103                .unwrap(),
104        );
105        Some(Self {
106            id: sid,
107            state,
108            clients,
109            clients_ids_cached: vec![],
110            messages: Default::default(),
111            thread,
112        })
113    }
114
115    fn close(mut self) -> Self {
116        self.cleanup();
117        self
118    }
119
120    fn id(&self) -> ServerId {
121        self.id
122    }
123
124    fn state(&self) -> ServerState {
125        if let Ok(state) = self.state.read() {
126            *state
127        } else {
128            ServerState::default()
129        }
130    }
131
132    fn clients(&self) -> &[ClientId] {
133        &self.clients_ids_cached
134    }
135
136    fn disconnect(&mut self, id: ClientId) {
137        if let Ok(mut clients) = self.clients.write() {
138            if let Some(client) = clients.remove(&id) {
139                client.close();
140            }
141        }
142    }
143
144    fn disconnect_all(&mut self) {
145        if let Ok(mut clients) = self.clients.write() {
146            for (_, client) in clients.drain() {
147                client.close();
148            }
149        }
150    }
151
152    fn send(&mut self, id: ClientId, msg_id: MessageId, data: &[u8]) -> Option<Range<usize>> {
153        if self.state() != ServerState::Open {
154            return None;
155        }
156        if let Ok(mut clients) = self.clients.write() {
157            if let Some(client) = clients.get_mut(&id) {
158                if let Some(size) = client.send(msg_id, data) {
159                    return Some(size);
160                }
161            }
162        }
163        None
164    }
165
166    fn send_all(&mut self, id: MessageId, data: &[u8]) {
167        if self.state() != ServerState::Open {
168            return;
169        }
170        if let Ok(mut clients) = self.clients.write() {
171            for client in clients.values_mut() {
172                client.send(id, data);
173            }
174        }
175    }
176
177    fn read(&mut self) -> Option<(ClientId, MessageId, Vec<u8>)> {
178        self.messages.pop_front()
179    }
180
181    fn read_all(&mut self) -> Vec<MsgData> {
182        self.messages.drain(..).collect()
183    }
184
185    fn process(&mut self) {
186        if let Ok(mut clients) = self.clients.write() {
187            for (id, client) in clients.iter_mut() {
188                self.messages.extend(
189                    client
190                        .read_all()
191                        .into_iter()
192                        .map(|(mid, data)| (*id, mid, data)),
193                );
194            }
195            clients.retain(|_, client| client.state() != ClientState::Closed);
196            self.clients_ids_cached.clear();
197            for id in clients.keys() {
198                self.clients_ids_cached.push(*id);
199            }
200        }
201    }
202}