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
| 1 | use 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 | |
| 26 | use oxedyne_fe2o3_iop_db::api::Meta; |
| 27 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 28 | |
| 29 | use 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`. |
| 40 | pub 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 | |
| 68 | impl< |
| 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 | |
| 81 | impl< |
| 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 | |
| 94 | impl< |
| 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 | |
| 192 | impl< |
| 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 |