Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_shield/src/srv/msg/decode.rs

18.2 KiB, 142 runs

created by r1870400018:4218, 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//! Taking an incoming packet apart, in the order a hostile network makes
2//! sensible.
3//!
4//! The work splits in two, and the split is what lets a peer that dials and a
5//! peer that listens share it. [`Protocol::accept`] does everything that is
6//! true of any packet whatever it turns out to say: rate-limit the address it
7//! came from, look up what proof of work is being demanded of that address,
8//! validate the artefacts, and hand the chunk to the assembler. What comes back
9//! is either nothing -- the message is still incomplete, or the packet was
10//! dropped -- or a whole message, which [`Protocol::read`] then parses against
11//! the syntax. Deciding what to *do* with the commands in it is the caller's,
12//! because a server answers requests and a client hears answers, and neither
13//! wants the other's dispatch table.
14
15use crate::{
16 srv::{
17 constant,
18 msg::{
19 core::{
20 IdTypes,
21 MsgFmt,
22 MsgIds,
23 MsgPow,
24 },
25 packet::{
26 PacketMeta,
27 PacketValidationArtefactRelativeIndices,
28 },
29 protocol::{
30 Protocol,
31 ProtocolTypes,
32 },
33 },
34 pow::PowPristine,
35 },
36};
37
38use oxedyne_fe2o3_core::{
39 prelude::*,
40 byte::FromBytes,
41};
42use oxedyne_fe2o3_crypto::keys::PublicKey;
43use oxedyne_fe2o3_hash::pow::PowVars;
44use oxedyne_fe2o3_iop_crypto::keys::KeyManager;
45use oxedyne_fe2o3_namex::InNamex;
46use oxedyne_fe2o3_syntax::{
47 core::SyntaxRef,
48 msg::Msg,
49};
50use oxedyne_fe2o3_text::string::Stringer;
51
52use std::net::{
53 IpAddr,
54 SocketAddr,
55};
56
57
58#[derive(Clone, Debug)]
59pub struct Accepted<
60 const MIDL: usize,
61 const UIDL: usize,
62 MID: oxedyne_fe2o3_jdat::id::NumIdDat<MIDL>,
63 UID: oxedyne_fe2o3_jdat::id::NumIdDat<UIDL>,
64> {
65 pub meta: PacketMeta<MIDL, UIDL, MID, UID>,
66 pub byts: Vec<u8>,
67}
68
69#[derive(Clone, Debug)]
70pub struct Received<
71 const SL: usize,
72 const UL: usize,
73 SID: oxedyne_fe2o3_jdat::id::NumIdDat<SL>,
74 UID: oxedyne_fe2o3_jdat::id::NumIdDat<UL>,
75> {
76 pub fmt: MsgFmt,
77 pub pow: MsgPow,
78 pub ids: MsgIds<SL, UL, SID, UID>,
79 pub msg: Msg,
80}
81
82impl<
83 const C: usize,
84 const ML: usize,
85 const SL: usize,
86 const UL: usize,
87 P: ProtocolTypes<ML, SL, UL> + 'static,
88>
89 Protocol<C, ML, SL, UL, P>
90{
91 pub fn accept(
92 mut self,
93 buf: &[u8],
94 src_addr: SocketAddr,
95 trg_ip: IpAddr,
96 )
97 -> Outcome<Option<Accepted<
98 ML,
99 UL,
100 <P::ID as IdTypes<ML, SL, UL>>::M,
101 <P::ID as IdTypes<ML, SL, UL>>::U,
102 >>>
103 {
104 let n = buf.len();
105 {
106 let mut unlocked_timer = lock_write!(self.timer);
107 unlocked_timer.update();
108 }
109 debug!(async_log::stream(), "incoming [{}]:", n);
110 for line in dump!(" {:02x}", &buf[..n], 32) {
111 debug!(async_log::stream(), "{}", line);
112 }
113 // Packet:
114 // validation
115 // artefacts
116 // |
117 // n1 n2 | n
118 // +-------------+--------------------------------+-------------+
119 // | | +----+ +----+
120 // | |
121 // | | |
122 // meta message validation
123 // chunk artefacts
124 //
125 // 1. Read meta data.
126 let (meta, n1) = res!(PacketMeta::<
127 ML,
128 UL,
129 <P::ID as IdTypes<ML, SL, UL>>::M,
130 <P::ID as IdTypes<ML, SL, UL>>::U,
131 >::from_bytes(&buf[..n])); // Decode packet meta.
132 debug!(async_log::stream(), "meta [{}]:", n1);
133 for line in Stringer::new(fmt!("{:?}", meta)).to_lines(" ") {
134 debug!(async_log::stream(), "{}", line);
135 }
136 //
137 // 1. First line of defence: rate limiting and blacklisting against the source address. We
138 // don't know if the sender of the packet is who they say they are, they could be
139 // address spoofing. The threat of primary concern is DDOS, so we are looking for any
140 // excuse to drop a packet before committing more resources or degrading service for
141 // good users. This check creates a new AddressLog entry if the source address is
142 // unknown and the request is an HREQ1. This precedes validation because we want to
143 // collect any custom validation parameters for this address.
144 if res!(crate::srv::guard::addr::drop_packet(
145 &*self.agrd,
146 self.hreq_exp,
147 meta.typ,
148 &src_addr,
149 )) {
150 debug!(async_log::stream(), "Address guard dropping packet.");
151 return Ok(None); // Drop silently.
152 }
153 if res!(self.ugrd.drop_packet(&meta.uid, self.accept_unknown)) { // Accesses the user log.
154 debug!(async_log::stream(), "User guard dropping packet.");
155 return Ok(None); // Drop silently.
156 }
157 debug!(async_log::stream(), "");
158 // A packet claiming a chunk it did not bring is a packet whose validation artefacts
159 // would be read from somebody else's bytes. Refuse it here rather than slicing past
160 // the end of the buffer.
161 let n2 = n1 + (meta.chnk.chunk_size as usize);
162 if n2 > n
163 || n - n2 < PacketValidationArtefactRelativeIndices::BYTE_PREFIX_LEN
164 {
165 debug!(async_log::stream(),
166 "Dropping packet of {} bytes: its header claims a {}-byte chunk after {} \
167 bytes of metadata, leaving no room for the validation artefacts.",
168 n, meta.chnk.chunk_size, n1);
169 return Ok(None); // Drop silently.
170 }
171 let (afact_rel_ind, _) =
172 res!(PacketValidationArtefactRelativeIndices::from_bytes(&buf[n2..n]));
173
174 // Get the (locked) shared address and user maps, and unlock them in tight scopes when we
175 // need to read or write.
176 let (akey, locked_amap) = res!(self.agrd.get_locked_map(&src_addr));
177 let (ukey, locked_umap) = res!(self.ugrd.get_locked_map(&meta.uid));
178
179 debug!(async_log::stream(), "");
180 // What are our proof of work requirements for the packet?
181 let powvars = match self.packval.pow {
182 Some(..) => {
183 let zbits = {
184 let unlocked_amap = lock_read!(locked_amap);
185 if let Some(alog) = unlocked_amap.get(&akey) {
186 let unlocked_timer = lock_read!(self.timer);
187 let zbits = res!(
188 self.gpzparams.required_global_zbits(unlocked_timer.avg_rps()),
189 IO,
190 );
191 if zbits >= alog.data.my_zbits {
192 zbits
193 } else {
194 alog.data.my_zbits
195 }
196 } else {
197 return Err(err!(
198 "No AddressLog entry for {:?}, which should have been created \
199 by the AddressGuard::drop_packet call.", src_addr;
200 Bug, Missing));
201 }
202 };
203 let code = {
204 let unlocked_umap = lock_read!(locked_umap);
205 if let Some(ulog) = unlocked_umap.get(&ukey) {
206 ulog.data.code.clone().unwrap_or([0; C])
207 } else {
208 return Err(err!(
209 "No UserLog entry for {:?}, which should have been created \
210 by the UserGuard::drop_packet call.", meta.uid;
211 Bug, Missing));
212 }
213 };
214 let pristine = res!(PowPristine::<
215 C,
216 {constant::POW_PREFIX_LEN},
217 {constant::POW_PREIMAGE_LEN},
218 >::new_rx(
219 code,
220 src_addr.ip(),
221 trg_ip,
222 self.pow_time_horiz,
223 ));
224 trace!(async_log::stream(), "POW Pristine rx:");
225 res!(pristine.trace());
226
227 Some(PowVars {
228 zbits,
229 pristine,
230 })
231 },
232 _ => None,
233 };
234 // Insert my record of your public signing key into the packet signer for the purpose of
235 // verification.
236 match &mut self.packval.sig {
237 Some(signer) => {
238 let unlocked_umap = lock_read!(locked_umap);
239 if let Some(ulog) = unlocked_umap.get(&ukey) {
240 let signer_nid = signer.local_id();
241 // The current signing scheme may differ from that for the public signing key I
242 // have on record, check it.
243 match &ulog.data.sigtpk_opt {
244 Some(sigtpk) => {
245 if sigtpk.sts.id != signer_nid {
246 return Err(err!(
247 "Local scheme id, {:?}, for public signing key of user, {:02x?}, does not \
248 match the nid for the current packet signing scheme, {:?}.",
249 sigtpk.sts.id, meta.uid, signer_nid;
250 Name, Mismatch));
251 }
252 // Update the signer with the public key I have for you.
253 *signer = res!(signer.clone_with_keys(Some(&sigtpk.key[..]), None));
254 },
255 None => (),
256 }
257 } else {
258 return Err(err!(
259 "No UserLog entry for {:02x?}, which should have been created \
260 by the UserGuard::drop_packet call.", meta.uid;
261 Bug, Missing));
262 }
263 },
264 _ => (),
265 }
266
267 //////// Debugging only
268 match &afact_rel_ind.pow {
269 Some(range) => {
270 let artefact = &buf[n2 + range.start..n2 + range.end];
271 trace!(async_log::stream(), "POW rx:");
272 res!(self.packval.trace(
273 powvars.as_ref(),
274 artefact,
275 ));
276 },
277 None => {
278 debug!(async_log::stream(),
279 "Dropping packet from {}: proof of work required, none supplied.",
280 src_addr);
281 return Ok(None); // Drop silently.
282 },
283 }
284 ////////
285
286 let validation = res!(self.packval.validate(
287 &buf[..n],
288 n2,
289 afact_rel_ind,
290 powvars,
291 meta.typ,
292 ));
293 debug!(async_log::stream(), "{:?}", validation);
294 let validity = fmt!("pow {} sig {}", validation.pow_state(), validation.sig_state());
295
296 match validation.is_valid() {
297 // sigpk_opt = possible public signing key that may be included in the packet
298 // validation artefact.
299 Some((valid, sigpk_opt)) => if !valid {
300 // TODO Take action on an invalid signature provided by this address and user id.
301 trace!(async_log::stream(), "Dropping packet: {}", validity);
302 return Ok(None); // Drop silently.
303 } else {
304 // The packet signature was valid.
305 debug!(async_log::stream(), "The packet is valid: {}", validity);
306 match sigpk_opt {
307 Some((nid, sigpk_given)) => {
308 // A public signing key was supplied, and was used for verification. My
309 // existing record of your public signing key, if it exists, was not used.
310 let mut unlocked_umap = lock_write!(locked_umap);
311 if let Some(ulog) = unlocked_umap.get_mut(&ukey) {
312 match &ulog.data.sigtpk_opt {
313 Some(sigtpk) => { // I have a record of your current public signing key.
314 if sigtpk.key != sigpk_given {
315 // The key you supplied doesn't match the one I've got.
316 // I'll record the one I've got as old, and you'll be asked
317 // to sign with it. I won't regard the key you supplied as
318 // genuine until you are validated using the old key.
319 ulog.data.sigtpk_opt_old = Some(sigtpk.clone());
320 } else {
321 // The key you supplied perfectly matches the one I've got.
322 match &ulog.data.sigtpk_opt_old {
323 Some(_sigtpk_old) => {
324 // I don't recognise the public key that you used. It is possible
325 // that I simply missed the key update. So find the latest public
326 // key I do have, in order to ask the peer to sign HReq2 using it,
327 // so I can be sure this is the user I think it is.
328 if let Some(pk) = ulog.data.pack_sigpk_set.first() {
329 ulog.data.sign_pack_this = Some(pk.key.clone());
330 }
331 },
332 None => {
333 // The earlier call to self.ugrd.drop_packet may have created a new
334 // entry for an unrecognised uid, but with no public signing key,
335 // I have no prior record of this user. Whether I accept them as
336 // a new user depends on our policy.
337 if self.accept_unknown {
338 ulog.data.sigtpk_opt = Some(res!(PublicKey::now(
339 nid,
340 sigtpk.key.clone(),
341 )));
342 } else {
343 // TODO If arranging for periodic garbage collection of users
344 // who lack packet public keys is more efficient, don't delete
345 // user just yet.
346 return Ok(None);
347 }
348 },
349 }
350 }
351 },
352 None => (), // TODO FINISHME I can't remember what is supposed to happen here!!!
353 }
354 } else {
355 return Err(err!(
356 "No UserLog entry for {:?}, which should have been created \
357 by the UserGuard::drop_packet call.", meta.uid;
358 Bug, Missing));
359 }
360 },
361 None => (), // The packet signature was valid, using the public key I possess.
362 }
363 },
364 None => (),
365 }
366 // Ok, we're almost done on a packet level. Insert the message chunk into the message
367 // assembler, which returns the message when complete. However, I may also have to drop
368 // the packet if there is a problem.
369 debug!(async_log::stream(), "");
370 match res!(self.massembler.get_msg( // Message checkpoint, drop the partial message?
371 &meta,
372 &buf[n1..n2], // payload chunk
373 &self.ma_params,
374 )) { // Returns whether to drop the packet, and the potential syntax protocol message.
375 (false, None) => Ok(None), // Payload remains incomplete.
376 (false, Some(byts)) => Ok(Some(Accepted { meta, byts })),
377 (true, _) => { // Drop the message completely.
378 res!(self.massembler.remove(&meta.mid));
379 Ok(None)
380 },
381 }
382 }
383
384 pub fn read(
385 &self,
386 accepted: &Accepted<
387 ML,
388 UL,
389 <P::ID as IdTypes<ML, SL, UL>>::M,
390 <P::ID as IdTypes<ML, SL, UL>>::U,
391 >,
392 syntax: SyntaxRef,
393 )
394 -> Outcome<Received<
395 SL,
396 UL,
397 <P::ID as IdTypes<ML, SL, UL>>::S,
398 <P::ID as IdTypes<ML, SL, UL>>::U,
399 >>
400 {
401 let msgrx = Msg::new(syntax.clone());
402 let mut msgrx = res!(msgrx.from_bytes(&accepted.byts, None));
403 debug!(async_log::stream(), "msgrx [{}]: {}", accepted.byts.len(), msgrx);
404 let ids: MsgIds<
405 SL,
406 UL,
407 <P::ID as IdTypes<ML, SL, UL>>::S,
408 <P::ID as IdTypes<ML, SL, UL>>::U,
409 > = res!(MsgIds::from_msg(
410 accepted.meta.uid,
411 &mut msgrx,
412 ));
413 let pow = res!(MsgPow::from_msg(&mut msgrx));
414 // The MsgFmt captures the syntax protocol against which incoming and outgoing
415 // messages are validated, and the encoding for any outgoing messages.
416 let fmt = MsgFmt {
417 syntax,
418 encoding: constant::DEFAULT_MSG_ENCODING, // TODO allow client to change
419 };
420 Ok(Received {
421 fmt,
422 pow,
423 ids,
424 msg: msgrx,
425 })
426 }
427}