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 | |
| 15 | use 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 | |
| 38 | use oxedyne_fe2o3_core::{ |
| 39 | prelude::*, |
| 40 | byte::FromBytes, |
| 41 | }; |
| 42 | use oxedyne_fe2o3_crypto::keys::PublicKey; |
| 43 | use oxedyne_fe2o3_hash::pow::PowVars; |
| 44 | use oxedyne_fe2o3_iop_crypto::keys::KeyManager; |
| 45 | use oxedyne_fe2o3_namex::InNamex; |
| 46 | use oxedyne_fe2o3_syntax::{ |
| 47 | core::SyntaxRef, |
| 48 | msg::Msg, |
| 49 | }; |
| 50 | use oxedyne_fe2o3_text::string::Stringer; |
| 51 | |
| 52 | use std::net::{ |
| 53 | IpAddr, |
| 54 | SocketAddr, |
| 55 | }; |
| 56 | |
| 57 | |
| 58 | #[derive(Clone, Debug)] |
| 59 | pub 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)] |
| 70 | pub 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 | |
| 82 | impl< |
| 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 | } |