Oregami
Repositories/oxedyne/fe2o3

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

1use 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
38use oxedyne_fe2o3_core::{
39 prelude::*,
40 byte::{
41 Encoding,
42 IntoBytes,
43 ToBytes,
44 },
45};
46use oxedyne_fe2o3_hash::{
47 pow::{
48 PowCreateParams,
49 PowVars,
50 Pristine,
51 ProofOfWork,
52 ZeroBits,
53 },
54};
55use oxedyne_fe2o3_syntax::{
56 SyntaxRef,
57 msg::{
58 Msg,
59 MsgCmd,
60 },
61};
62
63use 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.
87pub 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
313impl<
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}