Skip to main content

feoxdb/storage/
write_buffer.rs

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/// Sharded write buffer for reducing contention
24/// Each thread consistently uses the same shard to improve cache locality
25#[repr(align(64))] // Cache line alignment
26pub struct ShardedWriteBuffer {
27    /// Buffered writes pending flush
28    buffer: Mutex<VecDeque<WriteEntry>>,
29
30    /// Number of entries in buffer
31    count: AtomicUsize,
32
33    /// Total size of buffered data
34    size: AtomicUsize,
35}
36
37/// Write entry for buffered operations
38pub 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
145/// Main write buffer coordinator
146pub struct WriteBuffer {
147    /// Sharded buffers to reduce contention between threads
148    sharded_buffers: Arc<Vec<CachePadded<ShardedWriteBuffer>>>,
149
150    /// Shared disk I/O handle
151    disk_io: Arc<RwLock<DiskIO>>,
152
153    /// Free space manager for sector allocation
154    free_space: Arc<RwLock<FreeSpaceManager>>,
155
156    /// Per-worker channels for targeted flush requests
157    worker_channels: Vec<Sender<FlushRequest>>,
158
159    /// Background worker handles
160    worker_handles: Mutex<Vec<JoinHandle<()>>>,
161
162    /// Periodic flush thread handle
163    periodic_flush_handle: Mutex<Option<JoinHandle<()>>>,
164
165    /// Shutdown flag
166    shutdown: Arc<AtomicBool>,
167
168    /// Shared statistics
169    stats: Arc<Statistics>,
170
171    /// Stable per-store key hasher for preserving per-key write order
172    shard_hasher: RandomState,
173
174    retirement_queue: Arc<RetirementQueue>,
175
176    /// Format version for record serialization
177    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        // Use half CPU count for both shards and workers
284        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    /// Add write operation to buffer (lock-free fast path)
314    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    /// Start background worker threads
356    pub fn start_workers(&mut self, num_workers: usize) {
357        // Ensure we have the right number of workers for shards
358        let num_shards = self.sharded_buffers.len();
359        let actual_workers = num_workers.clamp(1, num_shards);
360
361        // Create per-worker channels
362        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        // Start workers, each owning one shard
370        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        // Start periodic flush coordinator
392        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    /// Force flush and wait for completion
424    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    /// Shutdown write buffer
483    pub fn initiate_shutdown(&self) {
484        self.shutdown.store(true, Ordering::Release);
485
486        // Don't call force_flush here as it can block
487        // Workers will see the shutdown flag and exit gracefully
488    }
489
490    /// Complete shutdown - must be called after initiate_shutdown
491    pub fn complete_shutdown(&mut self) {
492        self.finish_shutdown();
493    }
494
495    pub(crate) fn finish_shutdown(&self) {
496        // Ensure shutdown flag is set
497        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        // Signal workers to stop and wait
504        let handles = std::mem::take(&mut *self.worker_handles.lock());
505        for handle in handles {
506            let _ = handle.join();
507        }
508
509        // Note: disk_io shutdown is handled by the Store's Drop implementation
510        // to ensure proper ordering
511    }
512
513    /// Legacy shutdown for compatibility
514    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
524/// Background worker for processing write buffer flushes
525fn 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        // Wait for flush request with timeout to check shutdown periodically
534        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                        // Exponential backoff with jitter ±10%
1055                        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;