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
| 1 | use crate::{ |
| 2 | prelude::*, |
| 3 | base::constant, |
| 4 | test::{ |
| 5 | data::{ |
| 6 | compare_values, |
| 7 | stopwatch, |
| 8 | }, |
| 9 | }, |
| 10 | }; |
| 11 | |
| 12 | use oxedyne_fe2o3_iop_db::api::{ |
| 13 | Database, |
| 14 | RestSchemesOverride, |
| 15 | }; |
| 16 | use oxedyne_fe2o3_jdat::{ |
| 17 | prelude::*, |
| 18 | chunk::ChunkConfig, |
| 19 | id::NumIdDat, |
| 20 | }; |
| 21 | |
| 22 | use std::{ |
| 23 | thread, |
| 24 | time::{ |
| 25 | Duration, |
| 26 | Instant, |
| 27 | }, |
| 28 | }; |
| 29 | |
| 30 | pub 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 | |
| 157 | pub 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 | |
| 185 | pub 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 | |
| 235 | pub 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 | |
| 322 | pub 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 | |
| 369 | pub 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 | } |