oxedyne/fe2o3/fe2o3_shield/src/srv/msg/encode.rs
10.4 KiB, 145 runs
created by r1870400018:4342, 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 | use crate::{ |
| 2 | srv::{ |
| 3 | constant, |
| 4 | cfg::ServerConfig, |
| 5 | schemes::{ |
| 6 | WireSchemes, |
| 7 | WireSchemeTypes, |
| 8 | }, |
| 9 | msg::{ |
| 10 | app::{ |
| 11 | AppMsg, |
| 12 | AppMsgKind, |
| 13 | }, |
| 14 | core::{ |
| 15 | IdentifiedMessage, |
| 16 | IdTypes, |
| 17 | MsgFmt, |
| 18 | MsgIds, |
| 19 | MsgPow, |
| 20 | }, |
| 21 | packet::{ |
| 22 | PacketChunkState, |
| 23 | PacketCount, |
| 24 | PacketMeta, |
| 25 | PacketValidator, |
| 26 | }, |
| 27 | protocol::{ |
| 28 | Protocol, |
| 29 | ProtocolTypes, |
| 30 | }, |
| 31 | }, |
| 32 | pow::{ |
| 33 | PowPristine, |
| 34 | }, |
| 35 | }, |
| 36 | }; |
| 37 | |
| 38 | use oxedyne_fe2o3_core::{ |
| 39 | prelude::*, |
| 40 | byte::{ |
| 41 | Encoding, |
| 42 | IntoBytes, |
| 43 | ToBytes, |
| 44 | }, |
| 45 | }; |
| 46 | use oxedyne_fe2o3_hash::{ |
| 47 | pow::{ |
| 48 | PowCreateParams, |
| 49 | PowVars, |
| 50 | Pristine, |
| 51 | ProofOfWork, |
| 52 | ZeroBits, |
| 53 | }, |
| 54 | }; |
| 55 | use oxedyne_fe2o3_syntax::{ |
| 56 | SyntaxRef, |
| 57 | msg::{ |
| 58 | Msg, |
| 59 | MsgCmd, |
| 60 | }, |
| 61 | }; |
| 62 | |
| 63 | use std::{ |
| 64 | clone::Clone, |
| 65 | net::{ |
| 66 | IpAddr, |
| 67 | SocketAddr, |
| 68 | UdpSocket, |
| 69 | }, |
| 70 | sync::Arc, |
| 71 | time::{ |
| 72 | SystemTime, |
| 73 | UNIX_EPOCH, |
| 74 | }, |
| 75 | }; |
| 76 | |
| 77 | |
| 78 | /// Rather than a generic and possibly more complex callback mechanism, the processing of server |
| 79 | /// command is customised so as to access only parameters needed from the server loop scope. |
| 80 | /// Incoming server commands are encoded in a `oxedyne_fe2o3_syntax::msg::MsgCmd` using the `Syntax` |
| 81 | /// defined in `oxedyne_fe2o3_shield::syntax`. Each must be associated with a `struct` below that |
| 82 | /// is accessed in `oxedyne_fe2o3_o3db_sync::bots::bot_server`. This must capture some basic information (i.e. |
| 83 | /// `MsgFmt` and `MsgIds`) as well as the command-specific data. The associated `struct` must have |
| 84 | /// its own custom method for processing the incoming command (e.g. |
| 85 | /// `oxedyne_fe2o3_o3db_sync::comm::wire::HReq1::process`), and should implement `ShieldCommand` in order to |
| 86 | /// access supporting methods. There are plenty of examples to copy and modify. |
| 87 | pub trait ShieldCommand< |
| 88 | const ML: usize, |
| 89 | const SL: usize, |
| 90 | const UL: usize, |
| 91 | ID: IdTypes<ML, SL, UL>, |
| 92 | >: |
| 93 | Default |
| 94 | + IdentifiedMessage |
| 95 | + IntoBytes |
| 96 | { |
| 97 | fn fmt(&self) -> &MsgFmt; |
| 98 | fn pow(&self) -> &MsgPow; |
| 99 | fn mid(&self) -> &MsgIds<SL, UL, ID::S, ID::U>; |
| 100 | fn syntax(&self) -> &SyntaxRef { &self.fmt().syntax } |
| 101 | fn encoding(&self) -> &Encoding { &self.fmt().encoding } |
| 102 | fn uid(&self) -> ID::U { self.mid().uid.clone() } |
| 103 | fn sid_opt(&self) -> Option<ID::S> { |
| 104 | self.mid().sid_opt.as_ref().clone().copied() |
| 105 | } |
| 106 | fn pow_zbits(&self) -> ZeroBits { self.pow().zbits } |
| 107 | fn pad_last(&self) -> bool { true } |
| 108 | //fn pow_code(&self) -> Option<[u8; constant::POW_CODE_LEN]> { self.pow().code } |
| 109 | fn inc_sigpk(&self) -> bool; // Include signature public key in outgoing validator? |
| 110 | fn deconstruct(&mut self, _mcmd: &mut MsgCmd) -> Outcome<()> { Ok(()) } |
| 111 | fn construct(self) -> Outcome<Msg>; |
| 112 | |
| 113 | fn build< |
| 114 | const C: usize, |
| 115 | // Proof of work validator. |
| 116 | const N: usize, // Hash pre-image size = pristine + nonce sizes. |
| 117 | const P0: usize, // Length of private prefix bytes (i.e. not included in artefact). |
| 118 | const P1: usize, // Length of pristine bytes (i.e. included in artefact). |
| 119 | PRIS: Pristine<P0, P1>, // Pristine supplied to hasher. |
| 120 | W: WireSchemeTypes + 'static, // Contains the chunker, the pow hasher and the signer. |
| 121 | >( |
| 122 | self, |
| 123 | mid: ID::M, |
| 124 | src_addr: IpAddr, |
| 125 | trg_addr: IpAddr, |
| 126 | code: [u8; C], |
| 127 | schms: WireSchemes<W>, |
| 128 | ) |
| 129 | -> Outcome<Vec<Vec<u8>>> |
| 130 | { |
| 131 | // Copy some self parameters before consumption by into_bytes |
| 132 | let msg_name = self.name(); |
| 133 | //let uid = res!(self.uid().to_bytes(Vec::new())); |
| 134 | let inc_sigpk = self.inc_sigpk(); |
| 135 | let pad_last = self.pad_last(); |
| 136 | |
| 137 | let uid = self.uid().clone(); |
| 138 | |
| 139 | let zbits = self.pow_zbits(); |
| 140 | let typ = self.typ(); |
| 141 | let msg_byts = res!(self.into_bytes(Vec::new())); |
| 142 | //let tstamp = res!(SystemTime::now().duration_since(SystemTime::UNIX_EPOCH)).as_secs(); |
| 143 | |
| 144 | let pristine = PowPristine::<C, P0, P1> { |
| 145 | code, |
| 146 | src_addr, |
| 147 | trg_addr, |
| 148 | timestamp: res!(SystemTime::now().duration_since(UNIX_EPOCH)), |
| 149 | time_horiz: constant::POW_TIME_HORIZON_SEC, |
| 150 | }; |
| 151 | trace!(async_log::stream(), "POW Pristine tx:"); |
| 152 | res!(pristine.trace()); |
| 153 | |
| 154 | let validator = PacketValidator { |
| 155 | pow: Some(res!(ProofOfWork::new(schms.powh.clone()))), |
| 156 | sig: Some(schms.sign.clone()), |
| 157 | }; |
| 158 | |
| 159 | let powparams = PowCreateParams { |
| 160 | pvars: PowVars { |
| 161 | zbits, |
| 162 | pristine, |
| 163 | }, |
| 164 | time_lim: constant::POW_CREATE_TIMEOUT, |
| 165 | count_lim: constant::POW_CREATE_COUNT_LIM, |
| 166 | }; |
| 167 | |
| 168 | let chunk_cfg = schms.chnk.clone(); |
| 169 | let chunker = ServerConfig::chunker(chunk_cfg.set_pad_last(pad_last)); |
| 170 | trace!(async_log::stream(), "{:?}", chunker); |
| 171 | |
| 172 | let size = chunker.cfg.chunk_size; |
| 173 | let meta_len = PacketMeta::<ML, UL, ID::M, ID::U>::BYTE_LEN; |
| 174 | let warning = if 2 * meta_len > size { |
| 175 | Some(errmsg!("Message meta of length {} bytes is more than half \ |
| 176 | the specified packet size of {}. Consider increasing the \ |
| 177 | packet size.", meta_len, size, |
| 178 | )) |
| 179 | } else { |
| 180 | None |
| 181 | }; |
| 182 | |
| 183 | let msg_len = msg_byts.len(); |
| 184 | let (mut chunks, _) = res!(chunker.chunk(&msg_byts)); |
| 185 | let nc = chunks.len(); |
| 186 | if nc > PacketCount::MAX as usize { |
| 187 | return Err(err!("Message type {} of length {} bytes, \ |
| 188 | when broken into chunks of {} bytes creates {} packets, \ |
| 189 | exceeding the limit of {}. Reduce the message length or \ |
| 190 | increase the packet size.", |
| 191 | msg_name, msg_byts.len(), size, nc, PacketCount::MAX; |
| 192 | Invalid, Configuration)); |
| 193 | } |
| 194 | |
| 195 | let mut packets = Vec::new(); |
| 196 | for i in 0..nc { |
| 197 | let chunk_len = chunks[i].len(); |
| 198 | let chnk = PacketChunkState { |
| 199 | index: res!(i.try_into()), |
| 200 | num_chunks: res!(nc.try_into()), |
| 201 | //chunk_size: res!(chunker.config().chunk_size().try_into()), |
| 202 | chunk_size: res!(chunk_len.try_into()), |
| 203 | pad_last: chunker.cfg.pad_last, |
| 204 | }; |
| 205 | let meta = PacketMeta { |
| 206 | typ, |
| 207 | ver: constant::VERSION, |
| 208 | mid, |
| 209 | uid, |
| 210 | chnk, |
| 211 | //tstamp, |
| 212 | }; |
| 213 | // 1. Header |
| 214 | let mut packet = res!(meta.to_bytes(Vec::new())); |
| 215 | let meta_len = packet.len(); |
| 216 | // 2. Data chunk |
| 217 | packet.append(&mut chunks[i]); |
| 218 | let len = packet.len(); |
| 219 | // 3. Validators |
| 220 | packet = res!(validator.to_bytes::<N, P0, P1, PowPristine<C, P0, P1>>( |
| 221 | packet, |
| 222 | &powparams, |
| 223 | //powparams.clone(), |
| 224 | inc_sigpk, |
| 225 | )); |
| 226 | let validator_len = packet.len() - len; |
| 227 | trace!(async_log::stream(), "Packet {} lengths: msg {}, meta {} chunk {} valid {} total {}", |
| 228 | i, msg_len, meta_len, chunk_len, validator_len, packet.len(), |
| 229 | ); |
| 230 | trace!(async_log::stream(), " Chunk: {}", chunks[i].len()); |
| 231 | packets.push(packet); |
| 232 | } |
| 233 | |
| 234 | if let Some(warning) = warning { |
| 235 | warn!(async_log::stream(), "{}", warning); |
| 236 | } |
| 237 | Ok(packets) |
| 238 | } |
| 239 | |
| 240 | fn send_udp( |
| 241 | src_sock: &UdpSocket, |
| 242 | trg_addr: &SocketAddr, |
| 243 | packets: Vec<Vec<u8>>, |
| 244 | ) |
| 245 | -> Outcome<()> |
| 246 | { |
| 247 | for packet in packets { |
| 248 | res!(src_sock.send_to(&packet, &trg_addr)); |
| 249 | } |
| 250 | Ok(()) |
| 251 | } |
| 252 | |
| 253 | fn build_standard< |
| 254 | const C: usize, |
| 255 | W: WireSchemeTypes + 'static, |
| 256 | >( |
| 257 | self, |
| 258 | mid: ID::M, |
| 259 | src_addr: IpAddr, |
| 260 | trg_addr: IpAddr, |
| 261 | code: [u8; C], |
| 262 | schms: WireSchemes<W>, |
| 263 | ) |
| 264 | -> Outcome<Vec<Vec<u8>>> |
| 265 | { |
| 266 | self.build::< |
| 267 | C, |
| 268 | {constant::POW_INPUT_LEN}, // N |
| 269 | {constant::POW_PREFIX_LEN}, // P0 |
| 270 | {constant::POW_PREIMAGE_LEN}, // P1 |
| 271 | PowPristine< |
| 272 | C, |
| 273 | {constant::POW_PREFIX_LEN}, |
| 274 | {constant::POW_PREIMAGE_LEN}, |
| 275 | >, |
| 276 | W, |
| 277 | >( |
| 278 | mid, |
| 279 | src_addr, |
| 280 | trg_addr, |
| 281 | code, |
| 282 | schms, |
| 283 | ) |
| 284 | } |
| 285 | |
| 286 | fn send< |
| 287 | const C: usize, |
| 288 | W: WireSchemeTypes + 'static, |
| 289 | >( |
| 290 | self, |
| 291 | mid: ID::M, |
| 292 | src: Arc<UdpSocket>, |
| 293 | trg_addr: &SocketAddr, |
| 294 | code: [u8; C], |
| 295 | schms: WireSchemes<W>, |
| 296 | ) |
| 297 | -> Outcome<()> |
| 298 | { |
| 299 | let packets = res!(self.build_standard::<C, W>( |
| 300 | mid, |
| 301 | res!(src.local_addr()).ip(), |
| 302 | trg_addr.ip(), |
| 303 | code, |
| 304 | schms, |
| 305 | )); |
| 306 | for packet in packets { |
| 307 | res!(src.send_to(&packet, trg_addr)); |
| 308 | } |
| 309 | Ok(()) |
| 310 | } |
| 311 | } |
| 312 | |
| 313 | impl< |
| 314 | const C: usize, |
| 315 | const ML: usize, |
| 316 | const SL: usize, |
| 317 | const UL: usize, |
| 318 | P: ProtocolTypes<ML, SL, UL> + 'static, |
| 319 | > |
| 320 | Protocol<C, ML, SL, UL, P> |
| 321 | where <P as ProtocolTypes<ML, SL, UL>>::W: 'static, |
| 322 | { |
| 323 | pub fn build_app( |
| 324 | &self, |
| 325 | syntax: SyntaxRef, |
| 326 | kind: AppMsgKind, |
| 327 | mid: <P::ID as IdTypes<ML, SL, UL>>::M, |
| 328 | payload: Vec<u8>, |
| 329 | src_ip: IpAddr, |
| 330 | trg_ip: IpAddr, |
| 331 | ) |
| 332 | -> Outcome<Vec<Vec<u8>>> |
| 333 | { |
| 334 | let cmd: AppMsg<ML, SL, UL, P::ID> = AppMsg { |
| 335 | fmt: MsgFmt { |
| 336 | syntax, |
| 337 | encoding: constant::DEFAULT_MSG_ENCODING, |
| 338 | }, |
| 339 | pow: MsgPow { zbits: self.tx_zbits }, |
| 340 | mid: MsgIds { |
| 341 | sid_opt: None, |
| 342 | uid: self.uid, |
| 343 | }, |
| 344 | kind, |
| 345 | payload, |
| 346 | }; |
| 347 | cmd.build_standard::<C, P::W>( |
| 348 | mid, |
| 349 | src_ip, |
| 350 | trg_ip, |
| 351 | self.code, |
| 352 | self.schms.clone(), |
| 353 | ) |
| 354 | } |
| 355 | } |