oxedyne/fe2o3/fe2o3_shield/src/srv/client.rs
9.4 KiB, 10 runs
created by r1870400018:20653, 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 dialling half of a Shield exchange. |
| 2 | //! |
| 3 | //! A [`Client`] owns a UDP socket, sends an application payload on it, and |
| 4 | //! hears the answer on the same socket. That is the whole of it, and the shape |
| 5 | //! is chosen rather than inherited: a peer behind a household router can send |
| 6 | //! a datagram out and receive the reply to it, and cannot receive one that |
| 7 | //! arrives out of the blue. An answer sent to a fresh socket, or dialled back |
| 8 | //! to a listening port, is an answer half a real network never gets. |
| 9 | //! |
| 10 | //! Nothing here is a lesser peer. The client validates what arrives with the |
| 11 | //! same [`Protocol`] a server does -- the same guards, the same proof-of-work |
| 12 | //! and signature checks, the same message assembler -- because a reply is as |
| 13 | //! much somebody else's bytes as a request is. What it does not do is dispatch |
| 14 | //! requests: it hears answers to its own questions, and drops everything else. |
| 15 | //! |
| 16 | //! ```ignore |
| 17 | //! let client = res!(Client::bind(bind_addr, protocol, syntax).await); |
| 18 | //! let answer = res!(client.ask(peer_addr, payload, constant::APP_REPLY_WAIT).await); |
| 19 | //! ``` |
| 20 | |
| 21 | use crate::srv::{ |
| 22 | constant, |
| 23 | msg::{ |
| 24 | app::{ |
| 25 | AppMsg, |
| 26 | AppMsgKind, |
| 27 | }, |
| 28 | core::IdTypes, |
| 29 | decode::Received, |
| 30 | encode::ShieldCommand, |
| 31 | protocol::{ |
| 32 | Protocol, |
| 33 | ProtocolTypes, |
| 34 | }, |
| 35 | }, |
| 36 | }; |
| 37 | |
| 38 | use oxedyne_fe2o3_core::{ |
| 39 | prelude::*, |
| 40 | rand::RanDef, |
| 41 | }; |
| 42 | use oxedyne_fe2o3_syntax::SyntaxRef; |
| 43 | |
| 44 | use std::{ |
| 45 | net::{ |
| 46 | IpAddr, |
| 47 | SocketAddr, |
| 48 | }, |
| 49 | sync::Arc, |
| 50 | time::Duration, |
| 51 | }; |
| 52 | |
| 53 | use tokio::net::UdpSocket; |
| 54 | |
| 55 | |
| 56 | /// A peer that dials: it sends payloads and hears the answers to them. |
| 57 | pub struct Client< |
| 58 | const C: usize, |
| 59 | const ML: usize, |
| 60 | const SL: usize, |
| 61 | const UL: usize, |
| 62 | P: ProtocolTypes<ML, SL, UL>, |
| 63 | > { |
| 64 | /// The socket questions leave on and answers arrive on. One socket, because |
| 65 | /// that is what makes the answers arrive at all. |
| 66 | sock: Arc<UdpSocket>, |
| 67 | /// Guards, validator, schemes and message assembler. |
| 68 | protocol: Protocol<C, ML, SL, UL, P>, |
| 69 | /// Syntax messages are built against and validated by. |
| 70 | syntax: SyntaxRef, |
| 71 | } |
| 72 | |
| 73 | impl< |
| 74 | const C: usize, |
| 75 | const ML: usize, |
| 76 | const SL: usize, |
| 77 | const UL: usize, |
| 78 | P: ProtocolTypes<ML, SL, UL> + 'static, |
| 79 | > |
| 80 | Client<C, ML, SL, UL, P> |
| 81 | where <P as ProtocolTypes<ML, SL, UL>>::W: 'static, |
| 82 | { |
| 83 | /// Bind a socket to dial from. |
| 84 | /// |
| 85 | /// A port of zero is the usual thing to ask for: a peer that only dials |
| 86 | /// wants whatever port the operating system has going spare, and the peers |
| 87 | /// it talks to learn the address from the packets themselves. |
| 88 | pub async fn bind( |
| 89 | addr: SocketAddr, |
| 90 | protocol: Protocol<C, ML, SL, UL, P>, |
| 91 | syntax: SyntaxRef, |
| 92 | ) |
| 93 | -> Outcome<Self> |
| 94 | { |
| 95 | let sock = match UdpSocket::bind(addr).await { |
| 96 | Ok(s) => s, |
| 97 | Err(e) => return Err(err!(e, |
| 98 | "Could not bind a Shield client socket at {}.", addr; |
| 99 | IO, Network, Init)), |
| 100 | }; |
| 101 | Ok(Self { |
| 102 | sock: Arc::new(sock), |
| 103 | protocol, |
| 104 | syntax, |
| 105 | }) |
| 106 | } |
| 107 | |
| 108 | /// Where this client is dialling from. |
| 109 | pub fn local_addr(&self) -> Outcome<SocketAddr> { |
| 110 | Ok(res!(self.sock.local_addr(), IO, Network)) |
| 111 | } |
| 112 | |
| 113 | /// The socket, for a caller that has its own reason to hold one. |
| 114 | pub fn socket(&self) -> Arc<UdpSocket> { |
| 115 | self.sock.clone() |
| 116 | } |
| 117 | |
| 118 | /// The address this client's packets will appear to have come from, at |
| 119 | /// `trg_addr`. |
| 120 | /// |
| 121 | /// A socket bound to the unspecified address -- which is what a peer that |
| 122 | /// does not care which interface it leaves by asks for -- reports that |
| 123 | /// address as its own, and the receiver sees the interface the routing |
| 124 | /// table actually chose. The proof of work is bound to both ends' |
| 125 | /// addresses, so the two have to agree on which one the sender has, and |
| 126 | /// the sender is the one that can find out: a throwaway socket connected |
| 127 | /// to the same destination is given the same source address by the same |
| 128 | /// routing table. |
| 129 | pub async fn source_ip(&self, trg_addr: SocketAddr) -> Outcome<IpAddr> { |
| 130 | let local = res!(self.local_addr()); |
| 131 | if !local.ip().is_unspecified() { |
| 132 | return Ok(local.ip()); |
| 133 | } |
| 134 | let any = match trg_addr { |
| 135 | SocketAddr::V4(..) => "0.0.0.0:0", |
| 136 | SocketAddr::V6(..) => "[::]:0", |
| 137 | }; |
| 138 | let probe = match UdpSocket::bind(any).await { |
| 139 | Ok(s) => s, |
| 140 | Err(e) => return Err(err!(e, |
| 141 | "Could not find which address packets to {} leave by.", trg_addr; |
| 142 | IO, Network)), |
| 143 | }; |
| 144 | if let Err(e) = probe.connect(trg_addr).await { |
| 145 | return Err(err!(e, |
| 146 | "Could not find which address packets to {} leave by.", trg_addr; |
| 147 | IO, Network)); |
| 148 | } |
| 149 | Ok(res!(probe.local_addr(), IO, Network).ip()) |
| 150 | } |
| 151 | |
| 152 | /// Send a payload and return the message identifier it went out under. |
| 153 | /// |
| 154 | /// Nothing is waited for. A caller with something to tell rather than to |
| 155 | /// ask stops here; one that wants the answer passes the identifier to |
| 156 | /// [`Client::hear`]. |
| 157 | pub async fn tell( |
| 158 | &self, |
| 159 | trg_addr: SocketAddr, |
| 160 | payload: Vec<u8>, |
| 161 | ) |
| 162 | -> Outcome<<P::ID as IdTypes<ML, SL, UL>>::M> |
| 163 | { |
| 164 | let mid = <P::ID as IdTypes<ML, SL, UL>>::M::randef(); |
| 165 | let src_ip = res!(self.source_ip(trg_addr).await); |
| 166 | let packets = res!(self.protocol.build_app( |
| 167 | self.syntax.clone(), |
| 168 | AppMsgKind::Request, |
| 169 | mid, |
| 170 | payload, |
| 171 | src_ip, |
| 172 | trg_addr.ip(), |
| 173 | )); |
| 174 | for packet in packets { |
| 175 | res!(self.sock.send_to(&packet, trg_addr).await, IO, Network); |
| 176 | } |
| 177 | Ok(mid) |
| 178 | } |
| 179 | |
| 180 | /// Wait for the answer to the question sent under `mid`. |
| 181 | /// |
| 182 | /// Anything else that turns up on the socket in the meantime is dropped and |
| 183 | /// the wait goes on: a packet that fails its guards or its validation, a |
| 184 | /// message under an identifier this peer never sent, a request from |
| 185 | /// somebody who mistook this socket for a server. None of those is the |
| 186 | /// answer, and none of them shortens the time the answer has to arrive in. |
| 187 | pub async fn hear( |
| 188 | &self, |
| 189 | mid: &<P::ID as IdTypes<ML, SL, UL>>::M, |
| 190 | wait: Duration, |
| 191 | ) |
| 192 | -> Outcome<Vec<u8>> |
| 193 | { |
| 194 | let deadline = tokio::time::Instant::now() + wait; |
| 195 | let mut buf = [0u8; constant::UDP_BUFFER_SIZE]; |
| 196 | loop { |
| 197 | let left = deadline.saturating_duration_since(tokio::time::Instant::now()); |
| 198 | if left.is_zero() { |
| 199 | return Err(err!( |
| 200 | "Nothing answered message {} within {:?}.", mid, wait; |
| 201 | IO, Network, Timeout)); |
| 202 | } |
| 203 | let (n, trg_addr) = match tokio::time::timeout( |
| 204 | left, |
| 205 | self.sock.recv_from(&mut buf), |
| 206 | ).await { |
| 207 | Err(_) => return Err(err!( |
| 208 | "Nothing answered message {} within {:?}.", mid, wait; |
| 209 | IO, Network, Timeout)), |
| 210 | Ok(Err(e)) => return Err(err!(e, |
| 211 | "While waiting for the answer to message {}.", mid; |
| 212 | IO, Network)), |
| 213 | Ok(Ok(pair)) => pair, |
| 214 | }; |
| 215 | // The address the answer was sent to is the one the question left |
| 216 | // by, which is what the sender bound its proof of work to. |
| 217 | let src_ip = res!(self.source_ip(trg_addr).await); |
| 218 | let accepted = match self.protocol.clone().accept(&buf[..n], trg_addr, src_ip) { |
| 219 | Ok(Some(a)) => a, |
| 220 | Ok(None) => continue, // Dropped, or the message is still incomplete. |
| 221 | Err(e) => { |
| 222 | warn!(async_log::stream(), |
| 223 | "While reading a packet from {}: {}", trg_addr, e); |
| 224 | continue; |
| 225 | }, |
| 226 | }; |
| 227 | if accepted.meta.mid != *mid { |
| 228 | debug!(async_log::stream(), |
| 229 | "A message under identifier {} arrived from {} while waiting on {}; \ |
| 230 | dropped.", accepted.meta.mid, trg_addr, mid); |
| 231 | continue; |
| 232 | } |
| 233 | let Received { msg, .. } = res!(self.protocol.read(&accepted, self.syntax.clone())); |
| 234 | for (cmd_name, mut msgcmd) in msg.cmds { |
| 235 | match AppMsgKind::from_cmd_name(cmd_name.as_str()) { |
| 236 | Some(AppMsgKind::Reply) => { |
| 237 | let mut scmd: AppMsg<ML, SL, UL, P::ID> = AppMsg { |
| 238 | kind: AppMsgKind::Reply, |
| 239 | ..Default::default() |
| 240 | }; |
| 241 | res!(scmd.deconstruct(&mut msgcmd)); |
| 242 | return Ok(scmd.payload); |
| 243 | }, |
| 244 | _ => debug!(async_log::stream(), |
| 245 | "A '{}' arrived from {} under the identifier {} was waiting on, \ |
| 246 | which is not an answer; dropped.", cmd_name, trg_addr, mid), |
| 247 | } |
| 248 | } |
| 249 | } |
| 250 | } |
| 251 | |
| 252 | /// Send a payload and wait for the answer to it. |
| 253 | pub async fn ask( |
| 254 | &self, |
| 255 | trg_addr: SocketAddr, |
| 256 | payload: Vec<u8>, |
| 257 | wait: Duration, |
| 258 | ) |
| 259 | -> Outcome<Vec<u8>> |
| 260 | { |
| 261 | let mid = res!(self.tell(trg_addr, payload).await); |
| 262 | self.hear(&mid, wait).await |
| 263 | } |
| 264 | } |