Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_shield/src/srv/server.rs

10.4 KiB, 141 runs

created by r1870400018:890, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! The UDP server loop: bind a socket, hear packets, and answer the ones that
2//! ask something.
3
4use crate::{
5 srv::{
6 constant,
7 context::ServerContext,
8 msg::{
9 app::{
10 Answer,
11 AppMsg,
12 AppMsgKind,
13 },
14 core::IdTypes,
15 decode::Received,
16 encode::ShieldCommand,
17 handshake::HReq1,
18 protocol::ProtocolTypes,
19 },
20 cmd::Command,
21 },
22};
23
24use oxedyne_fe2o3_core::{
25 prelude::*,
26 channels::{
27 Recv,
28 simplex,
29 Simplex,
30 },
31};
32use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
33use oxedyne_fe2o3_iop_db::api::Database;
34use oxedyne_fe2o3_iop_hash::api::Hasher;
35use oxedyne_fe2o3_syntax::SyntaxRef;
36
37use std::{
38 future::Future,
39 net::SocketAddr,
40 sync::Arc,
41 time::{
42 Duration,
43 Instant,
44 },
45};
46
47use tokio::net::UdpSocket;
48
49
50pub async fn answer_nothing(_payload: Vec<u8>, _src_addr: SocketAddr) -> Outcome<Answer> {
51 Ok(Answer::Nothing)
52}
53
54pub struct Server<
55 const C: usize,
56 const ML: usize,
57 const SL: usize,
58 const UL: usize,
59 P: ProtocolTypes<ML, SL, UL>,
60 // Database
61 ENC: Encrypter,
62 KH: Hasher,
63 DB: Database<UL, <P::ID as IdTypes<ML, SL, UL>>::U, ENC, KH>,
64> {
65 context: ServerContext<C, ML, SL, UL, P, ENC, KH, DB>,
66 syntax: SyntaxRef,
67 ma_gc_last: Instant,
68 ma_gc_int: Duration,
69 cmd_chan: Simplex<Command>,
70}
71
72impl<
73 const C: usize,
74 const ML: usize,
75 const SL: usize,
76 const UL: usize,
77 P: ProtocolTypes<ML, SL, UL> + 'static,
78 // Database
79 ENC: Encrypter + 'static,
80 KH: Hasher + 'static,
81 DB: Database<UL, <P::ID as IdTypes<ML, SL, UL>>::U, ENC, KH> + 'static,
82>
83 Server<C, ML, SL, UL, P, ENC, KH, DB>
84 where <P as ProtocolTypes<ML, SL, UL>>::W: 'static,
85{
86 pub fn new(
87 context: ServerContext<C, ML, SL, UL, P, ENC, KH, DB>,
88 syntax: SyntaxRef,
89 )
90 -> (Self, Simplex<Command>)
91 {
92 let cmd_chan = simplex();
93 let cmd_chan_clone = cmd_chan.clone();
94
95 (
96 Self {
97 context,
98 syntax,
99 ma_gc_last: Instant::now(),
100 ma_gc_int: constant::MSG_ASSEMBLY_GC_INTERVAL,
101 cmd_chan,
102 },
103 cmd_chan_clone,
104 )
105 }
106
107 pub async fn bind(&self) -> Outcome<Arc<UdpSocket>> {
108 let port = self.context.cfg.server_port_udp;
109 let ip = res!(self.context.cfg.bind_ip());
110 // The proof of work on every packet is bound to the address it was sent
111 // *to* as well as the address it came from, and a socket on the
112 // wildcard address cannot say which of this machine's addresses a
113 // datagram arrived at. A server bound there would therefore reject
114 // every packet, which is a worse way to find out than this one.
115 if ip.is_unspecified() {
116 return Err(err!(
117 "server_address is '{}', and a Shield server cannot listen on the \
118 wildcard address: the proof of work on each packet is bound to the \
119 address the packet was sent to, and a wildcard socket is not told \
120 which of this machine's addresses that was. Name one, or write \
121 'local' for whichever this machine has on its network.",
122 self.context.cfg.server_address;
123 Invalid, Configuration, Network));
124 }
125 let addr = SocketAddr::new(ip, port);
126 match UdpSocket::bind(addr).await {
127 Ok(sock) => Ok(Arc::new(sock)),
128 Err(e) => Err(err!(e,
129 "Could not bind the Shield UDP socket at {}.", addr;
130 IO, Network, Init)),
131 }
132 }
133
134 pub async fn start<H, F>(&mut self, handler: H) -> Outcome<()>
135 where
136 H: Fn(Vec<u8>, SocketAddr) -> F,
137 F: Future<Output = Outcome<Answer>>,
138 {
139 let trg = res!(self.bind().await);
140 self.run(trg, handler).await
141 }
142
143 pub async fn run<H, F>(
144 &mut self,
145 trg: Arc<UdpSocket>,
146 handler: H,
147 )
148 -> Outcome<()>
149 where
150 H: Fn(Vec<u8>, SocketAddr) -> F,
151 F: Future<Output = Outcome<Answer>>,
152 {
153 let trg_addr = res!(trg.local_addr(), IO, Network);
154 info!(async_log::stream(), "mode = {:?}", self.context.protocol.mode);
155 info!(async_log::stream(), "Listening on UDP at {}.", trg_addr);
156
157 let mut buf = [0u8; constant::UDP_BUFFER_SIZE];
158 'main: loop {
159 // Wake periodically even when nothing arrives, so the garbage collector and the
160 // command channel are not hostage to a quiet network.
161 match tokio::time::timeout(
162 constant::SERVER_EXT_SOCKET_CHECK_INTERVAL,
163 trg.recv_from(&mut buf),
164 ).await {
165 Err(_) => (), // Nothing arrived in this window.
166 Ok(Err(e)) => error!(async_log::stream(),
167 err!(e, "While trying to receive packet."; IO, Network)),
168 Ok(Ok((n, src_addr))) => {
169 if let Err(e) = self.serve(&buf[..n], src_addr, &trg, &handler).await {
170 error!(async_log::stream(), err!(e,
171 "While handling incoming packet from {}.", src_addr;
172 IO, Network));
173 }
174 },
175 }
176
177 // Message assembly garbage collection.
178 if self.ma_gc_last.elapsed() > self.ma_gc_int {
179 let result = self.context.protocol.massembler
180 .message_assembly_garbage_collection(&self.context.protocol.ma_params);
181 match result {
182 Err(e) => error!(async_log::stream(), err!(e,
183 "While attempting to collect message assembler garbage.";
184 IO, Network)),
185 Ok(_) => {}
186 }
187 self.ma_gc_last = Instant::now();
188 }
189
190 // Check internal command channel.
191 'cmd: loop {
192 match self.cmd_chan.try_recv() {
193 Recv::Empty => break 'cmd,
194 Recv::Result(Ok(Command::Finish)) => break 'main,
195 Recv::Result(Ok(cmd)) => {
196 test!(async_log::stream(), "Server command received: {:?}", cmd);
197 }
198 Recv::Result(Err(e)) => error!(async_log::stream(), err!(e,
199 "While reading command channel."; Channel, Read)),
200 }
201 }
202 }
203
204 Ok(())
205 }
206
207 async fn serve<H, F>(
208 &self,
209 buf: &[u8],
210 src_addr: SocketAddr,
211 trg: &Arc<UdpSocket>,
212 handler: &H,
213 )
214 -> Outcome<()>
215 where
216 H: Fn(Vec<u8>, SocketAddr) -> F,
217 F: Future<Output = Outcome<Answer>>,
218 {
219 let trg_ip = res!(trg.local_addr(), IO, Network).ip();
220 let protocol = self.context.protocol.clone();
221 let accepted = match res!(protocol.clone().accept(buf, src_addr, trg_ip)) {
222 Some(a) => a,
223 None => return Ok(()), // Dropped, or the message is still incomplete.
224 };
225 let mid = accepted.meta.mid;
226 let Received { fmt, pow, ids, msg } = res!(protocol.read(&accepted, self.syntax.clone()));
227
228 // Multiple commands in a single message are permitted.
229 for (cmd_name, mut msgcmd) in msg.cmds {
230 match cmd_name.as_str() {
231 "hreq1" => {
232 debug!(async_log::stream(), "HREQ1");
233 let mut scmd: HReq1<ML, SL, UL, P::ID> = HReq1 {
234 fmt: fmt.clone(),
235 pow: pow.clone(),
236 mid: ids.clone(),
237 ..Default::default()
238 };
239 // Each command type can implement its own custom process method, which
240 // captures only the parameters it needs.
241 let (akey, locked_amap) = res!(protocol.agrd.get_locked_map(&src_addr));
242 let mut unlocked_amap = lock_write!(locked_amap);
243 if let Some(alog) = unlocked_amap.get_mut(&akey) {
244 res!(scmd.respond(
245 &mut msgcmd,
246 &mut alog.data, // For pow parameters.
247 ));
248 }
249 },
250 other => match AppMsgKind::from_cmd_name(other) {
251 Some(AppMsgKind::Request) => {
252 let mut scmd: AppMsg<ML, SL, UL, P::ID> = AppMsg {
253 fmt: fmt.clone(),
254 pow: pow.clone(),
255 mid: ids.clone(),
256 kind: AppMsgKind::Request,
257 ..Default::default()
258 };
259 res!(scmd.deconstruct(&mut msgcmd));
260 let answer = res!(handler(scmd.payload, src_addr).await);
261 if let Answer::Reply(payload) = answer {
262 // The answer travels under the identifier the question arrived
263 // with, and back to the address it arrived from. Neither is a
264 // detail: a peer behind a router has no other address, and a
265 // peer holding two questions at once has no other way of
266 // telling the answers apart.
267 let packets = res!(protocol.build_app(
268 self.syntax.clone(),
269 AppMsgKind::Reply,
270 mid,
271 payload,
272 trg_ip,
273 src_addr.ip(),
274 ));
275 for packet in packets {
276 res!(trg.send_to(&packet, src_addr).await, IO, Network);
277 }
278 }
279 },
280 Some(AppMsgKind::Reply) => debug!(async_log::stream(),
281 "Dropping an application reply from {} to a question this peer \
282 did not ask.", src_addr),
283 None => return Err(err!(
284 "Unrecognised message command '{}'.", other;
285 Bug, Unimplemented)),
286 },
287 }
288 }
289 Ok(())
290 }
291}