Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/bot_zone.rs

35.8 KiB, 160 runs

created by r1870400018:745, 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::{
4 cfg::ZoneConfig,
5 constant,
6 index::ZoneInd,
7 },
8 bots::{
9 base::bot_deps::*,
10 worker::bot::WorkerType,
11 },
12 comm::{
13 channels::{
14 BotChannels,
15 ChooseBot,
16 ZoneWorkerChannels,
17 },
18 response::Responder,
19 },
20 file::{
21 core::FileEntry,
22 floc::FileNum,
23 state::{
24 FileState,
25 FileStateMap,
26 Present,
27 },
28 zdir::ZoneDir,
29 },
30};
31
32use oxedyne_fe2o3_core::{
33 channels::Recv,
34};
35use oxedyne_fe2o3_data::time::Timestamp;
36use oxedyne_fe2o3_jdat::id::NumIdDat;
37
38use std::{
39 collections::BTreeMap,
40 fs::{
41 self,
42 File,
43 },
44 sync::Arc,
45 time::Instant,
46};
47
48#[derive(Clone, Debug, Default)]
49pub struct Resource {
50 pub size: usize,
51 pub ancillary_size: usize,
52 pub time: Timestamp,
53}
54
55#[derive(Clone, Debug, Default)]
56pub struct ZoneState {
57 pub caches: Vec<Resource>,
58 pub files: Vec<Resource>,
59}
60
61#[derive(Debug)]
62pub struct ZoneBot<
63 const UIDL: usize,
64 UID: NumIdDat<UIDL>,
65 ENC: Encrypter,
66 KH: Hasher,
67 PR: Hasher,
68 CS: Checksummer,
69>{
70 // Identity
71 zind: ZoneInd,
72 // Bot
73 sem: Semaphore,
74 errc: Arc<Mutex<usize>>,
75 log_stream_id: String,
76 // Config
77 zdir: ZoneDir,
78 // Comms
79 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
80 // API
81 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
82 // State
83 active: bool,
84 fnum: FileNum,
85 igbot_live: bool,
86 inited: bool,
87 size: usize,
88 trep: Instant,
89 zstat: ZoneState,
90}
91
92impl<
93 const UIDL: usize,
94 UID: NumIdDat<UIDL> + 'static,
95 ENC: Encrypter + 'static,
96 KH: Hasher + 'static,
97 PR: Hasher,
98 CS: Checksummer,
99>
100 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for ZoneBot<UIDL, UID, ENC, KH, PR, CS>
101{
102 bot_methods!();
103
104 fn go(&mut self) {
105
106 sync_log::set_stream(self.log_stream_id());
107
108 if self.no_init() { return; }
109 self.now_listening();
110 loop {
111
112 if self.trep.elapsed() > self.cfg().zone_state_update_interval() {
113 // Automated state reporting.
114 self.trep = Instant::now();
115 if let Err(e) = self.chans().sup().send(
116 OzoneMsg::ZoneState(*self.zind, self.zone_state().clone())
117 ) {
118 self.result(&Err(err!(e,
119 "{}: Cannot send zone state update to supervisor.", self.ozid();
120 Channel, Write)));
121 }
122 }
123
124 if self.listen().must_end() { break; }
125 }
126 }
127
128 fn listen(&mut self) -> LoopBreak {
129 match self.chan_in().recv_timeout(self.cfg().zone_state_update_interval()) {
130 Recv::Result(Err(e)) => self.err_cannot_receive(err!(e,
131 "{}: Waiting for message.", self.ozid();
132 IO, Channel)),
133 Recv::Result(Ok(msg)) => match msg {
134 // COMMAND
135 OzoneMsg::Channels(chans, resp) => {
136 self.set_chans(chans);
137 match self.chans().get_bot(&self.ozid()) {
138 Err(e) => self.error(e),
139 Ok(chan) => {
140 let chan_clone = chan.clone();
141 self.set_chan_in(chan_clone);
142 },
143 }
144 let resp2 = Responder::new(Some(self.ozid()));
145 match self.broadcast(OzoneMsg::Channels(self.chans().clone(), resp2.clone())) {
146 Err(e) => self.error(e),
147 Ok(zwbots) => {
148 // Channels are handed out at start-up, to bots that may only just
149 // have been scheduled, so this is held to the control deadline.
150 for _ in 0..zwbots.total_bot_count() {
151 match resp2.recv_timeout(constant::CONTROL_REQUEST_TIMEOUT) {
152 Err(e) => self.error(e),
153 Ok(OzoneMsg::ChannelsReceived(_)) => (),
154 m => self.error(err!(
155 "Received {:?}, expecting ChannelsReceived confirmation.", m;
156 Channel)),
157 }
158 }
159 },
160 }
161 let result = resp.send(OzoneMsg::ChannelsReceived(self.ozid().clone()));
162 self.result(&result);
163 },
164 //OzoneMsg::ClearCache(_) | OzoneMsg::GetUsers(_) => {
165 OzoneMsg::ClearCache(_) => {
166 match self.fwd_msg_to_pool(&WorkerType::Cache, msg) {
167 Err(e) => self.error(e),
168 Ok(_) => (),
169 }
170 },
171 OzoneMsg::GcControl(gc_ctrl, resp) => {
172 let result = match self.fwd_msg_to_pool(&WorkerType::File,
173 OzoneMsg::GcControl(gc_ctrl, Responder::none(Some(self.ozid()))),
174 ) {
175 Err(e) => Err(e),
176 Ok(_) => Ok(OzoneMsg::Ok),
177 };
178 self.respond(result, &resp);
179 },
180 OzoneMsg::GetZoneDir(resp) => {
181 self.respond(Ok(OzoneMsg::ZoneDir(*self.zind(), self.zdir().clone())), &resp);
182 },
183 OzoneMsg::ZoneInit(zdir, zcfg, resp) => {
184 // The database is not ready until every zone says this, and a zone that
185 // failed says why, because the caller of `O3db::start` is waiting on it.
186 let result = self.init_zone(zdir, zcfg);
187 match &result {
188 Err(e) => self.error(e.clone()),
189 Ok(()) => info!(sync_log::stream(), "{}: Zone {} init complete",
190 self.ozid(), self.zind),
191 }
192 self.respond(result.map(|()| OzoneMsg::Ok), &resp);
193 },
194 // WORK
195 OzoneMsg::CacheSize(b, size, ancillary_size) => {
196 if b+1 > self.zone_state().caches.len() {
197 self.error(err!(
198 "{}: The BotPoolInd for a cache size update, {}, exceeds the \
199 number of cbot slots {}.", self.ozid(), b, self.zone_state().caches.len();
200 Bug, Mismatch, Index, Size));
201 } else {
202 match Timestamp::now() {
203 Err(e) => self.error(e),
204 Ok(time) => {
205 self.zone_state_mut().caches[b] = Resource {
206 size,
207 ancillary_size,
208 time,
209 };
210 },
211 }
212 }
213 },
214 OzoneMsg::DumpCacheRequest(_) => {
215 match self.fwd_msg_to_pool(&WorkerType::Cache, msg) {
216 Err(e) => self.error(e),
217 Ok(_) => (),
218 }
219 },
220 OzoneMsg::DumpFiles(resp) => {
221 match self.read_files() {
222 Err(e) => self.error(e),
223 Ok(map) => self.respond(Ok(OzoneMsg::Files(*self.zind(), map)), &resp),
224 }
225 },
226 OzoneMsg::DumpFileStatesRequest(_) => {
227 match self.fwd_msg_to_pool(&WorkerType::File, msg) {
228 Err(e) => self.error(e),
229 Ok(_) => (),
230 }
231 },
232 OzoneMsg::NewLiveFile(_, _) => {
233 match self.fwd_msg_to_pool(&WorkerType::Writer, msg) {
234 Err(e) => self.error(e),
235 Ok(_) => (),
236 }
237 },
238 OzoneMsg::NextLiveFile(resp) => {
239 // [4] Respond to the wbot request for the next live file in the sequence.
240 //
241 // The number is claimed rather than counted to, so that reusing one
242 // is unconstructible rather than merely avoided.
243 //
244 // The counter is seeded from the highest number the zone survey found
245 // on disk, which is what makes it correct today. It was not always:
246 // reissuing a number already in use replaced the sealed file's record
247 // state wholesale, so its superseded bytes could never be reclaimed
248 // and every supersession afterwards raised an error. Creating the file
249 // is atomic, so a number that is taken cannot be handed out however
250 // the counter came by it.
251 //
252 // It also makes a zone written by two processes safe, each seeding
253 // from the same disk and so reaching the same next number. Numbers
254 // are then not contiguous, which nothing depends on: a zone's files
255 // are found by reading the directory.
256 let mut fnum = self.fnum;
257 let mut claim = None;
258 // Bounded so that a directory in a state nobody expected stops the
259 // bot rather than spinning it.
260 for _ in 0..constant::LIVE_FILE_CLAIM_LIMIT {
261 fnum += 1;
262 match self.zdir().claim(fnum) {
263 Ok(true) => { claim = Some(Ok(fnum)); break; },
264 Ok(false) => (),
265 Err(e) => { claim = Some(Err(e)); break; },
266 }
267 }
268 let result = match claim {
269 Some(Ok(fnum)) => {
270 self.fnum = fnum;
271 Ok(OzoneMsg::UseLiveFile(fnum))
272 },
273 Some(Err(e)) => Err(err!(e,
274 "{}: No live file could be claimed in zone {:?}.",
275 self.ozid(), self.zdir();
276 IO, File, Create)),
277 None => Err(err!(
278 "{}: No live file number could be claimed in zone {:?} within {} \
279 attempts from {}; every number tried was already taken.",
280 self.ozid(), self.zdir(), constant::LIVE_FILE_CLAIM_LIMIT, self.fnum;
281 IO, File, Create, Excessive)),
282 };
283 // The writer asking is waiting on this, and so is the caller whose write
284 // needed the new file, so a failure is answered rather than only logged.
285 if let Err(e) = &result {
286 self.error(e.clone());
287 }
288 self.respond(result, &resp);
289 },
290 OzoneMsg::ShardFileSize(b, size) => {
291 if b+1 > self.zone_state().files.len() {
292 self.error(err!(
293 "{}: The BotPoolInd for a file state shard size update, {}, exceeds the \
294 number of fbot slots {}.", self.ozid(), b, self.zone_state().files.len();
295 Bug, Mismatch, Size));
296 } else {
297 match Timestamp::now() {
298 Err(e) => self.error(e),
299 Ok(time) => {
300 self.zone_state_mut().files[b] = Resource {
301 size,
302 ancillary_size: 0,
303 time,
304 };
305 },
306 }
307 }
308 },
309 _ => return self.listen_more(msg),
310 },
311 Recv::Empty => (),
312 }
313 LoopBreak(false)
314 }
315}
316
317impl<
318 const UIDL: usize,
319 UID: NumIdDat<UIDL> + 'static,
320 ENC: Encrypter + 'static,
321 KH: Hasher + 'static,
322 PR: Hasher,
323 CS: Checksummer,
324>
325 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for ZoneBot<UIDL, UID, ENC, KH, PR, CS>
326{
327 ozonebot_methods!();
328}
329
330impl<
331 const UIDL: usize,
332 UID: NumIdDat<UIDL> + 'static,
333 ENC: Encrypter + 'static,
334 KH: Hasher + 'static,
335 PR: Hasher,
336 CS: Checksummer,
337>
338 ZoneBot<UIDL, UID, ENC, KH, PR, CS>
339{
340 pub fn new(
341 args: BotInitArgs<UIDL, UID, ENC, KH, PR, CS>,
342 zind: ZoneInd,
343 )
344 -> Self
345 {
346 Self {
347 // Identity
348 zind,
349 // Bot
350 sem: args.sem,
351 errc: Arc::new(Mutex::new(0)),
352 log_stream_id: args.log_stream_id,
353 // Config
354 zdir: ZoneDir::default(),
355 // Comms
356 chan_in: args.chan_in,
357 // API
358 api: args.api,
359 // State
360 active: true,
361 fnum: 0,
362 igbot_live: true,
363 inited: false,
364 size: 0,
365 trep: Instant::now(),
366 zstat: ZoneState::default(),
367 }
368 }
369
370 fn zind(&self) -> &ZoneInd { &self.zind }
371 fn zdir(&self) -> &ZoneDir { &self.zdir }
372 fn zone_state(&self) -> &ZoneState { &self.zstat }
373 fn zone_state_mut(&mut self) -> &mut ZoneState { &mut self.zstat }
374
375 fn get_zwbots(&self) -> Outcome<&ZoneWorkerChannels<UIDL, UID, ENC, KH>> {
376 self.chans().get_zwbots(&self.zind())
377 }
378
379 fn fwd_msg_to_pool(
380 &self,
381 wtyp: &WorkerType,
382 msg: OzoneMsg<UIDL, UID, ENC, KH>,
383 )
384 -> Outcome<usize>
385 {
386 let pool = &res!(self.get_zwbots())[wtyp];
387 pool.send_to_all(msg)
388 }
389
390 fn broadcast(
391 &mut self,
392 msg: OzoneMsg<UIDL, UID, ENC, KH>,
393 )
394 -> Outcome<&ZoneWorkerChannels<UIDL, UID, ENC, KH>>
395 {
396 let zwbots = res!(self.get_zwbots());
397 res!(zwbots[&WorkerType::Cache].send_to_all(msg.clone()));
398 res!(zwbots[&WorkerType::File].send_to_all(msg.clone()));
399 res!(zwbots[&WorkerType::InitGarbage].send_to_all(msg.clone()));
400 res!(zwbots[&WorkerType::Reader].send_to_all(msg.clone()));
401 res!(zwbots[&WorkerType::Scan].send_to_all(msg.clone()));
402 res!(zwbots[&WorkerType::Writer].send_to_all(msg.clone()));
403 Ok(zwbots)
404 }
405
406 pub fn activate(mut self) -> Self {
407 self.active = true;
408 self
409 }
410
411// /// Survey the existing data and index files and send the file state maps to the zone file bots.
412// pub fn survey_files(&mut self) -> Outcome<Vec<FileStateMap>> {
413//
414// info!(sync_log::stream(), "{}: Surveying files for zone {}...", self.ozid(), self.zind());
415//
416// let mut shards = Vec::new();
417// let nf = self.cfg().num_fbots_per_zone();
418// let mut dir_size: usize = 0;
419// let mut max_data_fnum: u32 = 0;
420// let mut max_data_size: usize = 0;
421// // 1. Count the number of files in the zone directory.
422// // <deleted>
423//
424// trace!(sync_log::stream(), "{}: Initializing {} file state shards", self.ozid(), nf);
425//
426// // 2. We will record whether both a data and index file are present for each file
427// // number (i.e. Present::Pair) or whether just one is present (i.e.
428// // Present::Solo(FileType)).
429// for _ in 0..nf {
430// shards.push(FileStateMap::default());
431// }
432//
433// // 3. Loop through all objects in the zone directory, looking for files.
434// for entry in res!(std::fs::read_dir(&self.zdir().dir)) {
435// let entry = res!(entry);
436// let path = entry.path();
437// if path.is_file() {
438// trace!(sync_log::stream(), "{}: Processing file {:?}", self.ozid(), path);
439// // 4. We found a file, now extract the file number and type from the file name.
440// let (fnum, ftyp) = res!(ZoneDir::ozone_file_number_and_type(&path));
441// // 5. Get the file size.
442// let file = res!(std::fs::OpenOptions::new()
443// .read(true)
444// .open(&path)
445// );
446// let meta = res!(file.metadata());
447// let flen = meta.len() as usize;
448// dir_size += flen;
449// // 6. Keep a track of the highest data file number, in order to set the live
450// // filenumber counter for the zone later.
451// if ftyp == FileType::Data {
452// if fnum > max_data_fnum {
453// max_data_fnum = fnum;
454// max_data_size = flen;
455// }
456// }
457//
458// // 7. Determine and record the state of the files for this file number, for use
459// // later by the `init_caches` method.
460// let i = FileStateMap::shard_index(fnum, nf);
461// match shards[i].map_mut().get_mut(&fnum) {
462// None => {
463// let mut fs = FileState::default();
464// fs.set_present(Present::Solo(ftyp));
465// match ftyp {
466// FileType::Data => fs.set_data_file_size(flen),
467// FileType::Index => fs.set_index_file_size(flen),
468// }
469// shards[i].map_mut().insert(fnum, fs);
470// },
471// Some(fs) => {
472// match fs.present().clone() {
473// Present::Solo(ftyp2) => {
474// if ftyp2 != ftyp {
475// fs.set_present(Present::Pair);
476// match ftyp {
477// FileType::Data => fs.set_data_file_size(flen),
478// FileType::Index => fs.set_index_file_size(flen),
479// }
480// } else {
481// return Err(err!(
482// "The file {:?} has already been surveyed, the file system \
483// should not permit such duplicate files.", path,
484// ), Unreachable, Bug));
485// }
486// },
487// Present::Pair => return Err(err!(
488// "The file {:?} has already been surveyed, the file system \
489// should not permit such duplicate files.", path,
490// ), Unreachable, Bug)),
491// }
492// },
493// }
494// }
495// }
496// // 8. Set the live file number for the zone. The function ozone_file_number_and_type
497// // ensures that max_data_filenum does not exceed u32::MAX. The
498// // zone::Zone::init_live method increments the first live file in the zone by one,
499// // so to avoid restarts simply creating additional empty files, wind the counter
500// // back one when necessary.
501// let max_data_file_size_ratio =
502// (max_data_size as f64) / (self.cfg().data_file_max_bytes as f64);
503// if max_data_fnum > 0
504// && max_data_file_size_ratio < constant::LIVE_FILE_INIT_SIZE_RATIO_THRESHOLD {
505// max_data_fnum = max_data_fnum - 1;
506// }
507// self.fnum = max_data_fnum;
508// // 9. Set the directory size for the zone.
509// self.size = dir_size;
510//
511// Ok(shards)
512// }
513
514
515 /// Configures the zone, surveys its files, gives each writer a live file and loads the caches.
516 fn init_zone(
517 &mut self,
518 zdir: ZoneDir,
519 zcfg: ZoneConfig,
520 )
521 -> Outcome<()>
522 {
523 self.zone_state_mut().caches = vec![Resource::default(); zcfg.ncbots];
524 self.zone_state_mut().files = vec![Resource::default(); zcfg.nfbots];
525 res!(self.fwd_msg_to_pool(&WorkerType::Cache, OzoneMsg::SetCacheSizeLimit(zcfg.cache_size_lim)));
526 self.zdir = zdir.clone();
527 res!(self.broadcast(OzoneMsg::ZoneDir(*self.zind(), zdir)));
528 let shards = res!(self.survey_files());
529 if zcfg.init_load_caches {
530 res!(self.init_caches(shards));
531 }
532 Ok(())
533 }
534
535 /// Survey the existing data and index files and send the file state maps to the zone file bots.
536 pub fn survey_files(&mut self) -> Outcome<Vec<FileStateMap>> {
537
538 info!(sync_log::stream(), "{}: Surveying {} files...", self.ozid(), self.zind());
539
540 let mut shards = Vec::new();
541 let nf = self.cfg().num_fbots_per_zone();
542 let mut dir_size: usize = 0;
543 let mut max_data_fnum: u32 = 0;
544 // Track incomplete data files for WriterBot initialization.
545 let mut incomplete_files = Vec::new();
546
547 trace!(sync_log::stream(), "{}: Initializing {} file state shards", self.ozid(), nf);
548
549 // 2. We will record whether both a data and index file are present for each file
550 // number (i.e. Present::Pair) or whether just one is present (i.e.
551 // Present::Solo(FileType)).
552 for _ in 0..nf {
553 shards.push(FileStateMap::default());
554 }
555
556 // 3. Loop through all objects in the zone directory, looking for files.
557 for entry in res!(std::fs::read_dir(&self.zdir().dir)) {
558 let entry = res!(entry);
559 let path = entry.path();
560 if path.is_file() {
561 trace!(sync_log::stream(), "{}: Processing file {:?}", self.ozid(), path);
562 // 3.1 Remove any abandoned garbage collection transcription. The data file it
563 // was copied from is intact, so the leftover is of no use, and leaving it
564 // in place would fail the survey for the whole zone.
565 if ZoneDir::is_gc_temp_file(&path) {
566 warn!(sync_log::stream(),
567 "{}: Removing {:?}, an abandoned garbage collection temporary.",
568 self.ozid(), path);
569 res!(fs::remove_file(&path));
570 continue;
571 }
572 // 4. We found a file, now extract the file number and type from the file name.
573 let (fnum, ftyp) = match ZoneDir::ozone_file_number_and_type(&path) {
574 Ok(pair) => pair,
575 Err(e) => {
576 warn!(sync_log::stream(),
577 "{}: Ignoring {:?}, which is not an ozone data or index file, \
578 caused by {}.", self.ozid(), path, e);
579 continue;
580 },
581 };
582 // 5. Get the file size.
583 let file = res!(std::fs::OpenOptions::new()
584 .read(true)
585 .open(&path)
586 );
587 let meta = res!(file.metadata());
588 let flen = meta.len() as usize;
589 dir_size += flen;
590 // 6. Keep a track of the highest data file number, in order to set the live
591 // filenumber counter for the zone later.
592 if ftyp == FileType::Data {
593 if fnum > max_data_fnum {
594 max_data_fnum = fnum;
595 }
596 // Check if this is an incomplete data file.
597 let ratio = flen as f64 / self.cfg().data_file_max_bytes as f64;
598 if ratio < constant::LIVE_FILE_INIT_SIZE_RATIO_THRESHOLD {
599 incomplete_files.push((fnum, flen));
600 }
601 }
602
603 // 7. Determine and record the state of the files for this file number, for use
604 // later by the `init_caches` method.
605 let i = FileStateMap::shard_index(fnum, nf);
606 match shards[i].map_mut().get_mut(&fnum) {
607 None => {
608 let mut fs = FileState::default();
609 fs.set_present(Present::Solo(ftyp));
610 match ftyp {
611 FileType::Data => fs.set_data_file_size(flen),
612 FileType::Index => fs.set_index_file_size(flen),
613 }
614 shards[i].map_mut().insert(fnum, fs);
615 },
616 Some(fs) => {
617 match fs.present().clone() {
618 Present::Solo(ftyp2) => {
619 if ftyp2 != ftyp {
620 fs.set_present(Present::Pair);
621 match ftyp {
622 FileType::Data => fs.set_data_file_size(flen),
623 FileType::Index => fs.set_index_file_size(flen),
624 }
625 } else {
626 return Err(err!(
627 "The file {:?} has already been surveyed, the file system \
628 should not permit such duplicate files.", path;
629 Unreachable, Bug));
630 }
631 },
632 Present::Pair => return Err(err!(
633 "The file {:?} has already been surveyed, the file system \
634 should not permit such duplicate files.", path;
635 Unreachable, Bug)),
636 }
637 },
638 }
639 }
640 }
641
642 // Sort incomplete files by number descending.
643 incomplete_files.sort_by(|a, b| b.0.cmp(&a.0));
644
645 // 8. Seed the live file number counter for the zone with the highest file number
646 // found on disk. The function ozone_file_number_and_type ensures that
647 // max_data_fnum does not exceed u32::MAX. This must happen before the writer
648 // bots are given their live files, because init_writer_live_files advances the
649 // counter past every new file number it creates. The counter must never be left
650 // at or below a file number that is already in use: handing an in-use number back
651 // out as the next live file makes the receiving fbot replace that file's
652 // FileState, silently discarding every record entry it holds, after which no
653 // record in that file can ever be flagged old and its garbage is unreclaimable.
654 self.fnum = max_data_fnum;
655
656 // Initialise WriterBot live files.
657 res!(self.init_writer_live_files(&incomplete_files));
658
659 // 9. Set the directory size for the zone.
660 self.size = dir_size;
661
662 Ok(shards)
663 }
664
665 /// Assigns live files to WriterBots, re-using discovered incomplete files before creating
666 /// new ones, and leaves the zone live file counter above every number assigned.
667 fn init_writer_live_files(
668 &mut self,
669 incomplete_files: &[(FileNum, usize)],
670 )
671 -> Outcome<()>
672 {
673 let n_w = self.cfg().num_wbots_per_zone as usize;
674 let wbots = res!(self.get_zwbots())[&WorkerType::Writer].clone();
675
676 trace!(sync_log::stream(), "{}: Assigning live files to {} writer bots.", self.ozid(), n_w);
677
678 let resp = Responder::new(Some(self.ozid()));
679 for i in 0..n_w {
680 let fnum = if i < incomplete_files.len() {
681 // Use existing incomplete file. Its number is at or below the counter,
682 // which was seeded with the highest file number on disk, so the counter
683 // needs no adjustment.
684 incomplete_files[i].0
685 } else {
686 // There are not enough existing incomplete files to assign to wbots, create new
687 // file number.
688 self.fnum += 1;
689 self.fnum
690 };
691
692 let wbot = res!(wbots.get_bot(i));
693 res!(wbot.send(OzoneMsg::NewLiveFile(Some(fnum), resp.clone())));
694 }
695
696 // Start-up work, held to the control deadline: see constant::CONTROL_REQUEST_TIMEOUT.
697 let (_, msgs) = res!(resp.recv_number(n_w, constant::CONTROL_REQUEST_WAIT));
698 for msg in msgs {
699 match msg {
700 OzoneMsg::Error(e) => return Err(err!(e,
701 "{}: In response to NewLiveFile message.", self.ozid();
702 Channel)),
703 OzoneMsg::Ok => (),
704 msg => return Err(err!(
705 "{}: Unexpected response to NewLiveFile message: {:?}", self.ozid(), msg;
706 Channel, Unexpected)),
707 };
708 }
709
710 Ok(())
711 }
712
713 pub fn read_files(&mut self) -> Outcome<BTreeMap<String, FileEntry>> {
714 let mut map = BTreeMap::new();
715 let list = res!(fs::read_dir(self.zdir().dir.clone()));
716 for item in list {
717 let item = res!(item);
718 let path = item.path();
719 let file = res!(File::open(&path));
720 let metadata = res!(file.metadata());
721 let ftyp = metadata.file_type();
722 let typ = if ftyp.is_dir() {
723 fmt!("d")
724 } else if ftyp.is_symlink() {
725 fmt!("s")
726 } else {
727 fmt!("f")
728 };
729 let elapsed = res!(metadata.modified()).elapsed();
730 let mods = res!(elapsed).as_secs();
731 //let mods = res!(res!(metadata.modified()).elapsed()).as_secs();
732 let size = metadata.len();
733 let name = match path.file_name() {
734 Some(s) => match s.to_os_string().into_string() {
735 Ok(s) => s,
736 Err(_) => fmt!("{}", path.display()),
737 }
738 None => fmt!("{}", path.display()),
739 };
740 map.insert(name.clone(), FileEntry { typ, size, mods, name });
741 }
742 Ok(map)
743 }
744
745 /// Initialise the zone data caches based on file survey data. The survey collected all file
746 /// sizes, but these will be recalculated as data is added to the cache.
747 pub fn init_caches(
748 &mut self,
749 shards: Vec<FileStateMap>,
750 )
751 -> Outcome<()>
752 {
753 // 1. Prepare to make a bunch of requests to InitGarbageBots.
754 let resp = Responder::new(Some(self.ozid()));
755 let mut bot_requests = 0;
756
757 for mut shard in shards {
758
759 // 2. Loop through all previously surveyed files and check their status.
760 let mut missing_data_files: Option<Vec<u32>> = None;
761 for (fnum, fstate) in shard.map_mut() {
762 match fstate.present() {
763 Present::Pair => {
764 // We have a data file and associated index file, and can therefore save
765 // some time by reading the (generally smaller) index file and loading its
766 // keys and file locations into the zone cache. Send the caching request
767 // to a randomised InitGarbageBot.
768 // CACHE INDEX FILE - since both the index and data file is present.
769
770 // Reset file sizes and relay to the igbot as a checksum.
771 let dat_file_size = fstate.get_data_file_size();
772 let ind_file_size = fstate.get_index_file_size();
773 fstate.reset_data_file_size();
774 fstate.reset_index_file_size();
775
776 // Send caching request, with responder.
777 let (bot, j) = res!(self.get_zwbots())[&WorkerType::InitGarbage]
778 .choose_bot(&ChooseBot::Randomly);
779 if let Err(e) = bot.send(
780 OzoneMsg::CacheIndexFile {
781 fnum: *fnum,
782 dat_file_size,
783 ind_file_size,
784 resp: resp.clone(),
785 }
786 ) {
787 return Err(err!(e,
788 "Cannot send cache index file request to igbot {}", j;
789 Channel, Write));
790 } else {
791 bot_requests += 1;
792 }
793 },
794 Present::Solo(FileType::Data) => {
795 // We have only a data file, and can manually index the key and value pairs
796 // it contains. The caching request will result in the re-creation of the
797 // index file.
798 // CACHE DATA FILE - since the index file is not present.
799
800 // Reset file size and relay to the igbot as a checksum.
801 let dat_file_size = fstate.get_data_file_size();
802 fstate.reset_data_file_size();
803
804 //// Send caching request, with responder.
805 let (bot, j) = res!(self.get_zwbots())[&WorkerType::InitGarbage]
806 .choose_bot(&ChooseBot::Randomly);
807 if let Err(e) = bot.send(
808 OzoneMsg::CacheDataFile {
809 fnum: *fnum,
810 dat_file_size,
811 resp: resp.clone(),
812 }
813 ) {
814 return Err(err!(e,
815 "Cannot send cache data file request to igbot {}", j;
816 Channel, Write));
817 } else {
818 bot_requests += 1;
819 }
820 },
821 Present::Solo(FileType::Index) => {
822 // An index file without a data file is bad, it means data has been lost.
823 // Keep a track of all such instances and warn the user.
824 match missing_data_files {
825 None => missing_data_files = Some(vec![*fnum]),
826 Some(ref mut missing) => missing.push(*fnum),
827 }
828 },
829 }
830 }
831 // 3. Warn the user of missing data files.
832 if let Some(missing) = missing_data_files {
833 for fnum in missing {
834 warn!(sync_log::stream(), "{:?}: Data file {} is missing, suggesting loss of data.",
835 self.ozid(), fnum);
836 }
837 }
838 }
839
840 // 4. Wait for and collect all request responses. Loading one large file under a busy
841 // disk can take longer than a bot request is allowed, and this is start-up work that
842 // nothing is waiting to be served behind, so it is held to the control deadline.
843 for _ in 0..bot_requests {
844 match resp.recv_timeout(constant::CONTROL_REQUEST_TIMEOUT) {
845 Err(e) => return Err(err!(e,
846 "While collecting cache initialisation request responses.";
847 IO, Channel, Read)),
848 Ok(OzoneMsg::Ok) => (),
849 Ok(msg) => return Err(err!(
850 "Unrecognised cache initialisation request response: {:?}", msg;
851 Channel)),
852 }
853 }
854
855 Ok(())
856 }
857}