oxedyne/fe2o3/fe2o3_o3db_sync/src/test/data.rs
6.2 KiB, 56 runs
created by r1870400018:809, 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 | prelude::*, |
| 3 | base::index::ZoneInd, |
| 4 | comm::{ |
| 5 | msg::OzoneMsg, |
| 6 | response::Responder, |
| 7 | }, |
| 8 | data::{ |
| 9 | core::{ |
| 10 | Encode, |
| 11 | }, |
| 12 | }, |
| 13 | }; |
| 14 | |
| 15 | use oxedyne_fe2o3_iop_db::api::{ |
| 16 | RestSchemesOverride, |
| 17 | }; |
| 18 | use oxedyne_fe2o3_jdat::{ |
| 19 | prelude::*, |
| 20 | id::NumIdDat, |
| 21 | }; |
| 22 | //use oxedyne_fe2o3_hash::{ |
| 23 | // csum::{ |
| 24 | // ChecksummerDefAlt, |
| 25 | // ChecksumScheme, |
| 26 | // }, |
| 27 | //}; |
| 28 | |
| 29 | use std::time::Instant; |
| 30 | |
| 31 | pub fn compare_values(i: usize, v1: &Vec<u8>, vorig: &Dat) -> Outcome<()> { |
| 32 | // Get the associated data since we know vorig comes with a Dat byte wrapper |
| 33 | if let Some(v2) = vorig.bytes_ref() { |
| 34 | if v1.len() != v2.len() { |
| 35 | return Err(err!( |
| 36 | "Data {}: Original length = {}, retrieved length = {}.", |
| 37 | i, v2.len(), v1.len(); |
| 38 | Test, Size, Mismatch)); |
| 39 | } |
| 40 | for j in 0..v1.len() { |
| 41 | if v1[j] != v2[j] { |
| 42 | debug!(sync_log::stream(), "i = {}",i); |
| 43 | debug!(sync_log::stream(), "expected = {:02x?}",v2); |
| 44 | debug!(sync_log::stream(), "got = {:02x?}",v1); |
| 45 | return Err(err!( |
| 46 | "Data {}: Differs at position {}.", i, j; |
| 47 | Data, Mismatch)); |
| 48 | } |
| 49 | } |
| 50 | } else { |
| 51 | return Err(err!( |
| 52 | "Given test value with index {} is not a Dat containing bytes.", i; |
| 53 | Invalid, Input)); |
| 54 | } |
| 55 | Ok(()) |
| 56 | } |
| 57 | |
| 58 | /// Identifies sequences that are unique, starting from the last element. |
| 59 | pub fn find_unique(v: &Vec<Vec<u8>>) -> Vec<bool> { |
| 60 | let mut b = vec![true; v.len()]; |
| 61 | if v.len() > 1 { |
| 62 | for i in (1..v.len()).rev() { |
| 63 | for j in 0..i { |
| 64 | if v[i].len() == v[j].len() { |
| 65 | let mut equal = true; |
| 66 | for k in 0..v[i].len() { |
| 67 | if v[i][k] != v[j][k] { |
| 68 | equal = false; |
| 69 | break; |
| 70 | } |
| 71 | } |
| 72 | if equal { |
| 73 | b[j] = false; |
| 74 | } |
| 75 | } |
| 76 | } |
| 77 | } |
| 78 | } |
| 79 | b |
| 80 | } |
| 81 | |
| 82 | pub fn stopwatch( |
| 83 | elapsed: f64, |
| 84 | n: usize, |
| 85 | byts: usize, // total bytes in data |
| 86 | ) |
| 87 | -> (u64, u64) |
| 88 | { |
| 89 | let tps = (n as f64) / elapsed; |
| 90 | let bw = ((byts as f64) / elapsed) * (8.0 / 1_000_000.0); |
| 91 | test!(sync_log::stream(), "Performance metrics:"); |
| 92 | test!(sync_log::stream(), " Elapsed: {:10.4} [s]", elapsed); |
| 93 | test!(sync_log::stream(), " TPS: {:10.2}", tps); |
| 94 | test!(sync_log::stream(), " Bandwidth: {:10.3} [Mb]/[s]", bw); |
| 95 | |
| 96 | (tps as u64, bw as u64) |
| 97 | } |
| 98 | |
| 99 | pub fn encode_daticles( |
| 100 | ks: Vec<Dat>, |
| 101 | vs: Vec<Dat>, |
| 102 | byts: usize, // total bytes in data |
| 103 | ) |
| 104 | -> Outcome<(Vec<Vec<u8>>, Vec<Vec<u8>>, u64, u64)> |
| 105 | { |
| 106 | let n = ks.len(); |
| 107 | |
| 108 | test!(sync_log::stream(), "Start daticle encoding timing run..."); |
| 109 | let start = Instant::now(); |
| 110 | for i in 0..n { |
| 111 | let (_k, _v) = res!(Encode::encode_dat(ks[i].clone(), vs[i].clone())); |
| 112 | } |
| 113 | let elapsed = start.elapsed().as_secs_f64(); |
| 114 | let (tps, bw) = stopwatch(elapsed, n, byts); |
| 115 | |
| 116 | test!(sync_log::stream(), "Redoing encoding to collect bytes..."); |
| 117 | let mut kbyts = Vec::new(); |
| 118 | let mut vbyts = Vec::new(); |
| 119 | for i in 0..n { |
| 120 | let (k, v) = res!(Encode::encode_dat(ks[i].clone(), vs[i].clone())); |
| 121 | kbyts.push(k); |
| 122 | vbyts.push(v); |
| 123 | } |
| 124 | test!(sync_log::stream(), " Finished."); |
| 125 | |
| 126 | Ok((kbyts, vbyts, tps, bw)) |
| 127 | } |
| 128 | |
| 129 | //pub fn encode_write_messages< |
| 130 | // const UIDL: usize, |
| 131 | // UID: NumIdDat<UIDL> + 'static, |
| 132 | // C: Checksummer + 'static, |
| 133 | //>( |
| 134 | // meta: &Meta<UIDL, UID>, |
| 135 | // ks: Vec<Vec<u8>>, |
| 136 | // vs: Vec<Vec<u8>>, |
| 137 | // byts: usize, // total bytes in data |
| 138 | // csummer: ChecksummerDefAlt<ChecksumScheme, C>, |
| 139 | //) |
| 140 | // -> Outcome<(u64, u64)> |
| 141 | //{ |
| 142 | // let n = ks.len(); |
| 143 | // |
| 144 | // test!(sync_log::stream(), "Assemble wrapped KeyVals..."); |
| 145 | // let mut keyvals = Vec::new(); |
| 146 | // for i in 0..n { |
| 147 | // let mut chash = [0u8; 4]; |
| 148 | // for j in 0..4 { |
| 149 | // chash[j] = ks[i][j]; |
| 150 | // } |
| 151 | // let kv = KeyVal { |
| 152 | // key: Key::Complete(ks[i].clone()), |
| 153 | // val: vs[i].clone(), |
| 154 | // chash, // dummy |
| 155 | // meta: meta.clone(), |
| 156 | // cbpind: 3, // dummy |
| 157 | // }; |
| 158 | // keyvals.push(kv); |
| 159 | // } |
| 160 | // test!(sync_log::stream(), "Start write message encoding timing run..."); |
| 161 | // let start = Instant::now(); |
| 162 | // for kv in keyvals { |
| 163 | // res!(Encode::encode(kv, csummer.clone())); |
| 164 | // } |
| 165 | // let elapsed = start.elapsed().as_secs_f64(); |
| 166 | // let (tps, bw) = stopwatch(elapsed, n, byts); |
| 167 | // |
| 168 | // Ok((tps, bw)) |
| 169 | //} |
| 170 | |
| 171 | pub fn prepare_write_messages< |
| 172 | const UIDL: usize, |
| 173 | UID: NumIdDat<UIDL> + 'static, |
| 174 | ENC: Encrypter + 'static, |
| 175 | KH: Hasher + 'static, |
| 176 | PR: Hasher + 'static, |
| 177 | CS: Checksummer + 'static, |
| 178 | >( |
| 179 | db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>, |
| 180 | user: UID, |
| 181 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 182 | ks: Vec<Dat>, |
| 183 | vs: Vec<Dat>, |
| 184 | byts: usize, // total bytes in data |
| 185 | ) |
| 186 | -> Outcome<(usize, Vec<(OzoneMsg<UIDL, UID, ENC, KH>, ZoneInd)>)> |
| 187 | { |
| 188 | let n = ks.len(); |
| 189 | |
| 190 | test!(sync_log::stream(), "Encoding data for writing..."); |
| 191 | let mut msgs = Vec::new(); |
| 192 | let resp = Responder::none(None); |
| 193 | let start = Instant::now(); |
| 194 | for i in 0..n { |
| 195 | let msgs2 = res!(db.api().prepare_write_dat( |
| 196 | ks[i].clone(), |
| 197 | vs[i].clone(), |
| 198 | user, |
| 199 | schms2, |
| 200 | resp.clone(), |
| 201 | None, |
| 202 | )); |
| 203 | for msg in msgs2 { |
| 204 | msgs.push(msg); |
| 205 | } |
| 206 | } |
| 207 | let elapsed = start.elapsed().as_secs_f64(); |
| 208 | test!(sync_log::stream(), " n = {}, msgs.len = {}", n, msgs.len()); |
| 209 | let tps = (n as f64) / elapsed; |
| 210 | let bw = ((byts as f64) / elapsed) * (8.0 / 1_000_000.0); |
| 211 | test!(sync_log::stream(), "Write preparation performance metrics:"); |
| 212 | test!(sync_log::stream(), " Elapsed: {:10.4} [s]", elapsed); |
| 213 | test!(sync_log::stream(), " TPS: {:10.2}", tps); |
| 214 | test!(sync_log::stream(), " Bandwidth: {:10.3} [Mb]/[s]", bw); |
| 215 | |
| 216 | Ok((n, msgs)) |
| 217 | } |
| 218 |