1use ahash::RandomState;
2use bytes::Bytes;
3use crossbeam_channel::{bounded, Receiver, Sender};
4use crossbeam_utils::CachePadded;
5use parking_lot::{Mutex, RwLock};
6use std::collections::VecDeque;
7use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicUsize, Ordering};
8use std::sync::Arc;
9use std::thread::{self, JoinHandle};
10use std::time::Duration;
11
12use crate::constants::*;
13use crate::core::record::Record;
14use crate::error::{FeoxError, Result};
15use crate::stats::Statistics;
16use crate::storage::allocation_journal::ALLOCATION_JOURNAL_MAX_ENTRIES;
17use crate::storage::format::{get_format_ref, sector_holds_record, RecordFormat};
18use crate::storage::free_space::FreeSpaceManager;
19use crate::storage::io::DiskIO;
20use crate::storage::seq_token::{stamp_seq_token, SEQ_TOKEN_MIN_VERSION};
21use crate::test_hooks::{fail_at, RECORD_WRITE};
22
23#[repr(align(64))] pub struct ShardedWriteBuffer {
27 buffer: Mutex<VecDeque<WriteEntry>>,
29
30 count: AtomicUsize,
32
33 size: AtomicUsize,
35}
36
37pub struct WriteEntry {
39 pub op: Operation,
40 pub record: Arc<Record>,
41 pub work_status: AtomicU32,
42 pub retry_count: AtomicU32,
43}
44
45impl WriteEntry {
46 fn new(op: Operation, record: Arc<Record>) -> Self {
47 Self {
48 op,
49 record,
50 work_status: AtomicU32::new(0),
51 retry_count: AtomicU32::new(0),
52 }
53 }
54}
55
56const DELETE_MARKER_DURABLE: u32 = 1;
57const RESERVATION_DIRTY: u32 = 1 << 31;
58const RESERVATION_QUARANTINED: u32 = 1 << 30;
59const RESERVATION_FLAGS: u32 = RESERVATION_DIRTY | RESERVATION_QUARANTINED;
60const FINAL_FLUSH_RETRY_LIMIT: usize = 1_024;
61
62#[inline]
63fn reserved_sector(entry: &WriteEntry) -> Option<u64> {
64 let sector = entry.work_status.load(Ordering::Acquire) & !RESERVATION_FLAGS;
65 (sector != 0).then_some(sector as u64)
66}
67
68#[inline]
69fn reserve_sector(entry: &WriteEntry, sector: u64) {
70 let sector = u32::try_from(sector).expect("device sector exceeds supported range");
71 debug_assert_eq!(sector & RESERVATION_FLAGS, 0);
72 entry.work_status.store(sector, Ordering::Release);
73}
74
75#[inline]
76fn mark_reservation_dirty(entry: &WriteEntry) {
77 entry
78 .work_status
79 .fetch_or(RESERVATION_DIRTY, Ordering::AcqRel);
80}
81
82#[inline]
83fn mark_reservation_clean(entry: &WriteEntry) {
84 entry
85 .work_status
86 .fetch_and(!RESERVATION_DIRTY, Ordering::AcqRel);
87}
88
89#[inline]
90fn reservation_is_dirty(entry: &WriteEntry) -> bool {
91 entry.work_status.load(Ordering::Acquire) & RESERVATION_DIRTY != 0
92}
93
94#[inline]
95fn quarantine_reservation(entry: &WriteEntry) {
96 entry
97 .work_status
98 .fetch_or(RESERVATION_QUARANTINED, Ordering::AcqRel);
99}
100
101#[inline]
102fn reservation_is_quarantined(entry: &WriteEntry) -> bool {
103 entry.work_status.load(Ordering::Acquire) & RESERVATION_QUARANTINED != 0
104}
105
106#[inline]
107fn clear_reserved_sector(entry: &WriteEntry) {
108 entry.work_status.store(0, Ordering::Release);
109}
110
111struct PreparedWrite {
112 data: Vec<u8>,
113 sectors_needed: usize,
114 entry: WriteEntry,
115 sector: Option<u64>,
116}
117
118struct BatchOutcome {
119 result: Result<()>,
120 retries: Vec<WriteEntry>,
121}
122
123struct BatchFailure {
124 error: FeoxError,
125 clear_journal: bool,
126 indeterminate: bool,
127}
128
129struct RetirementQueue {
130 pending: Mutex<Vec<WriteEntry>>,
131 flush: Mutex<()>,
132 released_sectors: AtomicU64,
133}
134
135impl RetirementQueue {
136 fn new() -> Self {
137 Self {
138 pending: Mutex::new(Vec::new()),
139 flush: Mutex::new(()),
140 released_sectors: AtomicU64::new(0),
141 }
142 }
143}
144
145pub struct WriteBuffer {
147 sharded_buffers: Arc<Vec<CachePadded<ShardedWriteBuffer>>>,
149
150 disk_io: Arc<RwLock<DiskIO>>,
152
153 free_space: Arc<RwLock<FreeSpaceManager>>,
155
156 worker_channels: Vec<Sender<FlushRequest>>,
158
159 worker_handles: Mutex<Vec<JoinHandle<()>>>,
161
162 periodic_flush_handle: Mutex<Option<JoinHandle<()>>>,
164
165 shutdown: Arc<AtomicBool>,
167
168 stats: Arc<Statistics>,
170
171 shard_hasher: RandomState,
173
174 retirement_queue: Arc<RetirementQueue>,
175
176 format_version: u32,
178
179 fault_scope: usize,
180}
181
182#[derive(Debug)]
183struct FlushRequest {
184 response: Option<Sender<Result<bool>>>,
185 defer_retirements: bool,
186}
187
188struct WorkerContext {
189 worker_id: usize,
190 worker_count: usize,
191 disk_io: Arc<RwLock<DiskIO>>,
192 free_space: Arc<RwLock<FreeSpaceManager>>,
193 sharded_buffers: Arc<Vec<CachePadded<ShardedWriteBuffer>>>,
194 shutdown: Arc<AtomicBool>,
195 stats: Arc<Statistics>,
196 retirement_queue: Arc<RetirementQueue>,
197 format_version: u32,
198 fault_scope: usize,
199}
200
201impl ShardedWriteBuffer {
202 fn new(_shard_id: usize) -> Self {
203 Self {
204 buffer: Mutex::new(VecDeque::new()),
205 count: AtomicUsize::new(0),
206 size: AtomicUsize::new(0),
207 }
208 }
209
210 fn add_entries<const N: usize>(
211 &self,
212 entries: [WriteEntry; N],
213 shutdown: &AtomicBool,
214 ) -> Result<()> {
215 let mut buffer = self.buffer.lock();
216 if shutdown.load(Ordering::Acquire) {
217 return Err(FeoxError::ShuttingDown);
218 }
219
220 let entry_size = entries
221 .iter()
222 .map(|entry| entry.record.calculate_size())
223 .sum::<usize>();
224 buffer.extend(entries);
225
226 self.count.fetch_add(N, Ordering::Relaxed);
227 self.size.fetch_add(entry_size, Ordering::Relaxed);
228 Ok(())
229 }
230
231 fn drain_entries(&self) -> Vec<WriteEntry> {
232 let mut buffer = self.buffer.lock();
233 let entries: Vec<_> = buffer.drain(..).collect();
234
235 self.count.store(0, Ordering::Relaxed);
236 self.size.store(0, Ordering::Relaxed);
237
238 entries
239 }
240
241 fn requeue_entries(&self, entries: Vec<WriteEntry>, stats: &Arc<Statistics>, failed: bool) {
242 if entries.is_empty() {
243 return;
244 }
245
246 let count = entries.len();
247 let size = entries
248 .iter()
249 .map(|entry| entry.record.calculate_size())
250 .sum();
251 let mut buffer = self.buffer.lock();
252 for entry in entries.into_iter().rev() {
253 if failed {
254 let retries = entry.retry_count.fetch_add(1, Ordering::Relaxed) + 1;
255 if retries == WRITE_ENTRY_RETRY_ALARM {
256 stats.record_write_entry_stuck();
257 eprintln!(
258 "feox: write entry for a {} byte key has been retried {} times",
259 entry.record.key.len(),
260 retries
261 );
262 }
263 }
264 buffer.push_front(entry);
265 }
266 self.count.fetch_add(count, Ordering::Relaxed);
267 self.size.fetch_add(size, Ordering::Relaxed);
268 }
269
270 fn is_full(&self) -> bool {
271 self.count.load(Ordering::Relaxed) >= WRITE_BUFFER_SIZE
272 || self.size.load(Ordering::Relaxed) >= FEOX_WRITE_BUFFER_SIZE
273 }
274}
275
276impl WriteBuffer {
277 pub fn new(
278 disk_io: Arc<RwLock<DiskIO>>,
279 free_space: Arc<RwLock<FreeSpaceManager>>,
280 stats: Arc<Statistics>,
281 format_version: u32,
282 ) -> Self {
283 let num_shards = (num_cpus::get() / 2).max(1);
285
286 let sharded_buffers = Arc::new(
287 (0..num_shards)
288 .map(|shard_id| CachePadded::new(ShardedWriteBuffer::new(shard_id)))
289 .collect(),
290 );
291
292 Self {
293 sharded_buffers,
294 disk_io,
295 free_space,
296 worker_channels: Vec::new(),
297 worker_handles: Mutex::new(Vec::new()),
298 periodic_flush_handle: Mutex::new(None),
299 shutdown: Arc::new(AtomicBool::new(false)),
300 stats,
301 shard_hasher: RandomState::new(),
302 retirement_queue: Arc::new(RetirementQueue::new()),
303 format_version,
304 fault_scope: crate::test_hooks::new_fault_scope(),
305 }
306 }
307
308 #[cfg(test)]
309 pub(crate) fn fault_scope(&self) -> usize {
310 self.fault_scope
311 }
312
313 pub fn add_write(
315 &self,
316 op: Operation,
317 record: Arc<Record>,
318 _old_value_len: usize,
319 ) -> Result<()> {
320 let shard_id = self.get_shard_id(&record.key);
321 let buffer = &self.sharded_buffers[shard_id];
322 buffer.add_entries([WriteEntry::new(op, record)], &self.shutdown)?;
323 self.stats.record_write_buffered();
324 self.trigger_flush(shard_id, buffer);
325 Ok(())
326 }
327
328 pub(crate) fn add_replacement(&self, record: Arc<Record>, replaced: Arc<Record>) -> Result<()> {
329 debug_assert_eq!(record.key, replaced.key);
330 let shard_id = self.get_shard_id(&record.key);
331 let buffer = &self.sharded_buffers[shard_id];
332 buffer.add_entries(
333 [
334 WriteEntry::new(Operation::Update, record),
335 WriteEntry::new(Operation::Delete, replaced),
336 ],
337 &self.shutdown,
338 )?;
339 self.stats.record_writes_buffered(2);
340 self.trigger_flush(shard_id, buffer);
341 Ok(())
342 }
343
344 fn trigger_flush(&self, shard_id: usize, buffer: &ShardedWriteBuffer) {
345 if buffer.is_full() && !self.worker_channels.is_empty() {
346 let worker_id = shard_id % self.worker_channels.len();
347 let req = FlushRequest {
348 response: None,
349 defer_retirements: false,
350 };
351 let _ = self.worker_channels[worker_id].try_send(req);
352 }
353 }
354
355 pub fn start_workers(&mut self, num_workers: usize) {
357 let num_shards = self.sharded_buffers.len();
359 let actual_workers = num_workers.clamp(1, num_shards);
360
361 let mut receivers = Vec::new();
363 for _ in 0..actual_workers {
364 let (tx, rx) = bounded(2);
365 self.worker_channels.push(tx);
366 receivers.push(rx);
367 }
368
369 for (worker_id, flush_rx) in receivers.into_iter().enumerate() {
371 let ctx = WorkerContext {
372 worker_id,
373 worker_count: actual_workers,
374 disk_io: self.disk_io.clone(),
375 free_space: self.free_space.clone(),
376 sharded_buffers: self.sharded_buffers.clone(),
377 shutdown: self.shutdown.clone(),
378 stats: self.stats.clone(),
379 retirement_queue: self.retirement_queue.clone(),
380 format_version: self.format_version,
381 fault_scope: self.fault_scope,
382 };
383
384 let handle = thread::spawn(move || {
385 write_buffer_worker(ctx, flush_rx);
386 });
387
388 self.worker_handles.get_mut().push(handle);
389 }
390
391 let worker_channels = self.worker_channels.clone();
393 let shutdown = self.shutdown.clone();
394 let sharded_buffers = self.sharded_buffers.clone();
395 let retirement_queue = Arc::clone(&self.retirement_queue);
396
397 let periodic_handle = thread::spawn(move || {
398 let interval = WRITE_BUFFER_FLUSH_INTERVAL;
399
400 while !shutdown.load(Ordering::Acquire) {
401 thread::sleep(interval);
402
403 let retirements_pending = !retirement_queue.pending.lock().is_empty();
404 for (worker_id, channel) in worker_channels.iter().enumerate() {
405 let pending = (worker_id..sharded_buffers.len())
406 .step_by(worker_channels.len())
407 .any(|shard_id| {
408 sharded_buffers[shard_id].count.load(Ordering::Relaxed) > 0
409 });
410 if pending || (worker_id == 0 && retirements_pending) {
411 let _ = channel.try_send(FlushRequest {
412 response: None,
413 defer_retirements: false,
414 });
415 }
416 }
417 }
418 });
419
420 *self.periodic_flush_handle.get_mut() = Some(periodic_handle);
421 }
422
423 pub fn force_flush(&self) -> Result<()> {
425 let mut pending_workers: Vec<_> = (0..self.worker_channels.len()).collect();
426 let mut retry_delay_us = 50;
427 let format = get_format_ref(self.format_version);
428
429 loop {
430 let mut responses = Vec::with_capacity(pending_workers.len());
431 for worker_id in pending_workers.drain(..) {
432 let (tx, rx) = bounded(1);
433 self.worker_channels[worker_id]
434 .send(FlushRequest {
435 response: Some(tx),
436 defer_retirements: true,
437 })
438 .map_err(|_| FeoxError::ChannelError)?;
439 responses.push((worker_id, rx));
440 }
441
442 let mut first_error = None;
443 for (worker_id, rx) in responses {
444 match rx.recv() {
445 Ok(Ok(true)) => pending_workers.push(worker_id),
446 Ok(Ok(false)) => {}
447 Ok(Err(error)) => {
448 if first_error.is_none() {
449 first_error = Some(error);
450 }
451 }
452 Err(_) => {
453 if first_error.is_none() {
454 first_error = Some(FeoxError::ChannelError);
455 }
456 }
457 }
458 }
459
460 if let Some(error) = first_error {
461 return Err(error);
462 }
463 if flush_pending_deletions(
464 &self.retirement_queue,
465 &self.disk_io,
466 &self.free_space,
467 &self.stats,
468 format,
469 )? {
470 pending_workers = (0..self.worker_channels.len()).collect();
471 }
472 if pending_workers.is_empty() {
473 return Ok(());
474 }
475 if !pending_workers.is_empty() {
476 thread::sleep(Duration::from_micros(retry_delay_us));
477 retry_delay_us = (retry_delay_us * 2).min(1_000);
478 }
479 }
480 }
481
482 pub fn initiate_shutdown(&self) {
484 self.shutdown.store(true, Ordering::Release);
485
486 }
489
490 pub fn complete_shutdown(&mut self) {
492 self.finish_shutdown();
493 }
494
495 pub(crate) fn finish_shutdown(&self) {
496 self.shutdown.store(true, Ordering::Release);
498
499 if let Some(handle) = self.periodic_flush_handle.lock().take() {
500 let _ = handle.join();
501 }
502
503 let handles = std::mem::take(&mut *self.worker_handles.lock());
505 for handle in handles {
506 let _ = handle.join();
507 }
508
509 }
512
513 pub fn shutdown(&mut self) {
515 self.complete_shutdown();
516 }
517
518 #[inline]
519 fn get_shard_id(&self, key: &[u8]) -> usize {
520 self.shard_hasher.hash_one(key) as usize % self.sharded_buffers.len()
521 }
522}
523
524fn write_buffer_worker(ctx: WorkerContext, flush_rx: Receiver<FlushRequest>) {
526 let format = get_format_ref(ctx.format_version);
527
528 loop {
529 if ctx.shutdown.load(Ordering::Acquire) {
530 break;
531 }
532
533 let req = match flush_rx.recv_timeout(Duration::from_millis(500)) {
535 Ok(req) => req,
536 Err(crossbeam_channel::RecvTimeoutError::Timeout) => {
537 continue;
538 }
539 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => {
540 break;
541 }
542 };
543
544 let result = flush_worker_shards(&ctx, format, !req.defer_retirements);
545 if let Some(tx) = req.response {
546 let _ = tx.send(result);
547 }
548 }
549
550 if ctx.shutdown.load(Ordering::Acquire) {
551 let mut retry_delay_us = 50;
552 let mut retries = 0;
553 loop {
554 match flush_worker_shards(&ctx, format, true) {
555 Ok(false) => break,
556 Ok(true) => {
557 retries += 1;
558 if retries == FINAL_FLUSH_RETRY_LIMIT {
559 eprintln!("feox: final write-buffer flush left pending retirements");
560 break;
561 }
562 thread::sleep(Duration::from_micros(retry_delay_us));
563 retry_delay_us = (retry_delay_us * 2).min(1_000);
564 }
565 Err(error @ FeoxError::IndeterminateWrite(_)) => {
566 eprintln!("feox: final write-buffer flush failed: {error}");
567 break;
568 }
569 Err(error) => {
570 if !final_flush_error_is_retryable(&error) {
571 eprintln!("feox: final write-buffer flush failed: {error}");
572 break;
573 }
574 retries += 1;
575 if retries == FINAL_FLUSH_RETRY_LIMIT {
576 eprintln!(
577 "feox: final write-buffer flush failed after {retries} attempts: {error}"
578 );
579 break;
580 }
581 thread::sleep(Duration::from_micros(retry_delay_us));
582 retry_delay_us = (retry_delay_us * 2).min(1_000);
583 }
584 }
585 }
586 }
587}
588
589fn final_flush_error_is_retryable(error: &FeoxError) -> bool {
590 matches!(error, FeoxError::IoError(_) | FeoxError::StaleExtent)
591}
592
593fn flush_worker_shards(
594 ctx: &WorkerContext,
595 format: &dyn RecordFormat,
596 flush_retirements: bool,
597) -> Result<bool> {
598 let released_sectors = ctx
599 .retirement_queue
600 .released_sectors
601 .load(Ordering::Relaxed);
602 let mut has_retries = false;
603 let mut first_error = None;
604
605 for shard_id in (ctx.worker_id..ctx.sharded_buffers.len()).step_by(ctx.worker_count) {
606 let buffer = &ctx.sharded_buffers[shard_id];
607 let entries = buffer.drain_entries();
608 if entries.is_empty() {
609 continue;
610 }
611
612 let mut entries = entries.into_iter();
613 let mut shard_retries = Vec::new();
614 let mut shard_failed = false;
615
616 loop {
617 let batch = entries
618 .by_ref()
619 .take(ALLOCATION_JOURNAL_MAX_ENTRIES)
620 .collect::<Vec<_>>();
621 if batch.is_empty() {
622 break;
623 }
624
625 let BatchOutcome { result, retries } = process_write_batch(ctx, batch, format);
626 shard_retries.extend(retries);
627
628 if let Err(error) = result {
629 shard_failed = true;
630 shard_retries.extend(entries);
631 if first_error.is_none() {
632 first_error = Some(error);
633 }
634 break;
635 }
636 }
637
638 has_retries |= !shard_retries.is_empty();
639 buffer.requeue_entries(shard_retries, &ctx.stats, shard_failed);
640 ctx.stats.flush_count.fetch_add(1, Ordering::Relaxed);
641 }
642
643 if first_error
644 .as_ref()
645 .is_some_and(|error| matches!(error, FeoxError::OutOfSpace | FeoxError::AllocationFailed))
646 {
647 flush_pending_deletions(
648 &ctx.retirement_queue,
649 &ctx.disk_io,
650 &ctx.free_space,
651 &ctx.stats,
652 format,
653 )?;
654 if ctx
655 .retirement_queue
656 .released_sectors
657 .load(Ordering::Relaxed)
658 != released_sectors
659 {
660 return Ok(true);
661 }
662 } else if first_error.is_none() && flush_retirements {
663 match flush_pending_deletions(
664 &ctx.retirement_queue,
665 &ctx.disk_io,
666 &ctx.free_space,
667 &ctx.stats,
668 format,
669 ) {
670 Ok(retries) => has_retries |= retries,
671 Err(error) => first_error = Some(error),
672 }
673 }
674
675 match first_error {
676 Some(error) => Err(error),
677 None => Ok(has_retries),
678 }
679}
680
681fn flush_pending_deletions(
682 retirement_queue: &RetirementQueue,
683 disk_io: &Arc<RwLock<DiskIO>>,
684 free_space: &Arc<RwLock<FreeSpaceManager>>,
685 stats: &Arc<Statistics>,
686 format: &dyn RecordFormat,
687) -> Result<bool> {
688 let _flush_guard = retirement_queue.flush.lock();
689 let delete_operations = {
690 let mut pending = retirement_queue.pending.lock();
691 if pending.is_empty() {
692 return Ok(false);
693 }
694 std::mem::take(&mut *pending)
695 };
696
697 let mut retries = Vec::new();
698 let mut released_sectors = 0;
699 let result = process_deletions(
700 disk_io,
701 free_space,
702 stats,
703 format,
704 delete_operations,
705 &mut retries,
706 &mut released_sectors,
707 );
708 if released_sectors != 0 {
709 retirement_queue
710 .released_sectors
711 .fetch_add(released_sectors, Ordering::Relaxed);
712 }
713 let has_retries = !retries.is_empty();
714 if has_retries {
715 retirement_queue.pending.lock().extend(retries);
716 }
717 result.map(|_| has_retries)
718}
719
720fn process_deletions(
721 disk_io: &Arc<RwLock<DiskIO>>,
722 free_space: &Arc<RwLock<FreeSpaceManager>>,
723 stats: &Arc<Statistics>,
724 format: &dyn RecordFormat,
725 delete_operations: Vec<WriteEntry>,
726 retries: &mut Vec<WriteEntry>,
727 released_sectors: &mut u64,
728) -> Result<()> {
729 let capacity = delete_operations.len();
730 let mut first_error = None;
731 let mut marker_writes = Vec::with_capacity(capacity);
732 let mut marker_extents = Vec::with_capacity(capacity);
733 let mut release_operations = Vec::with_capacity(capacity);
734
735 for entry in delete_operations {
736 let sector = entry.record.sector.load(Ordering::Acquire);
737 if sector == 0 {
738 continue;
739 }
740 if entry.work_status.load(Ordering::Acquire) == DELETE_MARKER_DURABLE {
741 release_operations.push(entry);
742 continue;
743 }
744 if !entry.record.successor_is_durable_or_deleted() {
745 retries.push(entry);
746 continue;
747 }
748
749 entry.record.retire_extent();
750 if entry.record.extent_has_readers() {
751 retries.push(entry);
752 continue;
753 }
754 let sectors_needed = format_extent_size(&entry, format);
755 marker_extents.push((sector, sectors_needed));
756 marker_writes.push(entry);
757 }
758
759 if !marker_writes.is_empty() {
760 match disk_io.write().retire_extents(&marker_extents) {
761 Ok(()) => {
762 for entry in &marker_writes {
763 entry
764 .work_status
765 .store(DELETE_MARKER_DURABLE, Ordering::Release);
766 }
767 release_operations.append(&mut marker_writes);
768 }
769 Err(error) => {
770 stats.record_sector_release_failure();
771 eprintln!("feox: extent retirement failed: {error}");
772 retries.append(&mut release_operations);
773 retries.append(&mut marker_writes);
774 return Err(error);
775 }
776 }
777 }
778
779 let mut releasable = Vec::with_capacity(release_operations.len());
780 for entry in release_operations {
781 if entry.record.extent_has_readers() {
782 retries.push(entry);
783 continue;
784 }
785 releasable.push(entry);
786 }
787
788 releasable.sort_unstable_by_key(|entry| entry.record.sector.load(Ordering::Acquire));
789 let mut free_space_guard = free_space.write();
790 let mut group = Vec::with_capacity(releasable.len());
791 let mut group_end = 0;
792 for entry in releasable {
793 let sector = entry.record.sector.load(Ordering::Acquire);
794 let sectors_needed = format_extent_size(&entry, format) as u64;
795 if !group.is_empty() && sector != group_end {
796 release_retirement_group(
797 &mut group,
798 &mut free_space_guard,
799 stats,
800 format,
801 retries,
802 &mut first_error,
803 released_sectors,
804 );
805 }
806 group_end = sector + sectors_needed;
807 group.push(entry);
808 }
809 release_retirement_group(
810 &mut group,
811 &mut free_space_guard,
812 stats,
813 format,
814 retries,
815 &mut first_error,
816 released_sectors,
817 );
818
819 match first_error {
820 Some(error) => Err(error),
821 None => Ok(()),
822 }
823}
824
825fn release_retirement_group(
826 group: &mut Vec<WriteEntry>,
827 free_space: &mut FreeSpaceManager,
828 stats: &Statistics,
829 format: &dyn RecordFormat,
830 retries: &mut Vec<WriteEntry>,
831 first_error: &mut Option<FeoxError>,
832 released_sectors: &mut u64,
833) {
834 let Some(first) = group.first() else {
835 return;
836 };
837 let sector = first.record.sector.load(Ordering::Acquire);
838 let sectors_needed = group
839 .iter()
840 .map(|entry| format_extent_size(entry, format) as u64)
841 .sum::<u64>();
842
843 match free_space.release_sectors(sector, sectors_needed) {
844 Ok(()) => {
845 *released_sectors += sectors_needed;
846 stats
847 .disk_usage
848 .fetch_sub(sectors_needed * FEOX_BLOCK_SIZE as u64, Ordering::Relaxed);
849 group.clear();
850 }
851 Err(error) => {
852 stats.record_sector_release_failure();
853 eprintln!("feox: sector release failed for {sector}+{sectors_needed}: {error}");
854 if first_error.is_none() {
855 *first_error = Some(error);
856 }
857 retries.append(group);
858 }
859 }
860}
861
862fn format_extent_size(entry: &WriteEntry, format: &dyn RecordFormat) -> usize {
863 format
864 .total_size(entry.record.key.len(), entry.record.value_len)
865 .div_ceil(FEOX_BLOCK_SIZE)
866}
867
868fn process_write_batch(
869 ctx: &WorkerContext,
870 entries: Vec<WriteEntry>,
871 format: &dyn RecordFormat,
872) -> BatchOutcome {
873 let disk_io = &ctx.disk_io;
874 let free_space = &ctx.free_space;
875 let stats = &ctx.stats;
876 let format_version = ctx.format_version;
877 let fault_scope = ctx.fault_scope;
878 let retirement_queue = &ctx.retirement_queue;
879 let mut prepared_writes = Vec::new();
880 let mut batch_writes = Vec::new();
881 let mut delete_operations = Vec::new();
882 let mut retry_entries = Vec::new();
883 let mut first_error = None;
884
885 for entry in entries {
886 match entry.op {
887 Operation::Insert | Operation::Update => {
888 let sector = reserved_sector(&entry);
889 if entry.record.sector.load(Ordering::Acquire) == 0
890 && (entry.record.refcount.load(Ordering::Acquire) > 0 || sector.is_some())
891 {
892 match prepare_record_data(&entry.record, format, disk_io) {
893 Ok(data) => {
894 let sectors_needed = data.len().div_ceil(FEOX_BLOCK_SIZE);
895 prepared_writes.push(PreparedWrite {
896 data,
897 sectors_needed,
898 entry,
899 sector,
900 });
901 }
902 Err(error) => {
903 if first_error.is_none() {
904 first_error = Some(error);
905 }
906 retry_entries.push(entry);
907 }
908 }
909 }
910 }
911 Operation::Delete => {
912 delete_operations.push(entry);
913 }
914 _ => {}
915 }
916 }
917
918 let stamp = format_version >= SEQ_TOKEN_MIN_VERSION;
919 let has_deletions = !delete_operations.is_empty();
920 if has_deletions {
921 retirement_queue
922 .pending
923 .lock()
924 .extend(delete_operations.drain(..));
925 }
926
927 if !prepared_writes.is_empty() {
928 let mut free_space_guard = free_space.write();
929 for index in 0..prepared_writes.len() {
930 let sectors_needed = prepared_writes[index].sectors_needed;
931 let sector = match prepared_writes[index].sector {
932 Some(sector) => sector,
933 None => match free_space_guard.allocate_sectors(sectors_needed as u64) {
934 Ok(sector) => {
935 reserve_sector(&prepared_writes[index].entry, sector);
936 stats.disk_usage.fetch_add(
937 (sectors_needed * FEOX_BLOCK_SIZE) as u64,
938 Ordering::Relaxed,
939 );
940 prepared_writes[index].sector = Some(sector);
941 sector
942 }
943 Err(error) => {
944 drop(free_space_guard);
945 let _ = release_allocations(free_space, &prepared_writes, stats);
946 retry_entries.extend(prepared_writes.drain(..).map(|write| write.entry));
947 retry_entries.extend(delete_operations);
948 return BatchOutcome {
949 result: Err(error),
950 retries: retry_entries,
951 };
952 }
953 },
954 };
955 let write = &mut prepared_writes[index];
956 if stamp {
957 stamp_seq_token(&mut write.data, sector, format);
958 }
959 batch_writes.push((sector, Bytes::from(std::mem::take(&mut write.data))));
960 }
961 }
962
963 if !batch_writes.is_empty() {
964 let mut disk_guard = disk_io.write();
965 for write in &prepared_writes {
966 mark_reservation_dirty(&write.entry);
967 }
968
969 let journal_extents = prepared_writes
970 .iter()
971 .map(|write| {
972 (
973 write.sector.expect("prepared write has an allocation"),
974 write.sectors_needed,
975 )
976 })
977 .collect::<Vec<_>>();
978 let journal_active = !journal_extents.is_empty();
979
980 if journal_active {
981 match disk_guard.write_allocation_journal(&journal_extents) {
982 Ok(()) => crash_at("after_allocation_intent"),
983 Err(error @ FeoxError::IndeterminateWrite(_)) => {
984 return failed_batch_outcome(
985 &mut disk_guard,
986 free_space,
987 &mut prepared_writes,
988 delete_operations,
989 retry_entries,
990 stats,
991 BatchFailure {
992 error,
993 clear_journal: true,
994 indeterminate: true,
995 },
996 );
997 }
998 Err(error) => {
999 return failed_batch_outcome(
1000 &mut disk_guard,
1001 free_space,
1002 &mut prepared_writes,
1003 delete_operations,
1004 retry_entries,
1005 stats,
1006 BatchFailure {
1007 error,
1008 clear_journal: true,
1009 indeterminate: false,
1010 },
1011 );
1012 }
1013 }
1014 }
1015
1016 if has_deletions {
1017 crash_at("before_replacement_write");
1018 }
1019
1020 let mut attempts = 3;
1021 let mut delay_us = 100;
1022
1023 while attempts > 0 {
1024 let result = if fail_at(RECORD_WRITE, fault_scope) {
1025 Err(FeoxError::IoError(std::io::Error::other(
1026 "injected record write failure",
1027 )))
1028 } else {
1029 disk_guard.batch_write_bytes(&batch_writes)
1030 };
1031
1032 match result {
1033 Ok(()) => {
1034 break;
1035 }
1036 Err(error @ FeoxError::IndeterminateWrite(_)) => {
1037 return failed_batch_outcome(
1038 &mut disk_guard,
1039 free_space,
1040 &mut prepared_writes,
1041 delete_operations,
1042 retry_entries,
1043 stats,
1044 BatchFailure {
1045 error,
1046 clear_journal: journal_active,
1047 indeterminate: true,
1048 },
1049 );
1050 }
1051 Err(e) => {
1052 attempts -= 1;
1053 if attempts > 0 {
1054 let jitter = {
1056 use rand::Rng;
1057 let mut rng = rand::rng();
1058 (delay_us * rng.random_range(-10..=10)) / 100
1059 };
1060 let actual_delay = (delay_us + jitter).max(1);
1061 thread::sleep(Duration::from_micros(actual_delay as u64));
1062 delay_us *= 2;
1063 } else {
1064 return failed_batch_outcome(
1065 &mut disk_guard,
1066 free_space,
1067 &mut prepared_writes,
1068 delete_operations,
1069 retry_entries,
1070 stats,
1071 BatchFailure {
1072 error: e,
1073 clear_journal: journal_active,
1074 indeterminate: false,
1075 },
1076 );
1077 }
1078 }
1079 }
1080 }
1081
1082 if journal_active {
1083 crash_at("before_allocation_journal_clear");
1084 if let Err(error) = disk_guard.clear_allocation_journal() {
1085 let indeterminate = matches!(error, FeoxError::IndeterminateWrite(_));
1086 return failed_batch_outcome(
1087 &mut disk_guard,
1088 free_space,
1089 &mut prepared_writes,
1090 delete_operations,
1091 retry_entries,
1092 stats,
1093 BatchFailure {
1094 error,
1095 clear_journal: true,
1096 indeterminate,
1097 },
1098 );
1099 }
1100 }
1101
1102 if has_deletions {
1103 crash_at("after_replacement_write");
1104 }
1105 for write in &prepared_writes {
1106 write
1107 .entry
1108 .record
1109 .sector
1110 .store(write.sector.unwrap(), Ordering::Release);
1111 std::sync::atomic::fence(Ordering::Release);
1112 write.entry.record.clear_value();
1113 }
1114 stats.record_write_flushed(prepared_writes.len() as u64);
1115 }
1116
1117 let result = match first_error {
1118 Some(error) => Err(error),
1119 None => Ok(()),
1120 };
1121
1122 BatchOutcome {
1123 result,
1124 retries: retry_entries,
1125 }
1126}
1127
1128#[cfg(test)]
1129fn crash_at(point: &str) {
1130 if std::env::var("FEOX_TEST_CRASH_POINT").as_deref() == Ok(point) {
1131 std::process::exit(86);
1132 }
1133}
1134
1135#[cfg(not(test))]
1136#[inline]
1137fn crash_at(_: &str) {}
1138
1139fn failed_batch_outcome(
1140 disk_io: &mut DiskIO,
1141 free_space: &Arc<RwLock<FreeSpaceManager>>,
1142 prepared_writes: &mut Vec<PreparedWrite>,
1143 delete_operations: Vec<WriteEntry>,
1144 mut retry_entries: Vec<WriteEntry>,
1145 stats: &Statistics,
1146 failure: BatchFailure,
1147) -> BatchOutcome {
1148 stats.record_write_failed();
1149
1150 let error = if failure.indeterminate {
1151 quarantine_allocations(prepared_writes);
1152 failure.error
1153 } else {
1154 match cleanup_failed_allocations(
1155 disk_io,
1156 free_space,
1157 prepared_writes,
1158 stats,
1159 failure.clear_journal,
1160 ) {
1161 Ok(()) => failure.error,
1162 Err(cleanup_error) => {
1163 quarantine_allocations(prepared_writes);
1164 disk_io.poison_writes(cleanup_error)
1165 }
1166 }
1167 };
1168
1169 retry_entries.extend(prepared_writes.drain(..).map(|write| write.entry));
1170 retry_entries.extend(delete_operations);
1171 BatchOutcome {
1172 result: Err(error),
1173 retries: retry_entries,
1174 }
1175}
1176
1177fn release_allocations(
1178 free_space: &Arc<RwLock<FreeSpaceManager>>,
1179 allocations: &[PreparedWrite],
1180 stats: &Statistics,
1181) -> Result<()> {
1182 let mut first_error = None;
1183 let mut free_space_guard = free_space.write();
1184 for allocation in allocations {
1185 let Some(sector) = allocation.sector else {
1186 continue;
1187 };
1188 if reservation_is_dirty(&allocation.entry) {
1189 continue;
1190 }
1191 match free_space_guard.release_sectors(sector, allocation.sectors_needed as u64) {
1192 Ok(()) => {
1193 stats.disk_usage.fetch_sub(
1194 (allocation.sectors_needed * FEOX_BLOCK_SIZE) as u64,
1195 Ordering::Relaxed,
1196 );
1197 clear_reserved_sector(&allocation.entry);
1198 }
1199 Err(error) => {
1200 stats.record_sector_release_failure();
1201 if first_error.is_none() {
1202 first_error = Some(error);
1203 }
1204 }
1205 }
1206 }
1207 match first_error {
1208 Some(error) => Err(error),
1209 None => Ok(()),
1210 }
1211}
1212
1213fn quarantine_allocations(allocations: &[PreparedWrite]) {
1214 for allocation in allocations {
1215 if allocation.sector.is_some() {
1216 quarantine_reservation(&allocation.entry);
1217 }
1218 }
1219}
1220
1221fn cleanup_failed_allocations(
1222 disk_io: &mut DiskIO,
1223 free_space: &Arc<RwLock<FreeSpaceManager>>,
1224 allocations: &[PreparedWrite],
1225 stats: &Statistics,
1226 clear_journal: bool,
1227) -> Result<()> {
1228 let extents = allocations
1229 .iter()
1230 .filter(|allocation| !reservation_is_quarantined(&allocation.entry))
1231 .filter_map(|allocation| {
1232 allocation
1233 .sector
1234 .map(|sector| (sector, allocation.sectors_needed))
1235 })
1236 .collect::<Vec<_>>();
1237
1238 if extents.is_empty() {
1239 if clear_journal {
1240 disk_io.clear_allocation_journal()?;
1241 }
1242 } else if let Err(error) = disk_io.retire_extents(&extents) {
1243 stats.record_sector_release_failure();
1244 eprintln!("feox: failed-write scrub failed: {error}");
1245 return Err(error);
1246 }
1247
1248 let mut free_space = free_space.write();
1249 release_scrubbed_allocations(&mut free_space, allocations, stats)
1250}
1251
1252fn release_scrubbed_allocations(
1253 free_space: &mut FreeSpaceManager,
1254 allocations: &[PreparedWrite],
1255 stats: &Statistics,
1256) -> Result<()> {
1257 let mut ordered = allocations
1258 .iter()
1259 .filter(|allocation| !reservation_is_quarantined(&allocation.entry))
1260 .filter_map(|allocation| allocation.sector.map(|sector| (sector, allocation)))
1261 .collect::<Vec<_>>();
1262 ordered.sort_unstable_by_key(|(sector, _)| *sector);
1263
1264 let mut first_error = None;
1265 let mut group_start = 0;
1266 while group_start < ordered.len() {
1267 let sector = ordered[group_start].0;
1268 let mut group_end = group_start + 1;
1269 let mut end_sector = sector + ordered[group_start].1.sectors_needed as u64;
1270 while group_end < ordered.len() && ordered[group_end].0 == end_sector {
1271 end_sector += ordered[group_end].1.sectors_needed as u64;
1272 group_end += 1;
1273 }
1274 let sectors_needed = end_sector - sector;
1275
1276 match free_space.release_sectors(sector, sectors_needed) {
1277 Ok(()) => {
1278 stats
1279 .disk_usage
1280 .fetch_sub(sectors_needed * FEOX_BLOCK_SIZE as u64, Ordering::Relaxed);
1281 for (_, allocation) in &ordered[group_start..group_end] {
1282 mark_reservation_clean(&allocation.entry);
1283 clear_reserved_sector(&allocation.entry);
1284 }
1285 }
1286 Err(error) => {
1287 stats.record_sector_release_failure();
1288 eprintln!(
1289 "feox: failed-write sector release failed for {sector}+{sectors_needed}: {error}"
1290 );
1291 if first_error.is_none() {
1292 first_error = Some(error);
1293 }
1294 }
1295 }
1296 group_start = group_end;
1297 }
1298
1299 match first_error {
1300 Some(error) => Err(error),
1301 None => Ok(()),
1302 }
1303}
1304
1305fn prepare_record_data(
1306 record: &Record,
1307 format: &dyn RecordFormat,
1308 disk_io: &Arc<RwLock<DiskIO>>,
1309) -> Result<Vec<u8>> {
1310 let total_size = format.total_size(record.key.len(), record.value_len);
1311 let sectors_needed = total_size.div_ceil(FEOX_BLOCK_SIZE);
1312 let padded_size = sectors_needed * FEOX_BLOCK_SIZE;
1313
1314 match record.get_value() {
1315 Some(value) => Ok(serialize_record_data(record, format, &value, padded_size)),
1316 None => prepare_deferred_record_data(record, format, disk_io, padded_size),
1317 }
1318}
1319
1320fn serialize_record_data(
1321 record: &Record,
1322 format: &dyn RecordFormat,
1323 value: &[u8],
1324 padded_size: usize,
1325) -> Vec<u8> {
1326 let mut data = Vec::with_capacity(padded_size);
1327 data.extend_from_slice(&SECTOR_MARKER.to_le_bytes());
1328 data.extend_from_slice(&0u16.to_le_bytes());
1329 format.serialize_record_into(record, false, &mut data);
1330 data.extend_from_slice(value);
1331 data.resize(padded_size, 0);
1332
1333 data
1334}
1335
1336fn prepare_deferred_record_data(
1337 record: &Record,
1338 format: &dyn RecordFormat,
1339 disk_io: &Arc<RwLock<DiskIO>>,
1340 padded_size: usize,
1341) -> Result<Vec<u8>> {
1342 let mut source = record.value_source().ok_or(FeoxError::StaleExtent)?;
1343 loop {
1344 if let Some(value) = source.get_value() {
1345 return Ok(serialize_record_data(record, format, &value, padded_size));
1346 }
1347 if source.sector.load(Ordering::Acquire) != 0 {
1348 break;
1349 }
1350 source = source.value_source().ok_or(FeoxError::StaleExtent)?;
1351 }
1352
1353 let extent = source.acquire_extent().ok_or(FeoxError::StaleExtent)?;
1354 let sector = source.sector.load(Ordering::Acquire);
1355 if sector == 0 {
1356 return Err(FeoxError::StaleExtent);
1357 }
1358 if source.value_len != record.value_len || source.key != record.key {
1359 return Err(FeoxError::InvalidRecord);
1360 }
1361 let total_size = format.total_size(source.key.len(), source.value_len);
1362 let sectors = total_size.div_ceil(FEOX_BLOCK_SIZE);
1363 let mut data = disk_io.read().read_sectors_sync(sector, sectors as u64)?;
1364 drop(extent);
1365 if !sector_holds_record(&data, &source) {
1366 return Err(FeoxError::StaleExtent);
1367 }
1368
1369 let value_offset = format.value_offset(record.key.len());
1370 if data.len() != padded_size || value_offset > data.len() {
1371 return Err(FeoxError::InvalidRecord);
1372 }
1373 let mut header = Vec::with_capacity(value_offset);
1374 header.extend_from_slice(&SECTOR_MARKER.to_le_bytes());
1375 header.extend_from_slice(&0u16.to_le_bytes());
1376 format.serialize_record_into(record, false, &mut header);
1377 if header.len() != value_offset {
1378 return Err(FeoxError::InvalidRecord);
1379 }
1380 data[..value_offset].copy_from_slice(&header);
1381 Ok(data)
1382}
1383
1384impl Drop for WriteBuffer {
1385 fn drop(&mut self) {
1386 self.complete_shutdown();
1387 }
1388}
1389
1390#[cfg(test)]
1391#[path = "../tests/write_buffer_safety_tests.rs"]
1392mod tests;