1use 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
25fn 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
49fn 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 if let Some(episode) = episodes.get(&owner_ip).copied()
74 && now >= episode.expires_at
75 {
76 episodes.remove(&owner_ip);
77 }
78
79 if let Some(episode) = episodes.get(&owner_ip).copied() {
81 return Some(episode);
82 }
83
84 if stall_probability <= 0.0 && throttle_probability <= 0.0 {
86 return None;
87 }
88
89 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
112fn 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
125pub(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
155fn 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 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 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 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
207fn 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 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 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 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 else if sim_random::<f64>() < config.misdirect_write_probability {
248 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 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 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 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
309fn 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 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 if sim_random::<f64>() < sync_failure_prob {
328 tracing::info!("Sync failure injected for file {:?}", file_id);
329 inner.storage.sync_failures.insert((file_id, op_seq));
331 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 file_state.storage.sync();
343 }
344
345 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
352fn handle_open_complete(inner: &mut SimInner, file_id: FileId) {
354 let Some((op_seq, _)) = take_pending_op(inner, file_id, PendingOpType::Open) else {
356 tracing::trace!("OpenComplete for file {:?} (no pending op)", file_id);
358 return;
359 };
360
361 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
368fn handle_set_len_complete(inner: &mut SimInner, file_id: FileId, new_len: u64) {
370 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 if let Some(file_state) = inner.storage.files.get_mut(&file_id) {
378 file_state.storage.resize(new_len);
379 }
380
381 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
393impl SimWorld {
398 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 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 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 if inner.storage.deleted_paths.contains(&path_str) && !options.is_create() {
441 return Err(StorageError::NotFound { path: path_str });
442 }
443
444 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 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 file_state.position = file_state.storage.size();
455 } else {
456 file_state.position = 0;
458 }
459 file_state.options = options;
461 file_state.is_closed = false;
462 }
463 return Ok(existing_id);
464 }
465
466 if !options.is_create() && !options.is_create_new() {
468 return Err(StorageError::NotFound { path: path_str });
469 }
470
471 let file_id = FileId(inner.storage.next_file_id);
473 inner.storage.next_file_id += 1;
474
475 inner.storage.deleted_paths.remove(&path_str);
477
478 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 !file_state.pending_ops.contains_key(&op_seq)
869 } else {
870 true
872 }
873 }
874
875 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 file_state.storage.apply_crash(crash_probability);
1072
1073 let lost_ops: Vec<u64> = file_state.pending_ops.keys().copied().collect();
1075
1076 file_state.pending_ops.clear();
1078
1079 for op_seq in lost_ops {
1081 wakers_to_wake.push((*file_id, op_seq));
1082 }
1083
1084 if close_files {
1086 file_state.is_closed = true;
1087 }
1088 }
1089 }
1090
1091 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 #[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 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 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 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 #[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}