Oregami
Repositories/oxedyne/fe2o3

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
15use crate::{
16 prelude::*,
17 base::cfg::OzoneConfig,
18 comm::{
19 msg::OzoneMsg,
20 response::Responder,
21 },
22 test::hooks,
23};
24
25use oxedyne_fe2o3_core::channels::Simplex;
26use oxedyne_fe2o3_jdat::id::NumIdDat;
27
28use 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)]
57pub enum SyncPolicy {
58 EveryWrite, // `sync_on_write`
59 EveryN(u32), // `sync_every_n_writes`
60 Interval(Duration), // `sync_interval_ms`
61 Never,
62}
63
64impl 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.
83pub 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.
100pub 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
111impl<
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
178struct 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
196impl<
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
426impl<
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}