oxedyne/fe2o3/fe2o3_o3db_sync/src/bots/worker/syncer.rs
16.1 KiB, 9 runs
created by r1870400018:61241, 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 | //! A writer's durability barrier, run on a thread of its own. |
| 2 | //! |
| 3 | //! A writer bot appends a record and goes straight on to the next; its syncer makes each record |
| 4 | //! durable as the store's sync policy asks and only then releases it to its cache bot, which makes |
| 5 | //! it readable and gives the caller its final answer. The barrier used to run inside the writer, |
| 6 | //! so one fsync held up by the rest of the machine's traffic -- eleven seconds was measured on |
| 7 | //! 2026-09-23 -- held every record queued behind it past the caller's six-second deadline, and each |
| 8 | //! was reported as failed although each went on to land. |
| 9 | //! |
| 10 | //! Records are released in the order they were written, so a record reaches its cache bot, and |
| 11 | //! its accounting reaches the file bots, only after every record written before it. A file bot |
| 12 | //! collects a sealed file only once the file's accounted size has caught up with its size on |
| 13 | //! disk, so holding records here delays a collection and never races one. |
| 14 | |
| 15 | use crate::{ |
| 16 | prelude::*, |
| 17 | base::cfg::OzoneConfig, |
| 18 | comm::{ |
| 19 | msg::OzoneMsg, |
| 20 | response::Responder, |
| 21 | }, |
| 22 | test::hooks, |
| 23 | }; |
| 24 | |
| 25 | use oxedyne_fe2o3_core::channels::Simplex; |
| 26 | use oxedyne_fe2o3_jdat::id::NumIdDat; |
| 27 | |
| 28 | use std::{ |
| 29 | collections::VecDeque, |
| 30 | fs::File, |
| 31 | sync::{ |
| 32 | Arc, |
| 33 | atomic::{ |
| 34 | AtomicBool, |
| 35 | Ordering, |
| 36 | }, |
| 37 | mpsc::{ |
| 38 | self, |
| 39 | Receiver, |
| 40 | RecvTimeoutError, |
| 41 | Sender, |
| 42 | }, |
| 43 | }, |
| 44 | thread::{ |
| 45 | self, |
| 46 | JoinHandle, |
| 47 | }, |
| 48 | time::{ |
| 49 | Duration, |
| 50 | Instant, |
| 51 | }, |
| 52 | }; |
| 53 | |
| 54 | |
| 55 | /// When a record has to be durable before it is released. |
| 56 | #[derive(Clone, Copy, Debug, Eq, PartialEq)] |
| 57 | pub enum SyncPolicy { |
| 58 | EveryWrite, // `sync_on_write` |
| 59 | EveryN(u32), // `sync_every_n_writes` |
| 60 | Interval(Duration), // `sync_interval_ms` |
| 61 | Never, |
| 62 | } |
| 63 | |
| 64 | impl SyncPolicy { |
| 65 | /// The policy a configuration asks for, the strongest first where it asks for several. |
| 66 | pub fn of(cfg: &OzoneConfig) -> Self { |
| 67 | if cfg.sync_on_write { |
| 68 | Self::EveryWrite |
| 69 | } else if cfg.sync_every_n_writes > 0 { |
| 70 | Self::EveryN(cfg.sync_every_n_writes) |
| 71 | } else if cfg.sync_interval_ms > 0 { |
| 72 | Self::Interval(Duration::from_millis(cfg.sync_interval_ms)) |
| 73 | } else { |
| 74 | Self::Never |
| 75 | } |
| 76 | } |
| 77 | } |
| 78 | |
| 79 | /// What a writer hands its syncer, in the order it wrote. A `Pair` names the live data and index |
| 80 | /// files every later record is appended to; the pair it replaces is sealed, made durable before |
| 81 | /// any record written to the new one is released. A `Record` is one appended record with the |
| 82 | /// insert that releases it to its cache bot. |
| 83 | pub enum Handed< |
| 84 | const UIDL: usize, |
| 85 | UID: NumIdDat<UIDL>, |
| 86 | ENC: Encrypter, |
| 87 | KH: Hasher, |
| 88 | > { |
| 89 | Pair(File, File), // data, index |
| 90 | Record { |
| 91 | cbot: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 92 | insert: OzoneMsg<UIDL, UID, ENC, KH>, |
| 93 | resp: Responder<UIDL, UID, ENC, KH>, |
| 94 | policy: SyncPolicy, |
| 95 | }, |
| 96 | Finish, |
| 97 | } |
| 98 | |
| 99 | /// A writer's handle on its syncer. |
| 100 | pub struct Syncer< |
| 101 | const UIDL: usize, |
| 102 | UID: NumIdDat<UIDL>, |
| 103 | ENC: Encrypter, |
| 104 | KH: Hasher, |
| 105 | > { |
| 106 | tx: Sender<Handed<UIDL, UID, ENC, KH>>, |
| 107 | handle: Option<JoinHandle<()>>, |
| 108 | running: Arc<AtomicBool>, // cleared by the thread's barrier as it goes, however it goes |
| 109 | } |
| 110 | |
| 111 | impl< |
| 112 | const UIDL: usize, |
| 113 | UID: NumIdDat<UIDL> + 'static, |
| 114 | ENC: Encrypter + 'static, |
| 115 | KH: Hasher + 'static, |
| 116 | > |
| 117 | Syncer<UIDL, UID, ENC, KH> |
| 118 | { |
| 119 | pub fn start( |
| 120 | label: String, |
| 121 | log_stream_id: String, |
| 122 | ) |
| 123 | -> Outcome<Self> |
| 124 | { |
| 125 | // A plain channel rather than a `Simplex`, whose two ends travel together: a writer must |
| 126 | // learn from its own send that the syncer has gone, rather than queue records that nothing |
| 127 | // will ever release. |
| 128 | let (tx, rx) = mpsc::channel(); |
| 129 | let running = Arc::new(AtomicBool::new(true)); |
| 130 | let cleared = running.clone(); |
| 131 | let builder = thread::Builder::new() |
| 132 | .name(fmt!("{}.sync", label)) |
| 133 | .stack_size(constant::STACK_SIZE); |
| 134 | let handle = res!(builder.spawn(move || { |
| 135 | sync_log::set_stream(log_stream_id); |
| 136 | Barrier::new(label, rx, cleared).run(); |
| 137 | })); |
| 138 | Ok(Self { |
| 139 | tx, |
| 140 | handle: Some(handle), |
| 141 | running, |
| 142 | }) |
| 143 | } |
| 144 | |
| 145 | /// Is the syncer still taking records? Once it is not, nothing more written through its |
| 146 | /// writer can be confirmed. |
| 147 | pub fn is_running(&self) -> bool { |
| 148 | self.running.load(Ordering::SeqCst) |
| 149 | } |
| 150 | |
| 151 | pub fn hand(&self, item: Handed<UIDL, UID, ENC, KH>) -> Outcome<()> { |
| 152 | match self.tx.send(item) { |
| 153 | Ok(()) => Ok(()), |
| 154 | Err(_) => Err(err!( |
| 155 | "The durability barrier thread has stopped, so no write from here on can be \ |
| 156 | confirmed."; |
| 157 | Thread, Missing)), |
| 158 | } |
| 159 | } |
| 160 | |
| 161 | /// Releases everything handed over, with a last barrier where the policy owes one, and ends |
| 162 | /// the thread. |
| 163 | pub fn finish(&mut self) -> Outcome<()> { |
| 164 | // A syncer that has already stopped cannot be sent this, and the join says why. |
| 165 | let _ = self.tx.send(Handed::Finish); |
| 166 | match self.handle.take() { |
| 167 | None => Ok(()), |
| 168 | Some(handle) => match handle.join() { |
| 169 | Ok(()) => Ok(()), |
| 170 | Err(_) => Err(err!( |
| 171 | "The durability barrier thread panicked."; |
| 172 | Thread, Panic)), |
| 173 | }, |
| 174 | } |
| 175 | } |
| 176 | } |
| 177 | |
| 178 | struct Barrier< |
| 179 | const UIDL: usize, |
| 180 | UID: NumIdDat<UIDL>, |
| 181 | ENC: Encrypter, |
| 182 | KH: Hasher, |
| 183 | > { |
| 184 | label: String, |
| 185 | rx: Receiver<Handed<UIDL, UID, ENC, KH>>, |
| 186 | queue: VecDeque<Handed<UIDL, UID, ENC, KH>>, |
| 187 | pair: Option<(File, File)>, |
| 188 | policy: SyncPolicy, // the latest record's |
| 189 | dirty: bool, // the pair holds records no barrier has covered |
| 190 | since: u32, // records released since the last barrier |
| 191 | last: Option<Instant>, // when the last barrier completed |
| 192 | failed: Option<Instant>, // when the last barrier failed, until one completes |
| 193 | running: Arc<AtomicBool>, // shared with the writer's handle |
| 194 | } |
| 195 | |
| 196 | impl< |
| 197 | const UIDL: usize, |
| 198 | UID: NumIdDat<UIDL> + 'static, |
| 199 | ENC: Encrypter + 'static, |
| 200 | KH: Hasher + 'static, |
| 201 | > |
| 202 | Barrier<UIDL, UID, ENC, KH> |
| 203 | { |
| 204 | fn new( |
| 205 | label: String, |
| 206 | rx: Receiver<Handed<UIDL, UID, ENC, KH>>, |
| 207 | running: Arc<AtomicBool>, |
| 208 | ) |
| 209 | -> Self |
| 210 | { |
| 211 | Self { |
| 212 | label, |
| 213 | rx, |
| 214 | queue: VecDeque::new(), |
| 215 | pair: None, |
| 216 | policy: SyncPolicy::Never, |
| 217 | dirty: false, |
| 218 | since: 0, |
| 219 | last: None, |
| 220 | failed: None, |
| 221 | running, |
| 222 | } |
| 223 | } |
| 224 | |
| 225 | fn run(mut self) { |
| 226 | loop { |
| 227 | if self.queue.is_empty() { |
| 228 | // The interval policy owes the disk a barrier on its period whether or not another |
| 229 | // record arrives, so the writes that end a burst are made durable too. |
| 230 | let item = match self.owed_in() { |
| 231 | Some(wait) => match self.rx.recv_timeout(wait) { |
| 232 | Ok(item) => Some(item), |
| 233 | Err(RecvTimeoutError::Timeout) => { |
| 234 | self.barrier_or_log(); |
| 235 | continue; |
| 236 | }, |
| 237 | Err(RecvTimeoutError::Disconnected) => None, |
| 238 | }, |
| 239 | None => self.rx.recv().ok(), |
| 240 | }; |
| 241 | match item { |
| 242 | Some(item) => self.queue.push_back(item), |
| 243 | // The writer has gone without a word, so nothing more is coming. |
| 244 | None => return self.end(), |
| 245 | } |
| 246 | } |
| 247 | // Take whatever else is waiting, so that one barrier can cover all of it. |
| 248 | while let Ok(item) = self.rx.try_recv() { |
| 249 | self.queue.push_back(item); |
| 250 | } |
| 251 | while let Some(item) = self.queue.pop_front() { |
| 252 | match item { |
| 253 | Handed::Pair(dat, ind) => { |
| 254 | // Sealing is unconditional, whatever the policy: a file that is no |
| 255 | // longer live is durable before anything written after it is released. |
| 256 | if self.dirty { |
| 257 | self.barrier_or_log(); |
| 258 | } |
| 259 | self.pair = Some((dat, ind)); |
| 260 | }, |
| 261 | Handed::Record { cbot, insert, resp, policy } => |
| 262 | self.records(cbot, insert, resp, policy), |
| 263 | Handed::Finish => return self.end(), |
| 264 | } |
| 265 | } |
| 266 | // A syncer stopping as one that panicked would, with no last barrier. |
| 267 | if hooks::syncer_stops() { |
| 268 | return; |
| 269 | } |
| 270 | } |
| 271 | } |
| 272 | |
| 273 | /// Releases one record, and with it the run queued behind it where every record waits on a |
| 274 | /// barrier of its own: one barrier covers every record appended before it begins. That is |
| 275 | /// always so under `sync_on_write`, and under the interval and every-n policies while the last |
| 276 | /// barrier has failed. Batched only under the first, a disk failing its syncs slowly took one |
| 277 | /// barrier per record under the others, one after another, and with writes arriving faster |
| 278 | /// than it failed the queue grew without bound (2026-09-23, every-n 2026-09-24). |
| 279 | fn records( |
| 280 | &mut self, |
| 281 | cbot: Simplex<OzoneMsg<UIDL, UID, ENC, KH>>, |
| 282 | insert: OzoneMsg<UIDL, UID, ENC, KH>, |
| 283 | resp: Responder<UIDL, UID, ENC, KH>, |
| 284 | policy: SyncPolicy, |
| 285 | ) { |
| 286 | let mut run = vec![(cbot, insert, resp)]; |
| 287 | let alone = match policy { |
| 288 | SyncPolicy::EveryWrite => true, |
| 289 | SyncPolicy::EveryN(_) | SyncPolicy::Interval(_) => self.failed.is_some(), |
| 290 | SyncPolicy::Never => false, |
| 291 | }; |
| 292 | if alone { |
| 293 | while let Some(Handed::Record { policy: next, .. }) = self.queue.front() { |
| 294 | if *next != policy { |
| 295 | break; |
| 296 | } |
| 297 | if let Some(Handed::Record { cbot, insert, resp, .. }) = self.queue.pop_front() { |
| 298 | run.push((cbot, insert, resp)); |
| 299 | } |
| 300 | } |
| 301 | } |
| 302 | self.policy = policy; |
| 303 | self.dirty = true; |
| 304 | self.since = self.since.saturating_add(run.len() as u32); |
| 305 | // While the disk is failing a record waits on a barrier of its own, so that its caller |
| 306 | // hears so rather than being answered as if its period will make it durable. |
| 307 | let due = match policy { |
| 308 | SyncPolicy::EveryWrite => true, |
| 309 | SyncPolicy::EveryN(n) => self.since >= n, |
| 310 | SyncPolicy::Interval(p) => self.failed.is_some() || match self.last { |
| 311 | None => true, |
| 312 | Some(last) => last.elapsed() >= p, |
| 313 | }, |
| 314 | SyncPolicy::Never => false, |
| 315 | }; |
| 316 | let synced = if due { self.barrier() } else { Ok(()) }; |
| 317 | for (cbot, insert, resp) in run { |
| 318 | // The caller hears when the barrier failed, from the cache bot once the record is |
| 319 | // readable, as it hears of one that did not: told from here first, it could read its |
| 320 | // write back and not find it. The record is released all the same: it is in the |
| 321 | // files whatever the barrier said, and holding it back would leave the cache |
| 322 | // disagreeing with them and its bytes accounted to no one, where no collection could |
| 323 | // ever reclaim them. |
| 324 | let insert = match (&synced, insert) { |
| 325 | (Ok(()), insert) => insert, |
| 326 | (Err(e), OzoneMsg::Insert(k, v, c, f, i, m, r, _)) => |
| 327 | OzoneMsg::Insert(k, v, c, f, i, m, r, Some(Self::written(e.clone()))), |
| 328 | // Not an insert, so nothing the cache bot would answer with: told from here. |
| 329 | (Err(e), other) => { |
| 330 | Self::tell(&resp, e.clone()); |
| 331 | other |
| 332 | }, |
| 333 | }; |
| 334 | if let Err(e) = cbot.send(insert) { |
| 335 | let e = err!(e, |
| 336 | "{}: A written record could not be released to its cache bot.", self.label; |
| 337 | Channel, Write); |
| 338 | error!(sync_log::stream(), e.clone()); |
| 339 | Self::tell(&resp, e); |
| 340 | } |
| 341 | } |
| 342 | } |
| 343 | |
| 344 | /// Forces the current pair to stable storage. |
| 345 | fn barrier(&mut self) -> Outcome<()> { |
| 346 | if let Err(e) = self.sync_pair() { |
| 347 | self.failed = Some(Instant::now()); |
| 348 | return Err(e); |
| 349 | } |
| 350 | self.dirty = false; |
| 351 | self.since = 0; |
| 352 | self.last = Some(Instant::now()); |
| 353 | self.failed = None; |
| 354 | Ok(()) |
| 355 | } |
| 356 | |
| 357 | fn sync_pair(&self) -> Outcome<()> { |
| 358 | if let Some((dat, ind)) = &self.pair { |
| 359 | hooks::barrier_delay(); |
| 360 | if let Err(e) = Self::sync(dat) { |
| 361 | return Err(err!(e, |
| 362 | "{}: sync_data on the live data file failed, so records written to it are not \ |
| 363 | confirmed durable.", self.label; |
| 364 | IO, File, Write)); |
| 365 | } |
| 366 | if let Err(e) = Self::sync(ind) { |
| 367 | return Err(err!(e, |
| 368 | "{}: sync_data on the live index file failed, so records written to it are \ |
| 369 | not confirmed durable.", self.label; |
| 370 | IO, File, Write)); |
| 371 | } |
| 372 | } |
| 373 | Ok(()) |
| 374 | } |
| 375 | |
| 376 | /// `File::sync_data`, or the failure a test has asked for (`test::hooks`). A barrier stops at |
| 377 | /// its first failed sync, so a barrier the hook fails is counted once. |
| 378 | fn sync(file: &File) -> std::io::Result<()> { |
| 379 | if hooks::sync_fails() { |
| 380 | return Err(std::io::Error::other( |
| 381 | "the disk failed the sync (test::hooks::set_barrier_failure)")); |
| 382 | } |
| 383 | file.sync_data() |
| 384 | } |
| 385 | |
| 386 | /// A barrier nobody is waiting on, whose failure can only be logged. |
| 387 | fn barrier_or_log(&mut self) { |
| 388 | if let Err(e) = self.barrier() { |
| 389 | error!(sync_log::stream(), e); |
| 390 | } |
| 391 | } |
| 392 | |
| 393 | /// How long until the interval policy owes a barrier, while it owes one. A failed barrier is |
| 394 | /// retried a period after it failed: counted from the last one that completed, a disk refusing |
| 395 | /// every sync had this thread retry it as fast as the disk could refuse, a million times in |
| 396 | /// three seconds with an error logged for each (2026-09-23). |
| 397 | fn owed_in(&self) -> Option<Duration> { |
| 398 | match (self.dirty, self.policy, self.failed.or(self.last)) { |
| 399 | (true, SyncPolicy::Interval(p), Some(t)) => Some(p.saturating_sub(t.elapsed())), |
| 400 | (true, SyncPolicy::Interval(_), None) => Some(Duration::ZERO), |
| 401 | _ => None, |
| 402 | } |
| 403 | } |
| 404 | |
| 405 | fn end(&mut self) { |
| 406 | if self.dirty && self.policy != SyncPolicy::Never { |
| 407 | self.barrier_or_log(); |
| 408 | } |
| 409 | } |
| 410 | |
| 411 | fn tell(resp: &Responder<UIDL, UID, ENC, KH>, e: Error<ErrTag>) { |
| 412 | // A caller that has already given up has dropped its end, and there is no one else to |
| 413 | // tell. A responder with no channel is an internal write nobody waits on. |
| 414 | if resp.is_some() { |
| 415 | let _ = resp.send(OzoneMsg::Error(Self::written(e))); |
| 416 | } |
| 417 | } |
| 418 | |
| 419 | /// What a written record's caller is told went wrong after the write: tagged `Unconfirmed`, |
| 420 | /// so that it can be told from a write that never landed without reading the words. |
| 421 | fn written(e: Error<ErrTag>) -> Error<ErrTag> { |
| 422 | err!(e, "The record is written, but not confirmed durable."; Write, Unconfirmed) |
| 423 | } |
| 424 | } |
| 425 | |
| 426 | impl< |
| 427 | const UIDL: usize, |
| 428 | UID: NumIdDat<UIDL>, |
| 429 | ENC: Encrypter, |
| 430 | KH: Hasher, |
| 431 | > |
| 432 | Drop for Barrier<UIDL, UID, ENC, KH> |
| 433 | { |
| 434 | fn drop(&mut self) { |
| 435 | // However the thread ends, a panic included, and before the receiver is dropped: a writer |
| 436 | // never finds the channel gone while the flag still says the syncer is running. |
| 437 | self.running.store(false, Ordering::SeqCst); |
| 438 | } |
| 439 | } |