Oregami
Repositories/oxedyne/fe2o3

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

9.0 KiB, 42 runs

created by r1870400018:813, 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::msg::OzoneMsg,
5 file::{
6 core::FileAccess,
7 floc::{
8 FileNum,
9 StoredFileLocation,
10 },
11 zdir::ZoneDir,
12 },
13 test::{
14 data::{
15 prepare_write_messages,
16 stopwatch,
17 },
18 },
19};
20
21use oxedyne_fe2o3_iop_db::api::{
22 RestSchemesOverride,
23};
24use oxedyne_fe2o3_jdat::{
25 prelude::*,
26 id::NumIdDat,
27};
28use oxedyne_fe2o3_hash::csum::ChecksumScheme;
29
30use std::{
31 collections::BTreeMap,
32 fs::{
33 File,
34 OpenOptions,
35 },
36 io::{
37 Read,
38 Write,
39 },
40 path::PathBuf,
41 time::{
42 self,
43 Instant,
44 },
45};
46
47pub use humantime::format_rfc3339_seconds as timefmt;
48use hostname;
49
50pub fn delete_all_index_files(zone_dirs: &BTreeMap<ZoneInd, ZoneDir>) -> Outcome<()> {
51 for (_, zdir) in zone_dirs {
52 let dir = &zdir.dir;
53 test!(sync_log::stream(), "Removing all index files from {:?}", dir);
54 for entry in res!(std::fs::read_dir(dir)) {
55 let entry = res!(entry);
56 let path = entry.path();
57 if path.is_file() {
58 if let Some(os_str) = path.extension() {
59 if let Some("ind") = os_str.to_str() {
60 res!(std::fs::remove_file(path));
61 }
62 }
63 }
64 }
65 }
66 Ok(())
67}
68
69pub fn corrupt_an_index_file(zone_dirs: &BTreeMap<ZoneInd, ZoneDir>) -> Outcome<()> {
70 test!(sync_log::stream(), "Deliberately corrupting the first index file in the first zone.");
71 if let Some(zdir) = zone_dirs.get(&ZoneInd::new(0usize)) {
72 for entry in res!(std::fs::read_dir(&zdir.dir)) {
73 let entry = res!(entry);
74 let path = entry.path();
75 if path.is_file() {
76 let (fnum, typ) = res!(ZoneDir::ozone_file_number_and_type(&path));
77 if fnum == 1 && typ == FileType::Index {
78 let mut buf = Vec::new();
79 {
80 let mut file = res!(std::fs::OpenOptions::new()
81 .read(true)
82 .open(&path)
83 );
84 res!(file.read_to_end(&mut buf));
85 // Pick a single byte in the middle and negate it.
86 let i = buf.len()/2;
87 buf[i] = !buf[i];
88 }
89 let mut file = res!(std::fs::OpenOptions::new()
90 .write(true)
91 .open(&path)
92 );
93 res!(file.write_all(&buf));
94 }
95 }
96 }
97 } else {
98 return Err(err!(
99 "Could not obtain the directory for the first zone.";
100 Missing, Data));
101 }
102
103 Ok(())
104}
105
106pub fn save_single_file<
107 const UIDL: usize,
108 UID: NumIdDat<UIDL> + 'static,
109 ENC: Encrypter + 'static,
110 KH: Hasher + 'static,
111 PR: Hasher + 'static,
112 CS: Checksummer + 'static,
113>(
114 dir: PathBuf,
115 mut db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
116 user: UID,
117 schms2: Option<&RestSchemesOverride<ENC, KH>>,
118 ks: Vec<Dat>,
119 vs: Vec<Dat>,
120 byts: usize, // total bytes in data
121)
122 -> Outcome<(u64, u64)>
123{
124 test!(sync_log::stream(), "Performing control experiment by saving data directly to a single disk file.");
125 test!(sync_log::stream(), " RestSchemes: {:?}", schms2);
126
127 let mut path = dir.clone();
128 path.push("control_file");
129
130 let mut file = res!(OpenOptions::new()
131 .create(true)
132 .read(true)
133 .write(true)
134 .append(true)
135 .open(path)
136 );
137
138 let (n, msgs) = res!(prepare_write_messages(
139 &mut db,
140 user,
141 schms2,
142 ks,
143 vs,
144 byts,
145 ));
146
147 let mut datsize: usize = 0;
148 test!(sync_log::stream(), "Saving {} key-value data pairs.", n);
149 let start = Instant::now();
150 for (msg, _zind) in msgs {
151 match msg {
152 OzoneMsg::Write{kstored, vstored, ..} => {
153 datsize += res!(file.write(&kstored));
154 datsize += res!(file.write(&vstored));
155
156 res!(file.write(&kstored));
157 let sfloc = res!(StoredFileLocation::new(
158 0,
159 datsize as u64,
160 kstored.len() as u64,
161 vstored.len() as u64,
162 ChecksumScheme::new_crc32(),
163 ));
164 let istored = sfloc.buf;
165 res!(file.write(&istored));
166 },
167 _ => return Err(err!(
168 "Expecting OzoneMsg::Write, got {:?}", msg;
169 Unreachable, Bug)),
170 }
171 }
172
173 let elapsed = start.elapsed().as_secs_f64();
174 Ok(stopwatch(elapsed, n, byts))
175}
176
177pub fn save_multiple_files<
178 const UIDL: usize,
179 UID: NumIdDat<UIDL> + 'static,
180 ENC: Encrypter + 'static,
181 KH: Hasher + 'static,
182 PR: Hasher + 'static,
183 CS: Checksummer + 'static,
184>(
185 dir: PathBuf,
186 size_lim: usize,
187 mut db: &mut O3db<UIDL, UID, ENC, KH, PR, CS>,
188 user: UID,
189 schms2: Option<&RestSchemesOverride<ENC, KH>>,
190 ks: Vec<Dat>,
191 vs: Vec<Dat>,
192 byts: usize, // total bytes in data
193)
194 -> Outcome<(u64, u64)>
195{
196 test!(sync_log::stream(), "Performing control experiment by saving data directly to multiple disk files.");
197 test!(sync_log::stream(), " RestSchemes: {:?}", schms2);
198
199 let mut fnum: FileNum = 0;
200 let mut datsize: usize = 0;
201 let mut datfile: Option<File> = None;
202 let mut indfile: Option<File> = None;
203
204 let (n, msgs) = res!(prepare_write_messages(
205 &mut db,
206 user,
207 schms2,
208 ks,
209 vs,
210 byts,
211 ));
212
213 let mut i = 0;
214 test!(sync_log::stream(), "Saving {} key-value data pairs.", n);
215 let start = Instant::now();
216 for (msg, _zind) in msgs {
217 match msg {
218 OzoneMsg::Write{kstored, vstored, ..} => {
219 if i == 0 || datsize + kstored.len() + vstored.len() > size_lim {
220 fnum += 1;
221 datsize = 0;
222 let mut path = dir.clone();
223 path.push(ZoneDir::relative_file_path(&FileType::Data, fnum));
224 test!(sync_log::stream(), "Creating file {:?}", path);
225 datfile = Some(res!(ZoneDir::open_file(&path, &FileAccess::Writing)));
226 let mut path = dir.clone();
227 path.push(ZoneDir::relative_file_path(&FileType::Index, fnum));
228 test!(sync_log::stream(), "Creating file {:?}", path);
229 indfile = Some(res!(ZoneDir::open_file(&path, &FileAccess::Writing)));
230 }
231 if let Some(file) = &mut datfile {
232 datsize += res!(file.write(&kstored));
233 datsize += res!(file.write(&vstored));
234 }
235 if let Some(file) = &mut indfile {
236 res!(file.write(&kstored));
237 let sfloc = res!(StoredFileLocation::new(
238 fnum,
239 datsize as u64,
240 kstored.len() as u64,
241 vstored.len() as u64,
242 ChecksumScheme::new_crc32(),
243 ));
244 let istored = sfloc.buf;
245 res!(file.write(&istored));
246 }
247 //if i % 100 == 0 {
248 // debug!(sync_log::stream(), "Key {} data written {}", i, datsize);
249 //}
250 i += 1;
251 },
252 _ => return Err(err!(
253 "Expecting OzoneMsg::Write, got {:?}", msg;
254 Unreachable, Bug)),
255 }
256 }
257
258 let elapsed = start.elapsed().as_secs_f64();
259 Ok(stopwatch(elapsed, n, byts))
260}
261
262#[macro_export]
263macro_rules! table_row { () => { "|{:<60}|{:^15}|{:^15}|" } }
264#[macro_export]
265macro_rules! table_single_line { () => { "+{:-<60}+{:-^15}+{:-^15}+" } }
266#[macro_export]
267macro_rules! table_double_line { () => { "+{:=<60}+{:=^15}+{:=^15}+" } }
268
269pub fn append_table(
270 path: PathBuf,
271 table: Vec<(String, u64, u64)>,
272)
273 -> Outcome<()>
274{
275 let mut file = res!(std::fs::OpenOptions::new().create(true).append(true).open(path));
276
277 res!(file.write_all(fmt!("\n").as_bytes()));
278 for line in [
279 fmt!(table_single_line!(), "", "", ""),
280 fmt!(table_row!(),
281 fmt!("{} {:?}", timefmt(time::SystemTime::now()), res!(hostname::get())),
282 "TPS",
283 "BW [Mb]/[s]",
284 ),
285 fmt!(table_double_line!(), "", "", ""),
286 ] {
287 test!(sync_log::stream(), "{}", line);
288 res!(file.write_all(fmt!("{}\n", line).as_bytes()));
289 }
290 for (s, tps, bw) in table {
291 for line in [
292 fmt!(table_row!(), s, tps, bw),
293 fmt!(table_single_line!(), "", "", ""),
294 ] {
295 test!(sync_log::stream(), "{}", line);
296 res!(file.write_all(fmt!("{}\n", line).as_bytes()));
297 }
298 }
299 Ok(())
300}
301