Skip to main content

moonpool_sim/storage/
file.rs

1//! Simulated storage file implementation.
2
3use crate::sim::WeakSimWorld;
4use crate::sim::state::FileId;
5use futures::io::{AsyncRead, AsyncSeek, AsyncWrite};
6use moonpool_core::StorageFile;
7use std::io::{self, SeekFrom};
8use std::pin::Pin;
9use std::sync::Mutex;
10use std::task::{Context, Poll};
11
12use super::futures::{SetLenFuture, SyncFuture};
13use super::sim_shutdown_error;
14
15/// Simulated storage file for deterministic testing.
16///
17/// This provides a simulation-aware file handle that integrates with
18/// the deterministic simulation engine for testing storage I/O patterns.
19///
20/// ## State Tracking
21///
22/// The file tracks pending operations under a `Mutex`:
23/// - `pending_read`: Active read operation (`op_seq`, offset, len)
24/// - `pending_write`: Active write operation (`op_seq`, `bytes_written`)
25///
26/// ## Polling Pattern
27///
28/// Operations follow the schedule → wait → complete pattern:
29/// 1. First poll: Schedule operation with `SimWorld`, store pending state
30/// 2. Subsequent polls: Check completion, return Pending until done
31/// 3. Final poll: Clear pending state, return result
32#[derive(Debug)]
33pub struct SimStorageFile {
34    sim: WeakSimWorld,
35    file_id: FileId,
36    /// Pending read operation: (`op_seq`, offset, len)
37    pending_read: Mutex<Option<(u64, u64, usize)>>,
38    /// Pending write operation: (`op_seq`, `bytes_written`)
39    pending_write: Mutex<Option<(u64, usize)>>,
40}
41
42impl SimStorageFile {
43    /// Create a new simulated storage file.
44    pub(crate) fn new(sim: WeakSimWorld, file_id: FileId) -> Self {
45        Self {
46            sim,
47            file_id,
48            pending_read: Mutex::new(None),
49            pending_write: Mutex::new(None),
50        }
51    }
52}
53
54impl StorageFile for SimStorageFile {
55    async fn sync_all(&self) -> io::Result<()> {
56        SyncFuture::new(self.sim.clone(), self.file_id).await
57    }
58
59    async fn sync_data(&self) -> io::Result<()> {
60        // Simulation treats sync_all and sync_data identically
61        SyncFuture::new(self.sim.clone(), self.file_id).await
62    }
63
64    async fn size(&self) -> io::Result<u64> {
65        let sim = self.sim.upgrade().map_err(|_| sim_shutdown_error())?;
66        sim.file_size(self.file_id)
67            .map_err(|e| io::Error::other(e.to_string()))
68    }
69
70    async fn set_len(&self, size: u64) -> io::Result<()> {
71        SetLenFuture::new(self.sim.clone(), self.file_id, size).await
72    }
73}
74
75impl AsyncRead for SimStorageFile {
76    fn poll_read(
77        self: Pin<&mut Self>,
78        cx: &mut Context<'_>,
79        buf: &mut [u8],
80    ) -> Poll<io::Result<usize>> {
81        let sim = self.sim.upgrade().map_err(|_| sim_shutdown_error())?;
82
83        // Check for pending read operation
84        let pending = *self
85            .pending_read
86            .lock()
87            .expect("Mutex poisoned: prior task panicked");
88        if let Some((op_seq, offset, len)) = pending {
89            // Check if operation is complete
90            if sim.is_storage_op_complete(self.file_id, op_seq) {
91                // Clear pending state
92                *self
93                    .pending_read
94                    .lock()
95                    .expect("Mutex poisoned: prior task panicked") = None;
96
97                // Calculate how many bytes to actually read
98                let bytes_to_read = buf.len().min(len);
99                if bytes_to_read == 0 {
100                    return Poll::Ready(Ok(0));
101                }
102
103                // Read from file at the stored offset
104                let bytes_read =
105                    sim.read_from_file(self.file_id, offset, &mut buf[..bytes_to_read])?;
106
107                // Update file position
108                let new_position = offset + bytes_read as u64;
109                sim.set_file_position(self.file_id, new_position)?;
110
111                return Poll::Ready(Ok(bytes_read));
112            }
113
114            // Operation not complete, register waker and wait
115            sim.register_storage_waker(self.file_id, op_seq, cx.waker().clone());
116            return Poll::Pending;
117        }
118
119        // No pending read - start a new one
120
121        // Get current position
122        let position = sim.file_position(self.file_id)?;
123
124        // Get file size to check for EOF
125        let file_size = sim.file_size(self.file_id)?;
126
127        // Check for EOF
128        if position >= file_size {
129            return Poll::Ready(Ok(0)); // EOF - 0 bytes read
130        }
131
132        // Calculate bytes to read (don't read past EOF)
133        let remaining_in_file =
134            usize::try_from(file_size - position).expect("remaining bytes in file fit in usize");
135        let len = buf.len().min(remaining_in_file);
136
137        if len == 0 {
138            return Poll::Ready(Ok(0));
139        }
140
141        // Schedule the read operation
142        let op_seq = sim.schedule_read(self.file_id, position, len)?;
143
144        // Store pending state
145        *self
146            .pending_read
147            .lock()
148            .expect("Mutex poisoned: prior task panicked") = Some((op_seq, position, len));
149
150        // Register waker
151        sim.register_storage_waker(self.file_id, op_seq, cx.waker().clone());
152
153        Poll::Pending
154    }
155}
156
157impl AsyncWrite for SimStorageFile {
158    fn poll_write(
159        self: Pin<&mut Self>,
160        cx: &mut Context<'_>,
161        buf: &[u8],
162    ) -> Poll<io::Result<usize>> {
163        let sim = self.sim.upgrade().map_err(|_| sim_shutdown_error())?;
164
165        // Check for pending write operation
166        let pending = *self
167            .pending_write
168            .lock()
169            .expect("Mutex poisoned: prior task panicked");
170        if let Some((op_seq, bytes_written)) = pending {
171            // Check if operation is complete
172            if sim.is_storage_op_complete(self.file_id, op_seq) {
173                // Clear pending state
174                *self
175                    .pending_write
176                    .lock()
177                    .expect("Mutex poisoned: prior task panicked") = None;
178
179                // Update file position
180                let position = sim.file_position(self.file_id)?;
181                let new_position = position + bytes_written as u64;
182                sim.set_file_position(self.file_id, new_position)?;
183
184                return Poll::Ready(Ok(bytes_written));
185            }
186
187            // Operation not complete, register waker and wait
188            sim.register_storage_waker(self.file_id, op_seq, cx.waker().clone());
189            return Poll::Pending;
190        }
191
192        // No pending write - start a new one
193
194        if buf.is_empty() {
195            return Poll::Ready(Ok(0));
196        }
197
198        // Get current position
199        let position = sim.file_position(self.file_id)?;
200
201        // Schedule the write operation
202        let op_seq = sim.schedule_write(self.file_id, position, buf.to_vec())?;
203
204        // Store pending state
205        *self
206            .pending_write
207            .lock()
208            .expect("Mutex poisoned: prior task panicked") = Some((op_seq, buf.len()));
209
210        // Register waker
211        sim.register_storage_waker(self.file_id, op_seq, cx.waker().clone());
212
213        Poll::Pending
214    }
215
216    fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<io::Result<()>> {
217        // Flush is a no-op - durability comes from sync_all/sync_data
218        Poll::Ready(Ok(()))
219    }
220
221    fn poll_close(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<io::Result<()>> {
222        // Close is a no-op - file cleanup handled via Drop if needed
223        Poll::Ready(Ok(()))
224    }
225}
226
227impl AsyncSeek for SimStorageFile {
228    fn poll_seek(
229        self: Pin<&mut Self>,
230        _cx: &mut Context<'_>,
231        pos: SeekFrom,
232    ) -> Poll<io::Result<u64>> {
233        let sim = self.sim.upgrade().map_err(|_| sim_shutdown_error())?;
234
235        let current_position = sim.file_position(self.file_id)?;
236        let file_size = sim.file_size(self.file_id)?;
237
238        let target = match pos {
239            SeekFrom::Start(p) => p,
240            SeekFrom::End(offset) => {
241                if offset >= 0 {
242                    file_size.saturating_add(offset.unsigned_abs())
243                } else {
244                    file_size.saturating_sub(offset.unsigned_abs())
245                }
246            }
247            SeekFrom::Current(offset) => {
248                if offset >= 0 {
249                    current_position.saturating_add(offset.unsigned_abs())
250                } else {
251                    current_position.saturating_sub(offset.unsigned_abs())
252                }
253            }
254        };
255
256        sim.set_file_position(self.file_id, target)?;
257        Poll::Ready(Ok(target))
258    }
259}