Oregami
Repositories/oxedyne/fe2o3

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

6.9 KiB, 26 runs

created by r1870400018:876, 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::srv::{
2 msg::packet::{
3 PacketCount,
4 PacketMeta,
5 },
6};
7
8use oxedyne_fe2o3_core::{
9 prelude::*,
10 map::MapMut,
11};
12use oxedyne_fe2o3_jdat::id::NumIdDat;
13use oxedyne_fe2o3_hash::map::ShardMap;
14use oxedyne_fe2o3_iop_hash::api::{
15 Hasher,
16 HashForm,
17};
18
19use std::{
20 clone::Clone,
21 collections::BTreeMap,
22 fmt::Debug,
23 sync::RwLock,
24 time::{
25 Duration,
26 Instant,
27 },
28};
29
30
31#[derive(Debug)]
32pub struct MsgAssembler<
33 // ShardMap
34 const C: usize, // Capacity (maximum number of shards).
35 M: MapMut<HashForm, MsgState> + Clone + Debug,
36 H: Hasher + Send + Sync + 'static, // Key hasher.
37 const S: usize, // Key hasher salt length.
38> {
39 pub msgs: ShardMap<C, S, MsgState, M, H>,
40}
41
42impl<
43 // ShardMap
44 const C: usize, // Capacity (maximum number of shards).
45 M: MapMut<HashForm, MsgState> + Clone + Debug,
46 H: Hasher + Send + Sync + 'static, // Key hasher.
47 const S: usize, // Key hasher salt length.
48>
49 MsgAssembler<C, M, H, S>
50{
51 pub fn new(
52 n: u32,
53 salt: [u8; S],
54 init_map: M,
55 hasher: H,
56 )
57 -> Outcome<Self>
58 {
59 Ok(Self {
60 msgs: res!(ShardMap::new(
61 n,
62 salt,
63 init_map,
64 hasher,
65 )),
66 })
67 }
68
69 pub fn get_locked_map<
70 const MIDL: usize,
71 MID: NumIdDat<MIDL>,
72 >(
73 &self,
74 mid: &MID,
75 )
76 -> Outcome<(HashForm, &RwLock<M>)>
77 {
78 let key = self.msgs.key(&mid.to_byte_array());
79 let locked_map = res!(self.msgs.get_shard_using_hash(&key));
80 Ok((key, locked_map))
81 }
82
83 /// Insert the packet message chunk into the message assembler map. Returns whether the
84 /// message should be dropped, and possibly the entire message when it is complete.
85 pub fn get_msg<
86 const MIDL: usize,
87 const UIDL: usize,
88 MID: NumIdDat<MIDL>,
89 UID: NumIdDat<UIDL>,
90 >(
91 &self,
92 meta: &PacketMeta<MIDL, UIDL, MID, UID>,
93 buf: &[u8],
94 params: &MsgAssemblyParams,
95 )
96 -> Outcome<(bool, Option<Vec<u8>>)>
97 {
98 let (key, locked_map) = res!(self.get_locked_map(&meta.mid));
99 let mut unlocked_map = lock_write!(locked_map);
100 if !unlocked_map.contains_key(&key) {
101 unlocked_map.insert(key.clone(), MsgState::new(meta.chnk.num_chunks));
102 }
103 let (drop, msg_byt_opt) = match unlocked_map.get_mut(&key) {
104 Some(mstat) => mstat.insert_part(
105 meta,
106 buf,
107 params,
108 ),
109 None => return Ok((true, None)),
110 };
111 if drop || msg_byt_opt.is_some() {
112 unlocked_map.remove(&key);
113 }
114 Ok((drop, msg_byt_opt))
115 }
116
117 pub fn remove<
118 const MIDL: usize,
119 MID: NumIdDat<MIDL>,
120 >(
121 &self,
122 mid: &MID,
123 )
124 -> Outcome<()>
125 {
126 let (key, locked_map) = res!(self.get_locked_map(mid));
127 let mut unlocked_map = lock_write!(locked_map);
128 unlocked_map.remove(&key);
129 Ok(())
130 }
131
132 pub fn message_assembly_garbage_collection(
133 &self,
134 params: &MsgAssemblyParams,
135 )
136 -> Outcome<()>
137 {
138 for i in 0..self.msgs.n {
139 if let Some(locked_map) = &self.msgs.shards[i] {
140 let mut unlocked_map = lock_write!(locked_map);
141 unlocked_map.retain(
142 |_key, mstat|
143 !mstat.drop_on_time_check(params)
144 )
145 }
146 }
147 Ok(())
148 }
149}
150
151#[derive(Clone, Debug, Default)]
152pub struct MsgAssemblyParams {
153 pub msg_sunset: Duration,
154 pub idle_max: Duration,
155 pub rep_tot_lim: u8,
156 pub rep_max_lim: u8,
157}
158
159#[derive(Clone, Debug)]
160pub struct MsgState {
161 parts: BTreeMap<PacketCount, (Vec<u8>, u8)>, // The packet and how many times it has been written.
162 tot: PacketCount,
163 cnt: PacketCount, // Count of packets not yet received.
164 first: Instant,
165 last: Instant,
166 rep_tot: u8, // Total repetitions.
167 rep_max: u8, // Maximum repetition for any packet number.
168}
169
170impl Default for MsgState {
171 fn default() -> Self {
172 Self {
173 parts: BTreeMap::new(),
174 tot: 0,
175 cnt: 0,
176 first: Instant::now(),
177 last: Instant::now(),
178 rep_tot: 0,
179 rep_max: 0,
180 }
181 }
182}
183
184impl MsgState {
185
186 pub fn new(total_packets: PacketCount) -> Self {
187 Self {
188 tot: total_packets,
189 cnt: total_packets,
190 ..Default::default()
191 }
192 }
193
194 /// Inserts the packet payload into the message. Returns whether the entire partial message
195 /// should be dropped, and possibly the completed message.
196 pub fn insert_part<
197 const MIDL: usize,
198 const UIDL: usize,
199 MID: NumIdDat<MIDL>,
200 UID: NumIdDat<UIDL>,
201 >(
202 &mut self,
203 meta: &PacketMeta<MIDL, UIDL, MID, UID>,
204 buf: &[u8],
205 params: &MsgAssemblyParams,
206 )
207 -> (bool, Option<Vec<u8>>)
208 {
209 if self.cnt == self.tot {
210 self.first = Instant::now();
211 }
212 if self.drop_on_time_check(params) {
213 return (true, None);
214 }
215 self.last = Instant::now();
216 match self.parts.get_mut(&meta.chnk.index) {
217 Some((_part, n)) => { // Update repetition data but do not copy packet.
218 match n.checked_add(1) {
219 Some(n2) => if n2 > params.rep_max_lim {
220 return (true, None);
221 } else {
222 *n = n2;
223 self.rep_max = n2;
224 }
225 None => return (true, None),
226 }
227 match self.rep_tot.checked_add(1) {
228 Some(n2) => if n2 > params.rep_tot_lim {
229 return (true, None);
230 } else {
231 self.rep_tot = n2;
232 },
233 None => return (true, None),
234 }
235 return (false, None);
236 },
237 None => (),
238 }
239 self.parts.insert(meta.chnk.index, (buf.to_vec(), 0));
240 if self.cnt == 1 {
241 // Assemble full message bytes.
242 let mut v = Vec::new();
243 for (_id, (part, _n)) in self.parts.iter_mut() {
244 v.append(part);
245 }
246 return (false, Some(v));
247 } else {
248 self.cnt -= 1;
249 }
250 (false, None)
251 }
252
253 pub fn drop_on_time_check(
254 &mut self,
255 params: &MsgAssemblyParams,
256 )
257 -> bool
258 {
259 if self.first.elapsed() > params.msg_sunset ||
260 self.last.elapsed() > params.idle_max
261 {
262 true
263 } else {
264 false
265 }
266 }
267}