Oregami
Repositories/oxedyne/fe2o3

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
21use 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
38use oxedyne_fe2o3_core::{
39 prelude::*,
40 rand::RanDef,
41};
42use oxedyne_fe2o3_syntax::SyntaxRef;
43
44use std::{
45 net::{
46 IpAddr,
47 SocketAddr,
48 },
49 sync::Arc,
50 time::Duration,
51};
52
53use tokio::net::UdpSocket;
54
55
56/// A peer that dials: it sends payloads and hears the answers to them.
57pub 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
73impl<
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}