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
| 1 | use crate::srv::{ |
| 2 | msg::packet::{ |
| 3 | PacketCount, |
| 4 | PacketMeta, |
| 5 | }, |
| 6 | }; |
| 7 | |
| 8 | use oxedyne_fe2o3_core::{ |
| 9 | prelude::*, |
| 10 | map::MapMut, |
| 11 | }; |
| 12 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 13 | use oxedyne_fe2o3_hash::map::ShardMap; |
| 14 | use oxedyne_fe2o3_iop_hash::api::{ |
| 15 | Hasher, |
| 16 | HashForm, |
| 17 | }; |
| 18 | |
| 19 | use 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)] |
| 32 | pub 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 | |
| 42 | impl< |
| 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)] |
| 152 | pub 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)] |
| 160 | pub 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 | |
| 170 | impl 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 | |
| 184 | impl 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 | } |