Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/tests/dal.rs

7.4 KiB, 70 runs

created by r1870400018:836, 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 oxedyne_fe2o3_core::{
2 prelude::*,
3 alt::Override,
4 rand::Rand,
5 test::test_it,
6 time::Timer,
7};
8use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
9use oxedyne_fe2o3_hash::{
10 csum::ChecksumScheme,
11 hash::HashScheme,
12};
13use oxedyne_fe2o3_iop_db::api::{
14 Meta,
15 RestSchemesOverride,
16};
17use oxedyne_fe2o3_jdat::{
18 prelude::*,
19 usr::{
20 UsrKind,
21 UsrKindCode,
22 UsrKindId,
23 },
24};
25use oxedyne_fe2o3_net::id::Uid;
26use oxedyne_fe2o3_o3db_sync::{
27 base::{
28 constant,
29 index::ZoneInd,
30 },
31 comm::{
32 response::Wait,
33 },
34 dal::doc::{
35 Doc,
36 DocKey,
37 },
38 data::core::RestSchemesInput,
39 file::zdir::ZoneDir,
40 test::{
41 dbapi,
42 file::{
43 delete_all_index_files,
44 corrupt_an_index_file,
45 },
46 setup,
47 },
48};
49use oxedyne_fe2o3_test::error::delayed_error;
50
51use std::{
52 collections::BTreeMap,
53 path::Path,
54 thread,
55 time::Duration,
56};
57
58
59const wait: Wait = constant::USER_REQUEST_WAIT;
60
61pub fn test_docs(filter: &'static str) -> Outcome<()> {
62
63 res!(test_it(filter, &["Document data abstraction layer 000", "all", "docs"], || {
64
65 res!(std::fs::create_dir_all("./test_db"));
66 let db_root = res!(Path::new("./test_db").canonicalize());
67
68 let mut enckey = [0u8; 32];
69 Rand::fill_u8(&mut enckey);
70 let aes_gcm = res!(EncryptionScheme::new_aes_256_gcm_with_key(&enckey[..]));
71 let _sha3_256 = HashScheme::new_sha3_256();
72 let crc32 = ChecksumScheme::new_crc32();
73 let schms2: RestSchemesOverride<EncryptionScheme, HashScheme> =
74 RestSchemesOverride::default().set_encrypter(Override::Default(aes_gcm.clone()));
75 let schms2 = Some(&schms2);
76 let user = setup::Uid::default();
77 //let meta = Meta::<{ setup::UID_LEN }, setup::Uid>::new(setup::Uid::default());
78 let schms_input = RestSchemesInput::new(
79 Some(aes_gcm.clone()),
80 None::<HashScheme>,
81 None::<HashScheme>,
82 Some(crc32.clone()),
83 );
84 //let user1_name = fmt!("Alice83");
85 //let user1_email = fmt!("alice83@gmail.com");
86 //let user1_pass = fmt!("alice_pass");
87
88 let mut cfg = res!(setup::default_cfg());
89 cfg.cache_size_limit_bytes = 100_000;
90 cfg.rest_chunk_threshold = 700;
91 cfg.num_cbots_per_zone = 2;
92 cfg.num_zones = 3;
93 cfg.zone_overrides = mapdat!{
94 1u16 => mapdat!{
95 "dir" => "../test_db_zone_container",
96 "max_size" => 100u64,
97 },
98 3u16 => mapdat!{
99 "dir" => "",
100 "max_size" => 100u64,
101 },
102 }.get_map().unwrap();
103
104 let error_delay = 2;
105
106 test!(sync_log::stream(), "+---------------------------------------------+");
107 test!(sync_log::stream(), "| NEW OZONE SESSION |");
108 test!(sync_log::stream(), "| Wipe all traces of previous test. |");
109 test!(sync_log::stream(), "| Start database. |");
110 test!(sync_log::stream(), "| Store and retrieve data using the document |");
111 test!(sync_log::stream(), "| Data Abstraction Layer (DAL). |");
112 test!(sync_log::stream(), "| Gracefully shut down the database. |");
113 test!(sync_log::stream(), "+---------------------------------------------+");
114 // Wipe all traces of previous test.
115 // Start database.
116 let mut db = match setup::start_db(
117 db_root.clone(),
118 Some(cfg.clone()),
119 schms_input.clone(),
120 Some(fmt!("./test_db_zone_container")),
121 false,//true,
122 true,
123 ) {
124 // These pauses on errors are needed to capture tardy messages from asynchronous
125 // logging.
126 Err(e) => return Err(delayed_error(e, error_delay)),
127 Ok(db) => db,
128 };
129
130 test!(sync_log::stream(), "Listing files and cache before any user activity...");
131 res!(db.api().list_files(wait));
132 res!(db.api().dump_file_states(wait));
133 res!(db.api().dump_caches(wait));
134
135 let doc = res!(Doc::new_doc("/dir1/dir2",
136 mapdat!{
137 "field3" => "val3",
138 "field4" => "val4",
139 },
140 ));
141 let (key, val) = doc.into_dats();
142
143 let resp = db.responder();
144 res!(db.api().store_dat_using_responder(
145 key,
146 val,
147 user,
148 schms2.clone(),
149 resp.clone(),
150 ));
151 let mut timer = Timer::new();
152 res!(resp.recv_store_ack()); // written, durable and readable
153 debug!(sync_log::stream(), "Storage confirmation delay: {:?}", res!(timer.split_micros()));
154
155 test!(sync_log::stream(), "Listing files and cache after direct insertion...");
156 res!(db.api().list_files(wait));
157 res!(db.api().dump_file_states(wait));
158 res!(db.api().dump_caches(wait));
159
160 thread::sleep(Duration::from_secs(7));
161
162 let doc = res!(Doc::new_doc("/dir1/dir2/dir3",
163 mapdat!{
164 "field1" => "val1",
165 "dir4" => res!(DocKey::new_dir("/dir1/dir2/dir3/dir4")).into_dat(),
166 "field2" => "val2",
167 },
168 ));
169 let (key, val) = doc.clone().into_dats();
170 let mut timer = Timer::new();
171 let resp = res!(db.api().put(key, val, user, None));
172 res!(resp.recv_store_ack()); // written, durable and readable
173 debug!(sync_log::stream(), "Storage confirmation: {:?}", res!(timer.split_micros()));
174
175 test!(sync_log::stream(), "Listing files and cache after insertion via server...");
176 res!(db.api().list_files(wait));
177 res!(db.api().dump_file_states(wait));
178 res!(db.api().dump_caches(wait));
179
180 let (key, val) = doc.clone().into_dats();
181 timer.reset();
182 let resp = res!(db.api().put(key, val, user, None));
183 res!(resp.recv_store_ack()); // written, durable and readable
184 debug!(sync_log::stream(), "Storage confirmation: {:?}", res!(timer.split_micros()));
185
186 test!(sync_log::stream(), "Listing files and cache after repeat insertion via server...");
187 res!(db.api().list_files(wait));
188 res!(db.api().dump_file_states(wait));
189 res!(db.api().dump_caches(wait));
190
191 ////test!(sync_log::stream(), "Pausing now to verify performance of health check.");
192 ////thread::sleep(Duration::from_secs(30));
193
194 ////// Ping all bots.
195 ////let missing = res!(db.api().ping_bots2(wait));
196 ////debug!(sync_log::stream(), "Ping not received from: {:?}.", missing);
197
198 let (key, val) = doc.clone().into_dats();
199 if let Some((dat, _)) = res!(db.api().get_wait(&key, None)) {
200 test!(sync_log::stream(), "Yes! It worked! Daticle is: {:?}", dat);
201 req!(val, dat, "(L: expected, R: actual)");
202 } else {
203 return Err(err!(
204 "Expected to find key {:?}, but found nothing.", key;
205 Data, Missing));
206 }
207
208 //test!(sync_log::stream(), "Listing files...");
209 //res!(db.api().list_files(wait));
210 //res!(db.api().dump_caches(wait));
211 test!(sync_log::stream(), "Shutting db down...");
212 // Gracefully shut down the database.
213 res!(db.shutdown());
214
215 Ok(())
216 }));
217
218
219 Ok(())
220}