Oregami
Repositories/oxedyne/fe2o3

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

1use 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
15use oxedyne_fe2o3_iop_db::api::{
16 RestSchemesOverride,
17};
18use oxedyne_fe2o3_jdat::{
19 prelude::*,
20 id::NumIdDat,
21};
22//use oxedyne_fe2o3_hash::{
23// csum::{
24// ChecksummerDefAlt,
25// ChecksumScheme,
26// },
27//};
28
29use std::time::Instant;
30
31pub 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.
59pub 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
82pub 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
99pub 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
171pub 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