Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/test/dbapi.rs

13.8 KiB, 131 runs

created by r1870400018:811, 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::constant,
4 test::{
5 data::{
6 compare_values,
7 stopwatch,
8 },
9 },
10};
11
12use oxedyne_fe2o3_iop_db::api::{
13 Database,
14 RestSchemesOverride,
15};
16use oxedyne_fe2o3_jdat::{
17 prelude::*,
18 chunk::ChunkConfig,
19 id::NumIdDat,
20};
21
22use std::{
23 thread,
24 time::{
25 Duration,
26 Instant,
27 },
28};
29
30pub fn simple<
31 const UIDL: usize,
32 UID: NumIdDat<UIDL> + 'static,
33 ENC: Encrypter + 'static,
34 KH: Hasher + 'static,
35 PR: Hasher + 'static,
36 CS: Checksummer + 'static,
37>(
38 db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
39 schms2: Option<&RestSchemesOverride<ENC, KH>>,
40 user: UID,
41)
42 -> Outcome<()>
43{
44 test!(sync_log::stream(), "Storing and retrieving some simple data.");
45
46 let resp = db.responder();
47 res!(db.api().store_dat_using_responder(
48 dat!("Not at post"),
49 dat!("TK421"),
50 user,
51 schms2,
52 resp.clone(),
53 ));
54 let (exists, n) = res!(resp.recv_store_ack());
55 if n != 1 {
56 return Err(err!("There should only be one chunk."; Test, Size));
57 }
58 if exists {
59 return Err(err!("This key should not exist."; Test, Unexpected));
60 }
61
62 thread::sleep(Duration::from_secs(1));
63 //res!(db.dump_caches(db.default_wait()));
64 //thread::sleep(Duration::from_secs(1));
65
66 // Store it again.
67 let resp = res!(db.api().store_using_schemes(
68 dat!("Not at post"),
69 dat!("TK421"),
70 user,
71 schms2,
72 ));
73 let (exists, n) = res!(resp.recv_store_ack());
74 if n != 1 {
75 return Err(err!("There should only be one chunk."; Test, Size));
76 }
77 if !exists {
78 return Err(err!(
79 "This key should exist.";
80 Test, Data, Missing));
81 }
82
83 // Now retrieve it.
84 let resp = db.responder();
85 res!(db.api().fetch_using_responder(&dat!("Not at post"), schms2, resp.clone()));
86 let expected = dat!("TK421");
87
88 {
89 let enc = db.api().schemes().encrypter();
90 let or_enc = schms2.map(|s| s.encrypter());
91
92 let result = res!(resp.recv_daticle(enc, or_enc));
93 match result {
94 (None, _) => return Err(err!(
95 "This key should exist, instead received {:?}.", result;
96 Test, Unexpected)),
97 (Some((dat, meta2)), _) => {
98 if dat != expected {
99 return Err(err!(
100 "Expected value {:?}, received {:?}.", expected, dat;
101 Test, Unexpected));
102 }
103 if meta2.user != user {
104 return Err(err!(
105 "Expected user {:?}, received {:?}.", user, meta2.user;
106 Test, Unexpected));
107 }
108 },
109 }
110 }
111
112 // Store some more data.
113 let data = mapdat![
114 "mass" => 100,
115 "prot" => 25,
116 "carb" => 30,
117 "fat" => 45,
118 ];
119 res!(db.api().store_blindly(
120 dat!("oats, uncooked"),
121 data,
122 user,
123 schms2,
124 ));
125 // It takes time to store the data. If we attempt to fetch it immediately, we may get nothing.
126 thread::sleep(Duration::from_millis(10)); // winding this down to 0 should raise an error below
127 // An alternative here is to use a db.responder() and wait until storage is complete.
128
129 // Now retrieve it.
130 let resp = res!(db.api().fetch_using_schemes(&dat!("oats, uncooked"), schms2));
131 let enc = db.api().schemes().encrypter();
132 let or_enc = schms2.map(|s| s.encrypter());
133 match res!(resp.recv_daticle(enc, or_enc)) {
134 (None, _) => return Err(err!("Could not find oats data."; Test, Data, Missing)),
135 (Some((Dat::Map(map), meta2)), _) => {
136 match map.get(&dat!("prot")) {
137 Some(Dat::I32(25)) => (),
138 result => return Err(err!(
139 "Unexpected value for field 'prot' in map: {:?}", result;
140 Test, Unexpected)),
141 }
142 if meta2.user != user {
143 return Err(err!(
144 "Expected user {:?}, received {:?}.", user, meta2.user;
145 Test, Unexpected));
146 }
147 },
148 (Some((dat, meta2)), _) => return Err(err!(
149 "Unexpected Dat {:?} and meta {:?} returned.", dat, meta2;
150 Test, Unexpected)),
151 }
152
153 test!(sync_log::stream(), "Wow, it worked");
154 Ok(())
155}
156
157pub fn simple_api<
158 const UIDL: usize,
159 UID: NumIdDat<UIDL> + 'static,
160 ENC: Encrypter + 'static,
161 KH: Hasher + 'static,
162 PR: Hasher + 'static,
163 CS: Checksummer + 'static,
164>(
165 db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
166 user: UID,
167)
168 -> Outcome<()>
169{
170 test!(sync_log::stream(), "Storing and retrieving some simple data using the Database insert and get api.");
171
172 let k = dat!("Meaning of life");
173 let v = dat!(42u8);
174 res!(db.insert(k.clone(), v.clone(), user, None));
175 let result = res!(db.get(&k, None));
176 if let Some((v2, _meta2)) = result {
177 req!(v, v2);
178 } else {
179 return Err(err!("Expected value."; Test, Missing, Data));
180 }
181
182 Ok(())
183}
184
185pub fn store_chunked_data<
186 const UIDL: usize,
187 UID: NumIdDat<UIDL> + 'static,
188 ENC: Encrypter + 'static,
189 KH: Hasher + 'static,
190 PR: Hasher + 'static,
191 CS: Checksummer + 'static,
192>(
193 db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
194 schms2: Option<&RestSchemesOverride<ENC, KH>>,
195 user: UID,
196 k: Dat,
197 v: Dat,
198)
199 -> Outcome<()>
200{
201 const CHNK_CFG: ChunkConfig = ChunkConfig {
202 threshold_bytes: 0,
203 chunk_size: 123,
204 dat_wrap: true,
205 pad_last: true,
206 };
207
208 test!(sync_log::stream(), "Store and read data that is chunked.");
209
210 let schms2 = match schms2 {
211 Some(schms2) => schms2.clone(),
212 None => RestSchemesOverride::<ENC, KH>::default(),
213 }.set_chunk_config(Some(CHNK_CFG));
214
215 let (resp, num_chunks) = res!(db.api().store_chunked(
216 k,
217 v,
218 user,
219 Some(&schms2),
220 ));
221 // The record count comes first, and every record is then answered. Counting the count as
222 // one of the answers, as this once did, returned before the last record was acknowledged.
223 let (exists, n) = res!(resp.recv_store_ack());
224 if n != num_chunks {
225 return Err(err!(
226 "The store reported {} records, the caller was told {}.", n, num_chunks;
227 Test, Mismatch));
228 }
229 if exists {
230 return Err(err!("This key should not exist."; Test, Unexpected));
231 }
232 Ok(())
233}
234
235pub fn fetch_chunked_data<
236 const UIDL: usize,
237 UID: NumIdDat<UIDL> + 'static,
238 ENC: Encrypter + 'static,
239 KH: Hasher + 'static,
240 PR: Hasher + 'static,
241 CS: Checksummer + 'static,
242>(
243 db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
244 key: &Dat,
245 val: &Vec<u8>,
246 user: UID,
247 schms2: Option<&RestSchemesOverride<ENC, KH>>,
248)
249 -> Outcome<()>
250{
251 // Now retrieve it, twice.
252 // 1. Retrieve values from caches,
253 // 2. Retrieve values from files.
254 for src in &["cache", "files"] {
255 test!(sync_log::stream(), "Retrieve values from {}.", src);
256 // First, fetch the bunch key.
257 let resp = res!(db.api().fetch_using_schemes(key, schms2));
258 let enc = db.api().schemes().encrypter();
259 let or_enc = schms2.map(|s| s.encrypter());
260 match res!(resp.recv_daticle(enc, or_enc)) {
261 (None, _) => return Err(err!(
262 "Could not find data for key {:?} {:02x?}.", key,key.as_bytes();
263 Test, Data, Missing)),
264 (Some((Dat::Tup5u64(tup), meta2)), _) => {
265 // Quick check on the metadata
266 if meta2.user != user {
267 return Err(err!(
268 "Expected user {:?}, received {:?}.", user, meta2.user;
269 Unexpected, Data, Mismatch));
270 }
271 // Fetch the chunks.
272 match res!(db.api().fetch_chunks(&Dat::Tup5u64(tup), schms2)) {
273 Dat::BU8(v) |
274 Dat::BU16(v) |
275 Dat::BU32(v) |
276 Dat::BU64(v) => {
277 if val.len() != v.len() {
278 return Err(err!(
279 "Original length = {}, retrieved length = {}.",
280 val.len(), v.len();
281 Test, Size, Mismatch));
282 }
283 for j in 0..val.len() {
284 if val[j] != v[j] {
285 return Err(err!(
286 "Failed at byte {} of {}. Original value {}, \
287 retrieved value {}.",
288 j, val.len(), val[j], v[j];
289 Test, Data, Mismatch));
290 }
291 }
292 },
293 dat => return Err(err!(
294 "Unexpected Dat {:?} returned.", dat;
295 Unexpected, Data)),
296 }
297 },
298 (Some((Dat::BU8(v), _)), _) |
299 (Some((Dat::BU16(v), _)), _) |
300 (Some((Dat::BU32(v), _)), _) |
301 (Some((Dat::BU64(v), _)), _) => {
302 // Should only get here if the data was not chunked.
303 if val.len() != v.len() {
304 return Err(err!(
305 "Original length = {}, retrieved length = {}.",
306 val.len(), v.len();
307 Test, Size, Mismatch));
308 }
309 },
310 (Some((dat, meta2)), _) => return Err(err!(
311 "Unexpected Dat {:?} and meta {:?} returned.", dat, meta2;
312 Test, Data, Unexpected)),
313 }
314 res!(db.api().clear_cache_values(constant::USER_REQUEST_WAIT));
315 }
316
317 test!(sync_log::stream(), "Wow, it worked");
318
319 Ok(())
320}
321
322pub fn store<
323 const UIDL: usize,
324 UID: NumIdDat<UIDL> + 'static,
325 ENC: Encrypter + 'static,
326 KH: Hasher + 'static,
327 PR: Hasher + 'static,
328 CS: Checksummer + 'static,
329>(
330 db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
331 user: UID,
332 schms2: Option<&RestSchemesOverride<ENC, KH>>,
333 ks: Vec<Dat>,
334 vs: Vec<Dat>,
335 byts: usize, // total bytes in data
336)
337 -> Outcome<(u64, u64)>
338{
339 test!(sync_log::stream(), "Starting some new live files to give us a clean slate...");
340 test!(sync_log::stream(), " RestSchemes: {:?}", schms2);
341 res!(db.api().new_live_files());
342
343 let n = ks.len();
344
345 // Almost all the data is stored in "firehose mode", where we don't ask for any database
346 // feedback on the result.
347 test!(sync_log::stream(), "Storing {} key-value data pairs.", n);
348 let start = Instant::now();
349 for i in 0..(n - 1) {
350 res!(db.api().store_blindly(
351 ks[i].clone(),
352 vs[i].clone(),
353 user,
354 schms2,
355 ));
356 }
357 // Store one final pair in order to get a responder to tell us when all writes have completed.
358 res!(res!(db.api().store_using_schemes(
359 ks[n - 1].clone(),
360 vs[n - 1].clone(),
361 user,
362 schms2,
363 )).recv_store_ack());
364
365 let elapsed = start.elapsed().as_secs_f64();
366 Ok(stopwatch(elapsed, n, byts))
367}
368
369pub fn fetch<
370 const UIDL: usize,
371 UID: NumIdDat<UIDL> + 'static,
372 ENC: Encrypter + 'static,
373 KH: Hasher + 'static,
374 PR: Hasher + 'static,
375 CS: Checksummer + 'static,
376>(
377 db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
378 schms2: Option<&RestSchemesOverride<ENC, KH>>,
379 ks: &Vec<Dat>,
380 mask: &Vec<bool>,
381 vs: &Vec<Dat>,
382 byts: usize,
383)
384 -> Outcome<Vec<(String, u64, u64)>>
385{
386 test!(sync_log::stream(), "Fetching data.");
387 // Now retrieve it, twice.
388 // 1. Retrieve values from caches,
389 // 2. Retrieve values from files.
390 let n = ks.len();
391 if mask.len() != n {
392 return Err(err!(
393 "Mask length {} does not match number of keys {}",
394 mask.len(), n;
395 Mismatch));
396 }
397
398 let mut result = Vec::new();
399 for src in &["cache", "files"] {
400 test!(sync_log::stream(), "Retrieve values from {}.", src);
401 let mut err_cnt: usize = 0;
402 let start = Instant::now();
403 for i in 0..n {
404 if mask[i] {
405 //debug!(sync_log::stream(), "Attempt to fetch key {}.", i);
406 let resp = res!(db.api().fetch_using_schemes(&ks[i], schms2));
407 let enc = db.api().schemes().encrypter();
408 let or_enc = schms2.map(|s| s.encrypter());
409 match res!(resp.recv_daticle(enc, or_enc)) {
410 (None, _) => err_cnt += 1,
411 (Some((Dat::Tup5u64(tup), _)), _) => {
412 // Fetch the chunks.
413 let dat = res!(db.api().fetch_chunks(&Dat::Tup5u64(tup.clone()), schms2));
414 let v1 = try_extract_dat!(dat, BU8, BU16, BU32, BU64);
415 res!(compare_values(i, &v1, &vs[i]));
416 },
417 (Some((dat, _)), _) => {
418 let v1 = try_extract_dat!(dat, BU8, BU16, BU32, BU64);
419 res!(compare_values(i, &v1, &vs[i]));
420 },
421 }
422 }
423 }
424 if err_cnt > 0 {
425 error!(sync_log::stream(), err!("Could not find data for {} items.", err_cnt; Test, Data, Missing));
426 } else {
427 test!(sync_log::stream(), "All data retrieved successfully.");
428 }
429 let elapsed = start.elapsed().as_secs_f64();
430 let (tps, bw) = stopwatch(elapsed, n, byts);
431 result.push((src.to_uppercase().to_string(), tps, bw));
432 if *src == "cache" {
433 res!(db.api().clear_cache_values(constant::USER_REQUEST_WAIT));
434 }
435 }
436
437 //res!(db.dump_file_states());
438
439 Ok(result)
440}