Skip to main content

moonpool_sim/sim/
storage_ops.rs

1//! Storage I/O operations for the simulation.
2//!
3//! This module contains storage-related event handlers and methods extracted from
4//! `world.rs` to improve code organization. It handles file operations like
5//! open, read, write, sync, and provides fault injection for testing storage reliability.
6
7use std::collections::HashMap;
8use std::net::IpAddr;
9use std::task::Waker;
10use std::time::Duration;
11use tracing::instrument;
12
13use crate::assert_reachable;
14use crate::storage::StorageError;
15
16use crate::chaos::fault_events::SimFaultEvent;
17
18use super::{
19    events::{Event, ScheduledEvent, StorageOperation},
20    rng::{sim_random, sim_random_range},
21    state::{DiskDegradationState, DiskEpisodeKind, FileId, PendingOpType, PendingStorageOp},
22    world::{SimInner, SimWorld},
23};
24
25// =============================================================================
26// Storage Event Handlers
27// =============================================================================
28
29/// Find and remove the first pending operation of the given type for a file.
30///
31/// Returns the sequence number and operation if found, or None if no such operation exists.
32fn take_pending_op(
33    inner: &mut SimInner,
34    file_id: FileId,
35    op_type: PendingOpType,
36) -> Option<(u64, PendingStorageOp)> {
37    let file_state = inner.storage.files.get_mut(&file_id)?;
38
39    let op_seq = file_state
40        .pending_ops
41        .iter()
42        .find(|(_, op)| op.op_type == op_type)
43        .map(|(&seq, _)| seq)?;
44
45    let op = file_state.pending_ops.remove(&op_seq)?;
46    Some((op_seq, op))
47}
48
49/// Advance an owning process's disk-degradation episode state machine and return
50/// the episode (if any) active for the operation now being scheduled.
51///
52/// Episodes are scoped per owner (keyed by process IP): a single stall/throttle
53/// window applies to every file that process owns, modelling device-level
54/// degradation that hits all of a machine's disks together (issue #147).
55///
56/// On each call: expire an episode whose window has elapsed, keep an episode
57/// that is still active, or — when none is active and a family is enabled —
58/// probabilistically enter a stall or throttle episode (stall takes precedence).
59///
60/// **Off-by-default guarantee:** when both probabilities are zero this returns
61/// without drawing any simulation RNG, keeping the RNG stream byte-identical to
62/// steady-state runs (so existing latency/determinism tests are unaffected).
63fn update_disk_episode(
64    episodes: &mut HashMap<IpAddr, DiskDegradationState>,
65    owner_ip: IpAddr,
66    now: Duration,
67    stall_probability: f64,
68    stall_duration: Duration,
69    throttle_probability: f64,
70    throttle_duration: Duration,
71) -> Option<DiskDegradationState> {
72    // Expire an episode whose window has elapsed.
73    if let Some(episode) = episodes.get(&owner_ip).copied()
74        && now >= episode.expires_at
75    {
76        episodes.remove(&owner_ip);
77    }
78
79    // An episode is still active: keep it, drawing no new randomness.
80    if let Some(episode) = episodes.get(&owner_ip).copied() {
81        return Some(episode);
82    }
83
84    // Off by default: never perturb the RNG stream when both families are disabled.
85    if stall_probability <= 0.0 && throttle_probability <= 0.0 {
86        return None;
87    }
88
89    // Probabilistically enter a new episode using a single draw.
90    let roll = sim_random::<f64>();
91    let episode = if roll < stall_probability {
92        assert_reachable!("disk: stall episode entered");
93        Some(DiskDegradationState {
94            kind: DiskEpisodeKind::Stall,
95            expires_at: now + stall_duration,
96        })
97    } else if roll < stall_probability + throttle_probability {
98        assert_reachable!("disk: throttle episode entered");
99        Some(DiskDegradationState {
100            kind: DiskEpisodeKind::Throttle,
101            expires_at: now + throttle_duration,
102        })
103    } else {
104        None
105    };
106    if let Some(episode) = episode {
107        episodes.insert(owner_ip, episode);
108    }
109    episode
110}
111
112/// Read the disk-episode knobs for an owning process into a copyable tuple, so
113/// the immutable borrow of `inner.storage` ends before `inner.storage.disk_episodes`
114/// is re-borrowed mutably for the episode update.
115fn disk_episode_knobs(inner: &SimInner, owner_ip: IpAddr) -> (f64, Duration, f64, Duration) {
116    let config = inner.storage.config_for(owner_ip);
117    (
118        config.disk_stall_probability,
119        config.disk_stall_duration,
120        config.disk_throttle_probability,
121        config.disk_throttle_duration,
122    )
123}
124
125/// Handle storage I/O events.
126///
127/// Storage events represent the completion of I/O operations.
128/// Processing applies faults and wakes waiting tasks.
129pub(crate) fn handle_storage_event(
130    inner: &mut SimInner,
131    file_id: u64,
132    operation: StorageOperation,
133) {
134    let file_id = FileId(file_id);
135
136    match operation {
137        StorageOperation::ReadComplete { len: _ } => {
138            handle_read_complete(inner, file_id);
139        }
140        StorageOperation::WriteComplete { len: _ } => {
141            handle_write_complete(inner, file_id);
142        }
143        StorageOperation::SyncComplete => {
144            handle_sync_complete(inner, file_id);
145        }
146        StorageOperation::OpenComplete => {
147            handle_open_complete(inner, file_id);
148        }
149        StorageOperation::SetLenComplete { new_len } => {
150            handle_set_len_complete(inner, file_id, new_len);
151        }
152    }
153}
154
155/// Handle read operation completion.
156fn handle_read_complete(inner: &mut SimInner, file_id: FileId) {
157    let read_fault_probability = inner.storage.files.get(&file_id).map_or(0.0, |f| {
158        inner.storage.config_for(f.owner_ip).read_fault_probability
159    });
160
161    // Find and remove the oldest pending read operation
162    let Some((op_seq, op)) = take_pending_op(inner, file_id, PendingOpType::Read) else {
163        tracing::warn!("ReadComplete for unknown file {:?}", file_id);
164        return;
165    };
166
167    let (offset, len) = (op.offset, op.len);
168
169    // Apply read fault injection - mark sectors as faulted based on probability
170    let mut read_faulted = false;
171    if read_fault_probability > 0.0
172        && let Some(file_state) = inner.storage.files.get_mut(&file_id)
173    {
174        let offset_usize = usize::try_from(offset).expect("offset fits in usize");
175        let start_sector = offset_usize / crate::storage::SECTOR_SIZE;
176        let end_sector = (offset_usize + len).div_ceil(crate::storage::SECTOR_SIZE);
177
178        for sector in start_sector..end_sector {
179            if sim_random::<f64>() < read_fault_probability {
180                file_state.storage.set_fault(sector);
181                read_faulted = true;
182                tracing::info!(
183                    "Read fault injected for file {:?}, sector {}",
184                    file_id,
185                    sector
186                );
187            }
188        }
189    }
190    if read_faulted {
191        let ip = inner.storage.files.get(&file_id).map(|f| f.owner_ip);
192        if let Some(ip) = ip {
193            inner.record_fault(SimFaultEvent::StorageReadFault {
194                ip: ip.to_string(),
195                file_id: file_id.0,
196            });
197        }
198    }
199
200    // Wake the waker for this operation
201    if let Some(waker) = inner.wakers.storage_ops.remove(&(file_id, op_seq)) {
202        tracing::trace!("Waking read waker for file {:?}, op {}", file_id, op_seq);
203        waker.wake();
204    }
205}
206
207/// Handle write operation completion.
208///
209/// Applies the write to storage with potential fault injection:
210/// - `phantom_write_probability`: write appears to succeed but isn't persisted
211/// - `misdirect_write_probability`: write lands at wrong location
212fn handle_write_complete(inner: &mut SimInner, file_id: FileId) {
213    let config = inner
214        .storage
215        .files
216        .get(&file_id)
217        .map(|f| inner.storage.config_for(f.owner_ip).clone())
218        .unwrap_or_default();
219
220    let owner_ip = inner.storage.files.get(&file_id).map(|f| f.owner_ip);
221
222    // Find and remove the oldest pending write operation
223    let Some((op_seq, op)) = take_pending_op(inner, file_id, PendingOpType::Write) else {
224        tracing::warn!("WriteComplete for unknown file {:?}", file_id);
225        return;
226    };
227
228    let (offset, data_opt) = (op.offset, op.data);
229
230    // Apply the write with potential fault injection
231    let mut write_fault_kind: Option<&str> = None;
232    if let Some(data) = data_opt
233        && let Some(file_state) = inner.storage.files.get_mut(&file_id)
234    {
235        // Check for phantom write (write appears to succeed but doesn't persist)
236        if sim_random::<f64>() < config.phantom_write_probability {
237            tracing::info!(
238                "Phantom write injected for file {:?}, offset {}, len {}",
239                file_id,
240                offset,
241                data.len()
242            );
243            file_state.storage.record_phantom_write(offset, &data);
244            write_fault_kind = Some("phantom");
245        }
246        // Check for misdirected write
247        else if sim_random::<f64>() < config.misdirect_write_probability {
248            // Pick a random different offset
249            let max_offset = file_state.storage.size().saturating_sub(data.len() as u64);
250            let mistaken_offset = if max_offset > 0 {
251                sim_random_range(0..max_offset)
252            } else {
253                0
254            };
255            tracing::info!(
256                "Misdirected write injected for file {:?}: intended={}, actual={}",
257                file_id,
258                offset,
259                mistaken_offset
260            );
261            if let Err(e) =
262                file_state
263                    .storage
264                    .apply_misdirected_write(offset, mistaken_offset, &data)
265            {
266                tracing::warn!("Failed to apply misdirected write: {}", e);
267            }
268            write_fault_kind = Some("misdirected");
269        }
270        // Normal write (not synced - may be lost on crash)
271        else if let Err(e) = file_state.storage.write(offset, &data, false) {
272            tracing::warn!("Write failed for file {:?}: {}", file_id, e);
273        } else {
274            // Check for write corruption - mark sectors as faulted after successful write
275            if config.write_fault_probability > 0.0 {
276                let offset_usize = usize::try_from(offset).expect("offset fits in usize");
277                let start_sector = offset_usize / crate::storage::SECTOR_SIZE;
278                let end_sector = (offset_usize + data.len()).div_ceil(crate::storage::SECTOR_SIZE);
279
280                for sector in start_sector..end_sector {
281                    if sim_random::<f64>() < config.write_fault_probability {
282                        file_state.storage.set_fault(sector);
283                        write_fault_kind = Some("corruption");
284                        tracing::info!(
285                            "Write fault injected for file {:?}, sector {}",
286                            file_id,
287                            sector
288                        );
289                    }
290                }
291            }
292        }
293    }
294    if let (Some(kind), Some(ip)) = (write_fault_kind, owner_ip) {
295        inner.record_fault(SimFaultEvent::StorageWriteFault {
296            ip: ip.to_string(),
297            file_id: file_id.0,
298            write_kind: kind.to_string(),
299        });
300    }
301
302    // Wake the waker for this operation
303    if let Some(waker) = inner.wakers.storage_ops.remove(&(file_id, op_seq)) {
304        tracing::trace!("Waking write waker for file {:?}, op {}", file_id, op_seq);
305        waker.wake();
306    }
307}
308
309/// Handle sync operation completion.
310///
311/// Applies `sync_failure_probability` fault injection.
312fn handle_sync_complete(inner: &mut SimInner, file_id: FileId) {
313    let sync_failure_prob = inner.storage.files.get(&file_id).map_or(0.0, |f| {
314        inner
315            .storage
316            .config_for(f.owner_ip)
317            .sync_failure_probability
318    });
319
320    // Find and remove the oldest pending sync operation
321    let Some((op_seq, _)) = take_pending_op(inner, file_id, PendingOpType::Sync) else {
322        tracing::warn!("SyncComplete for unknown file {:?}", file_id);
323        return;
324    };
325
326    // Check for sync failure
327    if sim_random::<f64>() < sync_failure_prob {
328        tracing::info!("Sync failure injected for file {:?}", file_id);
329        // Record the failure so SyncFuture can return an error
330        inner.storage.sync_failures.insert((file_id, op_seq));
331        // On sync failure, we don't call storage.sync()
332        // Data remains in pending state and may be lost on crash
333        let ip = inner.storage.files.get(&file_id).map(|f| f.owner_ip);
334        if let Some(ip) = ip {
335            inner.record_fault(SimFaultEvent::StorageSyncFault {
336                ip: ip.to_string(),
337                file_id: file_id.0,
338            });
339        }
340    } else if let Some(file_state) = inner.storage.files.get_mut(&file_id) {
341        // Successful sync - make all pending writes durable
342        file_state.storage.sync();
343    }
344
345    // Wake the waker for this operation
346    if let Some(waker) = inner.wakers.storage_ops.remove(&(file_id, op_seq)) {
347        tracing::trace!("Waking sync waker for file {:?}, op {}", file_id, op_seq);
348        waker.wake();
349    }
350}
351
352/// Handle open operation completion.
353fn handle_open_complete(inner: &mut SimInner, file_id: FileId) {
354    // Find and remove the oldest pending open operation
355    let Some((op_seq, _)) = take_pending_op(inner, file_id, PendingOpType::Open) else {
356        // File might not have pending open op (it was already "open" on creation)
357        tracing::trace!("OpenComplete for file {:?} (no pending op)", file_id);
358        return;
359    };
360
361    // Wake the waker for this operation
362    if let Some(waker) = inner.wakers.storage_ops.remove(&(file_id, op_seq)) {
363        tracing::trace!("Waking open waker for file {:?}, op {}", file_id, op_seq);
364        waker.wake();
365    }
366}
367
368/// Handle `set_len` operation completion.
369fn handle_set_len_complete(inner: &mut SimInner, file_id: FileId, new_len: u64) {
370    // Find and remove the oldest pending set_len operation
371    let Some((op_seq, _)) = take_pending_op(inner, file_id, PendingOpType::SetLen) else {
372        tracing::warn!("SetLenComplete for unknown file {:?}", file_id);
373        return;
374    };
375
376    // Resize the storage (preserves seed, written/fault bitmaps, and overlays)
377    if let Some(file_state) = inner.storage.files.get_mut(&file_id) {
378        file_state.storage.resize(new_len);
379    }
380
381    // Wake the waker for this operation
382    if let Some(waker) = inner.wakers.storage_ops.remove(&(file_id, op_seq)) {
383        tracing::trace!(
384            "Waking set_len waker for file {:?}, op {}, new_len={}",
385            file_id,
386            op_seq,
387            new_len
388        );
389        waker.wake();
390    }
391}
392
393// =============================================================================
394// Storage Methods for SimWorld
395// =============================================================================
396
397impl SimWorld {
398    /// Access storage configuration for the simulation.
399    ///
400    /// # Panics
401    ///
402    /// Panics if the simulation lock is poisoned by a prior task panic.
403    pub fn with_storage_config<F, R>(&self, f: F) -> R
404    where
405        F: FnOnce(&crate::storage::StorageConfiguration) -> R,
406    {
407        let inner = self
408            .inner
409            .read()
410            .expect("RwLock poisoned: prior task panicked");
411        f(&inner.storage.config)
412    }
413
414    /// Open a file in the simulation.
415    ///
416    /// Creates a new file or opens an existing one based on the options.
417    /// Files are tagged with the `owner_ip` for per-process storage fault injection.
418    /// Schedules an `OpenComplete` event and returns the file ID.
419    pub(crate) fn open_file(
420        &self,
421        path: &str,
422        options: moonpool_core::OpenOptions,
423        initial_size: u64,
424        owner_ip: IpAddr,
425    ) -> Result<FileId, StorageError> {
426        use crate::storage::InMemoryStorage;
427
428        let mut inner = self
429            .inner
430            .write()
431            .expect("RwLock poisoned: prior task panicked");
432        let path_str = path.to_string();
433
434        // Check create_new semantics - fail if file exists
435        if options.is_create_new() && inner.storage.path_to_file.contains_key(&path_str) {
436            return Err(StorageError::AlreadyExists { path: path_str });
437        }
438
439        // Check if file was deleted and create is not set
440        if inner.storage.deleted_paths.contains(&path_str) && !options.is_create() {
441            return Err(StorageError::NotFound { path: path_str });
442        }
443
444        // If file already exists and we're opening it, return existing file ID
445        if let Some(&existing_id) = inner.storage.path_to_file.get(&path_str) {
446            if let Some(file_state) = inner.storage.files.get_mut(&existing_id) {
447                // If truncate is set, reset the storage
448                if options.is_truncate() {
449                    let seed = sim_random::<u64>();
450                    file_state.storage = InMemoryStorage::new(0, seed);
451                    file_state.position = 0;
452                } else if options.is_append() {
453                    // For append mode, seek to end
454                    file_state.position = file_state.storage.size();
455                } else {
456                    // For normal reopen, reset position to start
457                    file_state.position = 0;
458                }
459                // Update options for the new open
460                file_state.options = options;
461                file_state.is_closed = false;
462            }
463            return Ok(existing_id);
464        }
465
466        // File doesn't exist - check if we're allowed to create it
467        if !options.is_create() && !options.is_create_new() {
468            return Err(StorageError::NotFound { path: path_str });
469        }
470
471        // Create new file
472        let file_id = FileId(inner.storage.next_file_id);
473        inner.storage.next_file_id += 1;
474
475        // Remove from deleted paths if re-creating
476        inner.storage.deleted_paths.remove(&path_str);
477
478        // Create in-memory storage with deterministic seed
479        let seed = sim_random::<u64>();
480        let storage = InMemoryStorage::new(initial_size, seed);
481
482        let file_state = super::state::StorageFileState::new(
483            file_id,
484            path_str.clone(),
485            options,
486            storage,
487            owner_ip,
488        );
489
490        inner.storage.files.insert(file_id, file_state);
491        inner.storage.path_to_file.insert(path_str, file_id);
492
493        // Schedule OpenComplete event with minimal latency
494        let open_latency = Duration::from_micros(1);
495        let scheduled_time = inner.current_time + open_latency;
496        let sequence = inner.next_sequence;
497        inner.next_sequence += 1;
498        let event = Event::Storage {
499            file_id: file_id.0,
500            operation: StorageOperation::OpenComplete,
501        };
502        inner
503            .event_queue
504            .schedule(ScheduledEvent::new(scheduled_time, event, sequence));
505
506        tracing::debug!("Opened file {:?} with id {:?}", path, file_id);
507        Ok(file_id)
508    }
509
510    /// Check if a file exists at the given path.
511    pub(crate) fn file_exists(&self, path: &str) -> bool {
512        let inner = self
513            .inner
514            .read()
515            .expect("RwLock poisoned: prior task panicked");
516        let path_str = path.to_string();
517        inner.storage.path_to_file.contains_key(&path_str)
518            && !inner.storage.deleted_paths.contains(&path_str)
519    }
520
521    /// Delete a file at the given path.
522    pub(crate) fn delete_file(&self, path: &str) -> Result<(), StorageError> {
523        let mut inner = self
524            .inner
525            .write()
526            .expect("RwLock poisoned: prior task panicked");
527        let path_str = path.to_string();
528
529        if let Some(file_id) = inner.storage.path_to_file.remove(&path_str) {
530            // Mark file as closed and remove it
531            if let Some(file_state) = inner.storage.files.get_mut(&file_id) {
532                file_state.is_closed = true;
533            }
534            inner.storage.files.remove(&file_id);
535            inner.storage.deleted_paths.insert(path_str);
536            tracing::debug!("Deleted file {:?}", path);
537            Ok(())
538        } else {
539            Err(StorageError::NotFound { path: path_str })
540        }
541    }
542
543    /// Rename a file from one path to another.
544    pub(crate) fn rename_file(&self, from: &str, to: &str) -> Result<(), StorageError> {
545        let mut inner = self
546            .inner
547            .write()
548            .expect("RwLock poisoned: prior task panicked");
549        let from_str = from.to_string();
550        let to_str = to.to_string();
551
552        if let Some(file_id) = inner.storage.path_to_file.remove(&from_str) {
553            // Update the path in the file state
554            if let Some(file_state) = inner.storage.files.get_mut(&file_id) {
555                file_state.path.clone_from(&to_str);
556            }
557            inner.storage.path_to_file.insert(to_str, file_id);
558            inner.storage.deleted_paths.remove(&from_str);
559            tracing::debug!("Renamed file {:?} to {:?}", from, to);
560            Ok(())
561        } else {
562            Err(StorageError::NotFound { path: from_str })
563        }
564    }
565
566    /// Schedule a read operation on a file.
567    ///
568    /// Returns an operation sequence number that can be used to check completion.
569    pub(crate) fn schedule_read(
570        &self,
571        file_id: FileId,
572        offset: u64,
573        len: usize,
574    ) -> Result<u64, StorageError> {
575        let mut inner = self
576            .inner
577            .write()
578            .expect("RwLock poisoned: prior task panicked");
579
580        let file_state = inner
581            .storage
582            .files
583            .get_mut(&file_id)
584            .ok_or(StorageError::InvalidFileHandle { file_id })?;
585
586        if file_state.is_closed {
587            return Err(StorageError::FileClosed { file_id });
588        }
589
590        let op_seq = file_state.next_op_seq;
591        file_state.next_op_seq += 1;
592
593        // Store the pending operation
594        file_state.pending_ops.insert(
595            op_seq,
596            PendingStorageOp {
597                op_type: PendingOpType::Read,
598                offset,
599                len,
600                data: None,
601            },
602        );
603
604        // Advance the disk-degradation episode, then compute latency.
605        let owner_ip = file_state.owner_ip;
606        let now = inner.current_time;
607        let (stall_p, stall_dur, throttle_p, throttle_dur) = disk_episode_knobs(&inner, owner_ip);
608        let episode = update_disk_episode(
609            &mut inner.storage.disk_episodes,
610            owner_ip,
611            now,
612            stall_p,
613            stall_dur,
614            throttle_p,
615            throttle_dur,
616        );
617        let config = inner.storage.config_for(owner_ip);
618        let latency = Self::calculate_storage_latency(config, len, false, episode, now);
619        let scheduled_time = now + latency;
620        let sequence = inner.next_sequence;
621        inner.next_sequence += 1;
622
623        let event = Event::Storage {
624            file_id: file_id.0,
625            operation: StorageOperation::ReadComplete {
626                len: u32::try_from(len).expect("read length fits in u32"),
627            },
628        };
629        inner
630            .event_queue
631            .schedule(ScheduledEvent::new(scheduled_time, event, sequence));
632
633        tracing::trace!(
634            "Scheduled read: file={:?}, offset={}, len={}, op_seq={}",
635            file_id,
636            offset,
637            len,
638            op_seq
639        );
640
641        Ok(op_seq)
642    }
643
644    /// Schedule a write operation on a file.
645    ///
646    /// Returns an operation sequence number that can be used to check completion.
647    pub(crate) fn schedule_write(
648        &self,
649        file_id: FileId,
650        offset: u64,
651        data: Vec<u8>,
652    ) -> Result<u64, StorageError> {
653        let mut inner = self
654            .inner
655            .write()
656            .expect("RwLock poisoned: prior task panicked");
657
658        let file_state = inner
659            .storage
660            .files
661            .get_mut(&file_id)
662            .ok_or(StorageError::InvalidFileHandle { file_id })?;
663
664        if file_state.is_closed {
665            return Err(StorageError::FileClosed { file_id });
666        }
667
668        let op_seq = file_state.next_op_seq;
669        file_state.next_op_seq += 1;
670        let len = data.len();
671
672        // Store the pending operation with the data
673        file_state.pending_ops.insert(
674            op_seq,
675            PendingStorageOp {
676                op_type: PendingOpType::Write,
677                offset,
678                len,
679                data: Some(data),
680            },
681        );
682
683        // Advance the disk-degradation episode, then compute latency.
684        let owner_ip = file_state.owner_ip;
685        let now = inner.current_time;
686        let (stall_p, stall_dur, throttle_p, throttle_dur) = disk_episode_knobs(&inner, owner_ip);
687        let episode = update_disk_episode(
688            &mut inner.storage.disk_episodes,
689            owner_ip,
690            now,
691            stall_p,
692            stall_dur,
693            throttle_p,
694            throttle_dur,
695        );
696        let config = inner.storage.config_for(owner_ip);
697        let latency = Self::calculate_storage_latency(config, len, true, episode, now);
698        let scheduled_time = now + latency;
699        let sequence = inner.next_sequence;
700        inner.next_sequence += 1;
701
702        let event = Event::Storage {
703            file_id: file_id.0,
704            operation: StorageOperation::WriteComplete {
705                len: u32::try_from(len).expect("write length fits in u32"),
706            },
707        };
708        inner
709            .event_queue
710            .schedule(ScheduledEvent::new(scheduled_time, event, sequence));
711
712        tracing::trace!(
713            "Scheduled write: file={:?}, offset={}, len={}, op_seq={}",
714            file_id,
715            offset,
716            len,
717            op_seq
718        );
719
720        Ok(op_seq)
721    }
722
723    /// Schedule a sync operation on a file.
724    ///
725    /// Returns an operation sequence number that can be used to check completion.
726    pub(crate) fn schedule_sync(&self, file_id: FileId) -> Result<u64, StorageError> {
727        let mut inner = self
728            .inner
729            .write()
730            .expect("RwLock poisoned: prior task panicked");
731
732        let file_state = inner
733            .storage
734            .files
735            .get_mut(&file_id)
736            .ok_or(StorageError::InvalidFileHandle { file_id })?;
737
738        if file_state.is_closed {
739            return Err(StorageError::FileClosed { file_id });
740        }
741
742        let op_seq = file_state.next_op_seq;
743        file_state.next_op_seq += 1;
744
745        // Store the pending operation
746        file_state.pending_ops.insert(
747            op_seq,
748            PendingStorageOp {
749                op_type: PendingOpType::Sync,
750                offset: 0,
751                len: 0,
752                data: None,
753            },
754        );
755
756        // Advance the disk-degradation episode, then compute sync latency. Sync
757        // is affected by stalls (the disk is frozen until expiry); a throttle
758        // scales IOPS/bandwidth, which sync latency does not use.
759        let owner_ip = file_state.owner_ip;
760        let now = inner.current_time;
761        let (stall_p, stall_dur, throttle_p, throttle_dur) = disk_episode_knobs(&inner, owner_ip);
762        let episode = update_disk_episode(
763            &mut inner.storage.disk_episodes,
764            owner_ip,
765            now,
766            stall_p,
767            stall_dur,
768            throttle_p,
769            throttle_dur,
770        );
771        let config = inner.storage.config_for(owner_ip);
772        let mut latency = crate::network::sample_latency(&config.sync_latency);
773        if let Some(DiskDegradationState {
774            kind: DiskEpisodeKind::Stall,
775            expires_at,
776        }) = episode
777        {
778            latency += expires_at.saturating_sub(now);
779        }
780        let scheduled_time = now + latency;
781        let sequence = inner.next_sequence;
782        inner.next_sequence += 1;
783
784        let event = Event::Storage {
785            file_id: file_id.0,
786            operation: StorageOperation::SyncComplete,
787        };
788        inner
789            .event_queue
790            .schedule(ScheduledEvent::new(scheduled_time, event, sequence));
791
792        tracing::trace!("Scheduled sync: file={:?}, op_seq={}", file_id, op_seq);
793
794        Ok(op_seq)
795    }
796
797    /// Schedule a `set_len` operation on a file.
798    ///
799    /// Returns an operation sequence number that can be used to check completion.
800    pub(crate) fn schedule_set_len(
801        &self,
802        file_id: FileId,
803        new_len: u64,
804    ) -> Result<u64, StorageError> {
805        let mut inner = self
806            .inner
807            .write()
808            .expect("RwLock poisoned: prior task panicked");
809
810        let file_state = inner
811            .storage
812            .files
813            .get_mut(&file_id)
814            .ok_or(StorageError::InvalidFileHandle { file_id })?;
815
816        if file_state.is_closed {
817            return Err(StorageError::FileClosed { file_id });
818        }
819
820        let op_seq = file_state.next_op_seq;
821        file_state.next_op_seq += 1;
822
823        // Store the pending operation
824        file_state.pending_ops.insert(
825            op_seq,
826            PendingStorageOp {
827                op_type: PendingOpType::SetLen,
828                offset: new_len,
829                len: 0,
830                data: None,
831            },
832        );
833
834        // Use write latency from per-process config
835        let owner_ip = file_state.owner_ip;
836        let config = inner.storage.config_for(owner_ip);
837        let latency = crate::network::sample_latency(&config.write_latency);
838        let scheduled_time = inner.current_time + latency;
839        let sequence = inner.next_sequence;
840        inner.next_sequence += 1;
841
842        let event = Event::Storage {
843            file_id: file_id.0,
844            operation: StorageOperation::SetLenComplete { new_len },
845        };
846        inner
847            .event_queue
848            .schedule(ScheduledEvent::new(scheduled_time, event, sequence));
849
850        tracing::trace!(
851            "Scheduled set_len: file={:?}, new_len={}, op_seq={}",
852            file_id,
853            new_len,
854            op_seq
855        );
856
857        Ok(op_seq)
858    }
859
860    /// Check if a storage operation is complete.
861    pub(crate) fn is_storage_op_complete(&self, file_id: FileId, op_seq: u64) -> bool {
862        let inner = self
863            .inner
864            .read()
865            .expect("RwLock poisoned: prior task panicked");
866        if let Some(file_state) = inner.storage.files.get(&file_id) {
867            // Operation is complete when it's no longer in pending_ops
868            !file_state.pending_ops.contains_key(&op_seq)
869        } else {
870            // File not found means operation is effectively "complete" (failed)
871            true
872        }
873    }
874
875    /// Check if a sync operation failed and clear the failure flag.
876    ///
877    /// Returns true if the sync failed due to fault injection.
878    pub(crate) fn take_sync_failure(&self, file_id: FileId, op_seq: u64) -> bool {
879        let mut inner = self
880            .inner
881            .write()
882            .expect("RwLock poisoned: prior task panicked");
883        inner.storage.sync_failures.remove(&(file_id, op_seq))
884    }
885
886    /// Register a waker for a storage operation.
887    pub(crate) fn register_storage_waker(&self, file_id: FileId, op_seq: u64, waker: Waker) {
888        let mut inner = self
889            .inner
890            .write()
891            .expect("RwLock poisoned: prior task panicked");
892        inner.wakers.storage_ops.insert((file_id, op_seq), waker);
893    }
894
895    /// Read data from a file at the given offset.
896    ///
897    /// This is called after `ReadComplete` to actually fetch the data.
898    pub(crate) fn read_from_file(
899        &self,
900        file_id: FileId,
901        offset: u64,
902        buf: &mut [u8],
903    ) -> Result<usize, StorageError> {
904        let inner = self
905            .inner
906            .read()
907            .expect("RwLock poisoned: prior task panicked");
908
909        let file_state = inner
910            .storage
911            .files
912            .get(&file_id)
913            .ok_or(StorageError::InvalidFileHandle { file_id })?;
914
915        if file_state.is_closed {
916            return Err(StorageError::FileClosed { file_id });
917        }
918
919        // Read from the in-memory storage
920        file_state
921            .storage
922            .read(offset, buf)
923            .map_err(|e| StorageError::Io {
924                file_id,
925                kind: e.kind(),
926                message: e.to_string(),
927            })?;
928
929        Ok(buf.len())
930    }
931
932    /// Get the current file position.
933    pub(crate) fn file_position(&self, file_id: FileId) -> Result<u64, StorageError> {
934        let inner = self
935            .inner
936            .read()
937            .expect("RwLock poisoned: prior task panicked");
938        inner
939            .storage
940            .files
941            .get(&file_id)
942            .map(|f| f.position)
943            .ok_or(StorageError::InvalidFileHandle { file_id })
944    }
945
946    /// Set the current file position.
947    pub(crate) fn set_file_position(
948        &self,
949        file_id: FileId,
950        position: u64,
951    ) -> Result<(), StorageError> {
952        let mut inner = self
953            .inner
954            .write()
955            .expect("RwLock poisoned: prior task panicked");
956        if let Some(file_state) = inner.storage.files.get_mut(&file_id) {
957            file_state.position = position;
958            Ok(())
959        } else {
960            Err(StorageError::InvalidFileHandle { file_id })
961        }
962    }
963
964    /// Get the size of a file.
965    pub(crate) fn file_size(&self, file_id: FileId) -> Result<u64, StorageError> {
966        let inner = self
967            .inner
968            .read()
969            .expect("RwLock poisoned: prior task panicked");
970        inner
971            .storage
972            .files
973            .get(&file_id)
974            .map(|f| f.storage.size())
975            .ok_or(StorageError::InvalidFileHandle { file_id })
976    }
977
978    /// Calculate storage latency using FDB formula, adjusted for any active
979    /// disk-degradation episode.
980    ///
981    /// Latency = `base_latency` + `iops_overhead` + `transfer_time`, where a
982    /// **throttle** episode divides effective IOPS/bandwidth by the configured
983    /// multipliers and a **stall** episode adds the time remaining until the
984    /// stall expires (the disk is frozen until then).
985    fn calculate_storage_latency(
986        config: &crate::storage::StorageConfiguration,
987        size: usize,
988        is_write: bool,
989        episode: Option<DiskDegradationState>,
990        now: Duration,
991    ) -> Duration {
992        // Sample base latency from config range
993        let base_range = if is_write {
994            &config.write_latency
995        } else {
996            &config.read_latency
997        };
998        let base = crate::network::sample_latency(base_range);
999
1000        // During a throttle episode, effective IOPS/bandwidth are divided by the
1001        // configured multipliers (>= 1.0); otherwise the divisors are 1.0.
1002        let (iops_divisor, bandwidth_divisor) = match episode {
1003            Some(DiskDegradationState {
1004                kind: DiskEpisodeKind::Throttle,
1005                ..
1006            }) => (
1007                config.disk_throttle_iops_multiplier.max(1.0),
1008                config.disk_throttle_bandwidth_multiplier.max(1.0),
1009            ),
1010            _ => (1.0, 1.0),
1011        };
1012
1013        // IOPS overhead: divisor/iops seconds per operation.
1014        // Precision loss is acceptable: iops is typically << 2^52.
1015        let iops_f64 = u32::try_from(config.iops).map_or(f64::from(u32::MAX), f64::from);
1016        let iops_overhead = Duration::from_secs_f64(iops_divisor / iops_f64);
1017
1018        // Transfer time: size * divisor / bandwidth seconds.
1019        // Precision loss is acceptable: simulated sizes/bandwidths fit in 2^52.
1020        let size_f64 = u32::try_from(size).map_or(f64::from(u32::MAX), f64::from);
1021        let bandwidth_f64 = u32::try_from(config.bandwidth).map_or(f64::from(u32::MAX), f64::from);
1022        let transfer = Duration::from_secs_f64(size_f64 * bandwidth_divisor / bandwidth_f64);
1023
1024        let steady = base + iops_overhead + transfer;
1025
1026        // During a stall the disk is frozen until expiry: the op waits out the
1027        // remaining window, then takes its normal steady-state latency.
1028        match episode {
1029            Some(DiskDegradationState {
1030                kind: DiskEpisodeKind::Stall,
1031                expires_at,
1032            }) => steady + expires_at.saturating_sub(now),
1033            _ => steady,
1034        }
1035    }
1036
1037    /// Simulate a crash affecting storage for a specific process.
1038    ///
1039    /// Only affects files owned by the given IP address:
1040    /// 1. Calls `apply_crash()` on matching `InMemoryStorage` instances
1041    /// 2. Clears pending operations (lost in crash)
1042    /// 3. Optionally marks files as closed
1043    /// 4. Wakes all storage wakers (operations will fail)
1044    ///
1045    /// Files owned by other IPs are unaffected.
1046    ///
1047    /// # Panics
1048    ///
1049    /// Panics if the simulation lock is poisoned by a prior task panic.
1050    #[instrument(skip(self))]
1051    pub fn simulate_crash_for_process(&self, ip: IpAddr, close_files: bool) {
1052        let mut inner = self
1053            .inner
1054            .write()
1055            .expect("RwLock poisoned: prior task panicked");
1056        let crash_probability = inner.storage.config_for(ip).crash_fault_probability;
1057
1058        // Collect all wakers to wake in one pass (to avoid borrow conflict)
1059        let mut wakers_to_wake = Vec::new();
1060        let file_ids: Vec<FileId> = inner
1061            .storage
1062            .files
1063            .iter()
1064            .filter(|(_, f)| f.owner_ip == ip)
1065            .map(|(id, _)| *id)
1066            .collect();
1067
1068        for file_id in &file_ids {
1069            if let Some(file_state) = inner.storage.files.get_mut(file_id) {
1070                // Apply crash to in-memory storage (may corrupt pending writes)
1071                file_state.storage.apply_crash(crash_probability);
1072
1073                // Collect lost op sequence numbers
1074                let lost_ops: Vec<u64> = file_state.pending_ops.keys().copied().collect();
1075
1076                // Clear pending ops - they're lost in crash
1077                file_state.pending_ops.clear();
1078
1079                // Collect waker keys for later removal
1080                for op_seq in lost_ops {
1081                    wakers_to_wake.push((*file_id, op_seq));
1082                }
1083
1084                // Optionally close files
1085                if close_files {
1086                    file_state.is_closed = true;
1087                }
1088            }
1089        }
1090
1091        // Wake all collected wakers (after file iteration is complete)
1092        for key in wakers_to_wake {
1093            if let Some(waker) = inner.wakers.storage_ops.remove(&key) {
1094                waker.wake();
1095            }
1096        }
1097
1098        inner.record_fault(SimFaultEvent::StorageCrash { ip: ip.to_string() });
1099
1100        tracing::info!(
1101            "Storage crash simulated for {}: {} files affected, close_files={}",
1102            ip,
1103            file_ids.len(),
1104            close_files
1105        );
1106    }
1107
1108    /// Wipe all storage for a specific process.
1109    ///
1110    /// Deletes all files owned by the given IP address. Used by `CrashAndWipe`
1111    /// reboot to simulate total data loss. After wipe, the process can create
1112    /// new files at the same paths.
1113    ///
1114    /// Files owned by other IPs are unaffected.
1115    ///
1116    /// # Panics
1117    ///
1118    /// Panics if the simulation lock is poisoned by a prior task panic.
1119    #[instrument(skip(self))]
1120    pub fn wipe_storage_for_process(&self, ip: IpAddr) {
1121        let mut inner = self
1122            .inner
1123            .write()
1124            .expect("RwLock poisoned: prior task panicked");
1125
1126        // Collect files owned by this IP
1127        let file_ids: Vec<(FileId, String)> = inner
1128            .storage
1129            .files
1130            .iter()
1131            .filter(|(_, f)| f.owner_ip == ip)
1132            .map(|(id, f)| (*id, f.path.clone()))
1133            .collect();
1134
1135        // Collect wakers to wake
1136        let mut wakers_to_wake = Vec::new();
1137
1138        for (file_id, path) in &file_ids {
1139            if let Some(file_state) = inner.storage.files.remove(file_id) {
1140                for op_seq in file_state.pending_ops.keys() {
1141                    wakers_to_wake.push((*file_id, *op_seq));
1142                }
1143            }
1144            inner.storage.path_to_file.remove(path);
1145            inner.storage.deleted_paths.insert(path.clone());
1146        }
1147
1148        // Wake all collected wakers
1149        for key in wakers_to_wake {
1150            if let Some(waker) = inner.wakers.storage_ops.remove(&key) {
1151                waker.wake();
1152            }
1153        }
1154
1155        inner.record_fault(SimFaultEvent::StorageWipe { ip: ip.to_string() });
1156
1157        tracing::info!("Storage wiped for {}: {} files deleted", ip, file_ids.len(),);
1158    }
1159
1160    /// Set storage configuration for a specific process.
1161    ///
1162    /// Files owned by this IP will use this configuration for fault injection
1163    /// and latency calculations. Takes effect immediately, even for files
1164    /// already open.
1165    ///
1166    /// # Panics
1167    ///
1168    /// Panics if the simulation lock is poisoned by a prior task panic.
1169    #[instrument(skip(self, config))]
1170    pub fn set_process_storage_config(
1171        &self,
1172        ip: IpAddr,
1173        config: crate::storage::StorageConfiguration,
1174    ) {
1175        let mut inner = self
1176            .inner
1177            .write()
1178            .expect("RwLock poisoned: prior task panicked");
1179        inner.storage.per_process_configs.insert(ip, config);
1180    }
1181}