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 | |
| 4 | use 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 | |
| 24 | use oxedyne_fe2o3_core::{ |
| 25 | prelude::*, |
| 26 | channels::{ |
| 27 | Recv, |
| 28 | simplex, |
| 29 | Simplex, |
| 30 | }, |
| 31 | }; |
| 32 | use oxedyne_fe2o3_iop_crypto::enc::Encrypter; |
| 33 | use oxedyne_fe2o3_iop_db::api::Database; |
| 34 | use oxedyne_fe2o3_iop_hash::api::Hasher; |
| 35 | use oxedyne_fe2o3_syntax::SyntaxRef; |
| 36 | |
| 37 | use std::{ |
| 38 | future::Future, |
| 39 | net::SocketAddr, |
| 40 | sync::Arc, |
| 41 | time::{ |
| 42 | Duration, |
| 43 | Instant, |
| 44 | }, |
| 45 | }; |
| 46 | |
| 47 | use tokio::net::UdpSocket; |
| 48 | |
| 49 | |
| 50 | pub async fn answer_nothing(_payload: Vec<u8>, _src_addr: SocketAddr) -> Outcome<Answer> { |
| 51 | Ok(Answer::Nothing) |
| 52 | } |
| 53 | |
| 54 | pub 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 | |
| 72 | impl< |
| 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 | } |