Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/worker/bot_writer.rs

26.9 KiB, 151 runs

created by r1870400018:759, 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 bots::{
5 base::bot_deps::*,
6 worker::{
7 syncer::{
8 Handed,
9 SyncPolicy,
10 Syncer,
11 },
12 worker_deps::*,
13 },
14 },
15 file::{
16 core::FileType,
17 floc::{
18 FileNum,
19 StoredFileLocation,
20 },
21 live::LivePair,
22 },
23 test::hooks,
24};
25
26use oxedyne_fe2o3_iop_db::api::Meta;
27use oxedyne_fe2o3_jdat::id::NumIdDat;
28
29use std::{
30 fs::File,
31 io::{
32 Seek,
33 SeekFrom,
34 Write,
35 },
36 sync::Arc,
37};
38
39/// Each `WriterBot` in a zone has its own `LivePair`.
40pub struct WriterBot<
41 const UIDL: usize,
42 UID: NumIdDat<UIDL>,
43 ENC: Encrypter,
44 KH: Hasher,
45 PR: Hasher,
46 CS: Checksummer,
47>{
48 // Identity
49 wind: WorkerInd,
50 wtyp: WorkerType,
51 // Bot
52 sem: Semaphore,
53 errc: Arc<Mutex<usize>>,
54 log_stream_id: String,
55 // Config
56 zdir: ZoneDir,
57 // Comms
58 chan_in: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>,
59 // API
60 api: OzoneApi<UIDL, UID, ENC, KH, PR, CS>,
61 // State
62 active: bool,
63 inited: bool,
64 lpair: LivePair,
65 syncer: Syncer<UIDL, UID, ENC, KH>, // makes each record durable, then releases it
66}
67
68impl<
69 const UIDL: usize,
70 UID: NumIdDat<UIDL> + 'static,
71 ENC: Encrypter + 'static,
72 KH: Hasher + 'static,
73 PR: Hasher,
74 CS: Checksummer,
75>
76 WorkerBot<UIDL, UID, ENC, KH, PR, CS> for WriterBot<UIDL, UID, ENC, KH, PR, CS>
77{
78 workerbot_methods!();
79}
80
81impl<
82 const UIDL: usize,
83 UID: NumIdDat<UIDL> + 'static,
84 ENC: Encrypter + 'static,
85 KH: Hasher + 'static,
86 PR: Hasher,
87 CS: Checksummer,
88>
89 OzoneBot<UIDL, UID, ENC, KH, PR, CS> for WriterBot<UIDL, UID, ENC, KH, PR, CS>
90{
91 ozonebot_methods!();
92}
93
94impl<
95 const UIDL: usize,
96 UID: NumIdDat<UIDL> + 'static,
97 ENC: Encrypter + 'static,
98 KH: Hasher + 'static,
99 PR: Hasher,
100 CS: Checksummer,
101>
102 Bot<{ BID_LEN }, Bid, OzoneMsg<UIDL, UID, ENC, KH>> for WriterBot<UIDL, UID, ENC, KH, PR, CS>
103{
104 bot_methods!();
105
106 fn go(&mut self) {
107
108 sync_log::set_stream(self.log_stream_id());
109
110 if self.no_init() { return; }
111 self.now_listening();
112 loop {
113 if self.wind().b() < self.cfg().num_bots_per_zone((&self).wtyp()) {
114 if self.listen().must_end() { break; }
115 } else {
116 // This bot is to be terminated. Forward incoming messages to the remaining bots of
117 // this type.
118 }
119 }
120 // Everything written is released, and made durable where the policy owes it, before this
121 // bot ends, so a database closed after a write has that write on disk.
122 let result = self.syncer.finish();
123 self.result(&result);
124 }
125
126 fn listen(&mut self) -> LoopBreak {
127 match self.chan_in().recv() {
128 Err(e) => self.err_cannot_receive(err!(e,
129 "{}: Waiting for message.", self.ozid();
130 IO, Channel)),
131 Ok(msg) => {
132 if let Some(msg) = self.listen_worker(msg) {
133 match msg {
134 // COMMAND
135 OzoneMsg::NewLiveFile(fnum_opt, resp) => {
136 let result = match fnum_opt {
137 Some(fnum) => {
138 // Direct initialization with provided file number.
139 self.lpair.fnum = fnum;
140 self.open_live_pair()
141 },
142 None => {
143 // Routine request for new live file.
144 self.new_live_pair().map(|_| ())
145 }
146 };
147 // Answered either way: a zone starting up waits on every writer.
148 if let Err(e) = &result {
149 self.error(e.clone());
150 }
151 self.respond(result.map(|()| OzoneMsg::Ok), &resp);
152 }
153 // WORK
154 OzoneMsg::Write{
155 kstored,
156 vstored,
157 klen_cache,
158 cind,
159 meta,
160 cbpind,
161 resp: resp_w1,
162 } => {
163 let resp = resp_w1.clone();
164 if let Err(e) = self.write(
165 kstored,
166 vstored,
167 klen_cache,
168 cind,
169 meta,
170 cbpind,
171 resp_w1,
172 ) {
173 // The caller is waiting on this answer. Only logged, a failure
174 // here reached it as a responder timeout that named no cause.
175 self.error(e.clone());
176 self.respond(Err(e), &resp);
177 }
178 }
179 //OzoneMsg::Delete(kv, resp_w1) => {
180 // let result = self.write(kv, resp_w1);
181 // self.result(result);
182 //},
183 _ => return self.listen_more(msg),
184 }
185 }
186 },
187 }
188 LoopBreak(false)
189 }
190}
191
192impl<
193 const UIDL: usize,
194 UID: NumIdDat<UIDL> + 'static,
195 ENC: Encrypter + 'static,
196 KH: Hasher + 'static,
197 PR: Hasher,
198 CS: Checksummer,
199>
200 WriterBot<UIDL, UID, ENC, KH, PR, CS>
201{
202 /// Starts the writer's durability barrier thread, so that a writer which could never confirm
203 /// a write fails the database's start instead of its first write.
204 pub fn new(
205 args: ZoneWorkerInitArgs<UIDL, UID, ENC, KH, PR, CS>,
206 )
207 -> Outcome<Self>
208 {
209 let syncer = res!(Syncer::start(fmt!("{}", args.api.ozid), args.log_stream_id.clone()));
210 Ok(Self {
211 // Identity
212 wind: args.wind,
213 wtyp: args.wtyp,
214 // Bot
215 sem: args.sem,
216 errc: Arc::new(Mutex::new(0)),
217 log_stream_id: args.log_stream_id,
218 // Config
219 zdir: ZoneDir::default(),
220 // Comms
221 chan_in: args.chan_in,
222 api: args.api,
223 // State
224 active: false,
225 inited: false,
226 lpair: LivePair::default(),
227 syncer,
228 })
229 }
230
231 fn lpair(&self) -> &LivePair { &self.lpair }
232 fn lpair_mut(&mut self) -> &mut LivePair { &mut self.lpair }
233
234 //fn cached_livefile_len(&self) -> usize { self.lpair.dat_size as usize }
235
236 /// This is the main writer method which:
237 ///
238 /// 1. Appends a checksum to the key and value bytes then appends these to the current zone
239 /// `LivePair` data file (creating the next `LivePair` if necessary in order not to exceed
240 /// the file size limit).
241 /// 2. Inserts the key, location and possibly the value into the zone data cache.
242 /// 3. Appends the key and location to the `LivePair` index file, each with an appended
243 /// checksum.
244 ///
245 /// Returns on a write error, with nothing written to the cache or index file.
246 ///
247 ///```ignore
248 ///
249 /// Appended to data file Appended to index file
250 /// +---------------------+ -+ +- +---------------------+
251 /// | | | | | |
252 /// | key | | | | key |
253 /// | | | | | |
254 /// +---------------------+ +- StoredKey ---+ +---------------------+
255 /// | meta | | | | meta |
256 /// +---------------------+ | | +---------------------+
257 /// | checksum | | | | checksum |
258 /// +---------------------+ -+ +- +---------------------+ -+
259 /// | | | | | start | | part of a
260 /// | | | | +---------------------+ +-- FileLocation
261 /// | | | StoredIndex -+ | klen | |
262 /// | | | | +---------------------+ |
263 /// | value | | | | vlen | |
264 /// | | +- StoredValue | +---------------------+ -+
265 /// | | | | | checksum |
266 /// | | | +- +---------------------+
267 /// | | |
268 /// | | |
269 /// +---------------------+ |
270 /// | checksum | |
271 /// +---------------------+ -+
272 ///
273 ///```
274 fn write(
275 &mut self,
276 mut kbyts: Vec<u8>,
277 vstored: Vec<u8>,
278 klen_cache: usize,
279 cind: Option<usize>,
280 meta: Meta<UIDL, UID>,
281 cbpind: usize, // cbot pool index
282 resp_w1: Responder<UIDL, UID, ENC, KH>,
283 )
284 -> Outcome<()>
285 {
286 // A record appended now could be neither confirmed nor withdrawn, so a writer whose syncer
287 // has stopped refuses it before anything reaches the files. Appended first, it was
288 // reported failed and came back at the next start (2026-09-23).
289 if !self.syncer.is_running() {
290 return Err(err!(
291 "{}: The durability barrier thread has stopped, so the write was refused before \
292 anything was written.", self.ozid();
293 Thread, Missing, Write));
294 }
295 let start = res!(self.write_to_file(FileType::Data, vec![&kbyts[..], &vstored[..]]));
296
297 // Define the location.
298 let sfloc = res!(StoredFileLocation::new(
299 self.lpair().fnum,
300 start,
301 kbyts.len() as u64,
302 vstored.len() as u64,
303 self.api().schemes().checksummer().clone(),
304 ));
305 let istored = &sfloc.buf;
306
307 // Append key and location to the current index file.
308 res!(self.write_to_file(FileType::Index, vec![&kbyts[..], &istored[..]]));
309
310 // [11] The record is in the live pair, so say so before anything waits on the disk. The
311 // caller holds this answer to the short deadline, which then measures whether the
312 // writer is alive rather than how busy the machine's disk happens to be.
313 self.respond(Ok(OzoneMsg::Written), &resp_w1);
314
315 // [12] The syncer releases the record to a cbot once the barrier the policy asks for is
316 // behind it, so no observer sees a key before it is as durable as configured. The
317 // cbot then makes it readable and gives the caller its final answer.
318 let cbots = res!(self.cbots());
319 let cbot = res!(cbots.get_bot(cbpind)).clone();
320 kbyts.drain(..constant::CACHE_HASH_BYTES); // remove data pathway hash used to identify cbot
321 kbyts.truncate(klen_cache); // remove metadata
322 let resp = resp_w1.clone();
323 let insert = OzoneMsg::Insert(
324 kbyts,
325 Some(vstored),
326 cind,
327 sfloc.ref_file_location().clone(),
328 istored.len(),
329 meta,
330 resp_w1, // The cbot responds to the caller.
331 None,
332 );
333 let policy = SyncPolicy::of(self.cfg());
334 if let Err(e) = self.syncer.hand(Handed::Record { cbot, insert, resp, policy }) {
335 // The syncer stopped after the check above, and the record is in the files.
336 return Err(err!(e,
337 "{}: The record is written, but the durability barrier thread stopped before it \
338 could take it, so it is not confirmed durable.", self.ozid();
339 Thread, Write, Unconfirmed));
340 }
341
342 Ok(())
343 }
344
345 /// Hands the syncer the live pair every record from here on is appended to. It syncs through
346 /// handles of its own, which reach the same open files.
347 fn hand_pair(&self, lpair: &LivePair) -> Outcome<()> {
348 if hooks::pair_hand_fails() {
349 return Err(err!(
350 "{}: Live pair {} could not be duplicated for the syncer (test::hooks).",
351 self.ozid(), lpair.fnum;
352 IO, File));
353 }
354 let dat = match &lpair.dat.file {
355 Some(file) => res!(file.try_clone()),
356 None => return Err(err!(
357 "{}: The live data file {:?} is not open.", self.ozid(), lpair.dat.path;
358 Bug, Missing)),
359 };
360 let ind = match &lpair.ind.file {
361 Some(file) => res!(file.try_clone()),
362 None => return Err(err!(
363 "{}: The live index file {:?} is not open.", self.ozid(), lpair.ind.path;
364 Bug, Missing)),
365 };
366 self.syncer.hand(Handed::Pair(dat, ind))
367 }
368
369 /// Takes the live file the zone assigned at start-up, new or partly written.
370 fn open_live_pair(&mut self) -> Outcome<()> {
371 let lpair = res!(self.zdir().open_live(self.lpair.fnum));
372 // The syncer has the pair before the writer does, as at a rollover.
373 res!(self.hand_pair(&lpair));
374 self.lpair.close();
375 self.lpair = lpair;
376 self.register_live_file(self.lpair.fnum)
377 }
378
379 /// Closes a pair opened for a rollover that did not happen, and removes its files, which
380 /// nothing was written to.
381 fn abandon(&self, mut lpair: LivePair) {
382 lpair.close();
383 if lpair.dat.size > 0 || lpair.ind.size > 0 {
384 return; // not new after all, so not this writer's to remove
385 }
386 for path in [&lpair.dat.path, &lpair.ind.path] {
387 if let Err(e) = std::fs::remove_file(path) {
388 warn!(sync_log::stream(),
389 "{}: Could not remove {:?}, created for a rollover that did not happen: {}",
390 self.ozid(), path, e);
391 }
392 }
393 }
394
395 /// Tells the file's bot that the file is live, as a rollover does for the file it opens. The
396 /// bot otherwise first heard of a new file from its first record, which reaches it through the
397 /// syncer and a cache bot, so a writer could seal the file before the bot knew it existed and
398 /// the seal failed; and a partly written file taken over at start-up was never flagged live,
399 /// leaving it open to collection while it was still being written. Its accounting starts
400 /// empty, since the records already in such a file are counted as the zone loads them.
401 fn register_live_file(&self, fnum: FileNum) -> Outcome<()> {
402 let resp = Responder::new(Some(self.ozid()));
403 let bots = res!(self.fbots());
404 let (bot, _) = bots.choose_bot(&ChooseBot::ByFile(fnum));
405 res!(bot.send(OzoneMsg::OpenNewLiveFileState {
406 fnum_new: fnum,
407 new_dat_size: 0,
408 new_ind_size: 0,
409 resp: resp.clone(),
410 }));
411 // Start-up work, held to the control deadline: see constant::CONTROL_REQUEST_TIMEOUT.
412 match resp.recv_timeout(constant::CONTROL_REQUEST_TIMEOUT) {
413 Err(e) => Err(err!(e,
414 "{}: While registering live file {} with its file bot.", self.ozid(), fnum;
415 IO, Channel, Read)),
416 Ok(OzoneMsg::Ok) => Ok(()),
417 Ok(OzoneMsg::Error(e)) => Err(err!(e,
418 "{}: The file bot could not register live file {}.", self.ozid(), fnum;
419 IO, File)),
420 Ok(msg) => Err(err!(
421 "{}: Unrecognised response to registering live file {}: {:?}", self.ozid(), fnum, msg;
422 Channel, Unexpected)),
423 }
424 }
425
426 fn new_live_pair(&mut self) -> Outcome<(FileNum, u64)> {
427 let fnum_old = self.lpair().fnum;
428 // [3] Ask zbot for next live file number.
429 let resp = Responder::new(Some(self.ozid()));
430 match self.zbot() {
431 None => return Err(err!(
432 "{}: Could not get zbot work channel.", self.ozid();
433 Missing, Data)),
434 Some(zbot) =>
435 res!(zbot.send(OzoneMsg::NextLiveFile(resp.clone()))),
436 }
437 let fnum_new = match resp.recv_timeout(constant::BOT_REQUEST_TIMEOUT) {
438 Err(e) => return Err(err!(e,
439 "While getting next live file info from zbot.";
440 IO, Channel, Read)),
441 Ok(OzoneMsg::UseLiveFile(fnum)) => fnum,
442 Ok(OzoneMsg::Error(e)) => return Err(err!(e,
443 "{}: The zone could not give this writer a new live file.", self.ozid();
444 IO, File, Create)),
445 Ok(msg) => return Err(err!(
446 "Unrecognised new live file request response: {:?}", msg;
447 Bug, Invalid, Input)),
448 };
449
450 // The syncer is handed the new pair before this writer switches to it, and nothing after
451 // the hand-off can fail the switch. Switched first, a writer whose hand-off failed went
452 // on appending to a pair its syncer did not hold: records confirmed durable that no
453 // barrier had covered, in a file its file bot never flagged live, while the file it had
454 // left stayed flagged live for good (2026-09-23). The file being sealed is made durable
455 // by the syncer before it releases anything written to its successor, whatever the
456 // configured policy, so crash loss stays bounded by the live file's tail. This writer
457 // does not wait for that: the disk is the syncer's to wait on, and a writer held by it
458 // would hold every record queued behind the rollover.
459 let lpair = res!(self.zdir().open_live(fnum_new));
460 if let Err(e) = self.hand_pair(&lpair) {
461 self.abandon(lpair);
462 return Err(err!(e,
463 "{}: New live file {} could not be handed to the durability barrier thread, so \
464 this writer stays on file {}.", self.ozid(), fnum_new, fnum_old;
465 IO, File));
466 }
467 self.lpair.close();
468 self.lpair = lpair;
469 let start = self.lpair().dat.size;
470
471 // [5] Tell the fbot for the previous live file of the change and wait for the response.
472 let resp_w3 = Responder::new(Some(self.ozid()));
473 let bots = res!(self.fbots());
474 let (bot, _) = bots.choose_bot(&ChooseBot::ByFile(fnum_old));
475 res!(bot.send(OzoneMsg::CloseOldLiveFileState {
476 fnum_old,
477 fnum_new,
478 new_dat_size: self.lpair.dat.size,
479 new_ind_size: self.lpair.ind.size,
480 resp: resp_w3.clone(),
481 }));
482 // [9] Wait to hear when the new live file is ready to go.
483 match resp_w3.recv_timeout(constant::BOT_REQUEST_TIMEOUT) {
484 Err(e) => return Err(err!(e,
485 "While advising fbot to update live file states.";
486 IO, Channel, Read)),
487 Ok(OzoneMsg::Ok) => (),
488 Ok(OzoneMsg::Error(e)) => return Err(err!(e,
489 "{}: The file bot could not seal live file {} for file {}.",
490 self.ozid(), fnum_old, fnum_new;
491 IO, File)),
492 Ok(msg) => return Err(err!(
493 "Unrecognised response after advising fbot to update live file states: {:?}", msg;
494 Channel)),
495 }
496
497 Ok((fnum_new, start))
498 }
499
500 /// Writes the given byte vector to the current `LivePair` for this zone, either to the data or
501 /// index file. This method starts a new `LivePair` data file when writing the given value
502 /// would cause the current file to exceed the data file size limit. The data to be written is
503 /// not modified in any way.
504 ///
505 /// # Arguments
506 /// * `typ` - `FileType::Data` or `FileType::Index`.
507 /// * `v` - bytes to write to the `LivePair`.
508 ///
509 /// # Errors
510 /// * The data length cannot exceed the maximum data file size.
511 /// * The process did not write all the data to file.
512 /// * The partially written data space could not be recovered.
513 ///
514 /// Returns the starting position in the file for the data sequence.
515 fn write_to_file(
516 &mut self,
517 typ: FileType,
518 v: Vec<&[u8]>,
519 )
520 -> Outcome<u64>
521 {
522 // [2.1] Tally total length of data value.
523 let mut vlen = 0;
524 for vi in &v {
525 vlen += vi.len();
526 }
527 let max_file_len = self.cfg().data_file_max_bytes as usize;
528 if vlen > max_file_len {
529 return Err(err!(
530 "Attempt to store {:?} value of length {} bytes \
531 exceeds the config setting of {}.",
532 typ, vlen, max_file_len;
533 IO, File, Write, Input, TooBig));
534 }
535
536 // [2.2] If the current live file is full, start a new one before anything is written.
537 let mut new_file = false;
538 let mut start = self.lpair().dat.size;
539 let file_len = start as usize;
540 if self.lpair.dat.file.is_none() || (typ == FileType::Data && (vlen + file_len > max_file_len)) {
541 new_file = true;
542 let (_, start2) = res!(self.new_live_pair());
543 start = start2;
544 }
545
546 // [10] Write the data to the live file.
547 match typ {
548 FileType::Data => {
549 let mut bytes_written = 0;
550 for vi in v {
551 // [10.1] Value bytes written here. TODO concat and write once?
552 match self.lpair_mut().dat.file.as_mut() {
553 Some(file) => {
554 match file.write(vi) {
555 Err(e) => {
556 error!(sync_log::stream(), err!(e,
557 "{}: while writing to file, rewinding.", self.ozid();
558 IO, File, Write));
559 break;
560 },
561 Ok(n) => bytes_written += n,
562 }
563 },
564 None => return Err(err!(
565 "{}: The data file should not be None.", self.ozid();
566 Unreachable)),
567 }
568 }
569 if bytes_written < vlen {
570 // [10.2] We have a problem, try and rewind the data file pointer.
571 let msg = format!(
572 "{}: Only {} of {} bytes was written to data file {:?}, but \
573 the process has been aborted with no adverse effect on the \
574 integrity of the database{}",
575 self.ozid(), bytes_written, vlen, self.lpair().dat.path,
576 if new_file { " (although a new file was started)" }
577 else { "." },
578 );
579 // [10.3] Rewind the data file pointer.
580 let len = self.lpair().dat.size - 1;
581 match self.lpair_mut().dat.file.as_mut() {
582 Some(file) => res!(Self::rewind_file_pos(
583 file,
584 len,
585 bytes_written,
586 format!("{} An attempt to recover file space also failed, \
587 again with no impact on database integrity", msg,
588 ))),
589 None => return Err(err!(
590 "{}: The data file should not be None.", self.ozid();
591 Unreachable, Bug)),
592 }
593
594 return Err(err!(
595 "{} The {} bytes of file space were fully recovered.",
596 msg, bytes_written;
597 IO, File, Write));
598 } else {
599 // [10.2] Good write, refresh the cached data file length.
600 self.lpair_mut().dat.size = res!(self.lpair().dat.get_file_len());
601 }
602 Ok(start)
603 },
604 FileType::Index => {
605 let mut bytes_written = 0;
606 for vi in v {
607 // [10.4] Index bytes written here.
608 match self.lpair_mut().ind.file.as_mut() {
609 Some(file) => match file.write(vi) {
610 Err(e) => {
611 error!(sync_log::stream(), err!(e,
612 "{}: while writing to file, rewinding.", self.ozid();
613 IO, File, Write));
614 break;
615 },
616 Ok(n) => bytes_written += n,
617 },
618 None => return Err(err!(
619 "{}: The index file should not be None.", self.ozid();
620 Unreachable, Bug)),
621 }
622 }
623 if bytes_written < vlen {
624 return Err(err!(
625 "{}: Only {} of {} bytes was written to index file {:?}, but \
626 the process has been aborted with no adverse effect on the \
627 integrity of the database{} The corruption will be detected on \
628 next start up, triggering a more laborious scan of the associated \
629 data file and a re-write of the index file.",
630 self.ozid(), bytes_written, vlen, self.lpair().dat.path,
631 if new_file { " (although a new file was started)" }
632 else { "." };
633 IO, File, Write));
634 // Retain the corrupted index data, it will be detected and dealt
635 // with on re-start.
636 }
637 Ok(start)
638 },
639 }
640 }
641
642 fn rewind_file_pos(
643 file: &mut File,
644 orig_pos: u64,
645 bytes_written: usize,
646 msg: String,
647 )
648 -> Outcome<()>
649 {
650 match file.seek(SeekFrom::Start(orig_pos)) {
651 Err(e) => Err(err!(e, "{}.", msg; IO, File, Seek)),
652 Ok(actual_pos) => {
653 if actual_pos != orig_pos {
654 Err(err!(
655 "{}, the file cursor only rewound {} of the required {} bytes.",
656 msg, actual_pos-orig_pos, bytes_written;
657 IO, File, Seek))
658 } else {
659 Ok(())
660 }
661 },
662 }
663 }
664}
665