oxygengine_network_backend_native/
server.rs1use 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}