moonpool_sim/storage/
file.rs1use 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#[derive(Debug)]
33pub struct SimStorageFile {
34 sim: WeakSimWorld,
35 file_id: FileId,
36 pending_read: Mutex<Option<(u64, u64, usize)>>,
38 pending_write: Mutex<Option<(u64, usize)>>,
40}
41
42impl SimStorageFile {
43 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 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 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 if sim.is_storage_op_complete(self.file_id, op_seq) {
91 *self
93 .pending_read
94 .lock()
95 .expect("Mutex poisoned: prior task panicked") = None;
96
97 let bytes_to_read = buf.len().min(len);
99 if bytes_to_read == 0 {
100 return Poll::Ready(Ok(0));
101 }
102
103 let bytes_read =
105 sim.read_from_file(self.file_id, offset, &mut buf[..bytes_to_read])?;
106
107 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 sim.register_storage_waker(self.file_id, op_seq, cx.waker().clone());
116 return Poll::Pending;
117 }
118
119 let position = sim.file_position(self.file_id)?;
123
124 let file_size = sim.file_size(self.file_id)?;
126
127 if position >= file_size {
129 return Poll::Ready(Ok(0)); }
131
132 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 let op_seq = sim.schedule_read(self.file_id, position, len)?;
143
144 *self
146 .pending_read
147 .lock()
148 .expect("Mutex poisoned: prior task panicked") = Some((op_seq, position, len));
149
150 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 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 if sim.is_storage_op_complete(self.file_id, op_seq) {
173 *self
175 .pending_write
176 .lock()
177 .expect("Mutex poisoned: prior task panicked") = None;
178
179 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 sim.register_storage_waker(self.file_id, op_seq, cx.waker().clone());
189 return Poll::Pending;
190 }
191
192 if buf.is_empty() {
195 return Poll::Ready(Ok(0));
196 }
197
198 let position = sim.file_position(self.file_id)?;
200
201 let op_seq = sim.schedule_write(self.file_id, position, buf.to_vec())?;
203
204 *self
206 .pending_write
207 .lock()
208 .expect("Mutex poisoned: prior task panicked") = Some((op_seq, buf.len()));
209
210 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 Poll::Ready(Ok(()))
219 }
220
221 fn poll_close(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<io::Result<()>> {
222 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}