1use crate::buffer::{BufferFuture, BufferRef, MutableBufferRef};
6use crate::buffer_allocator::{BufferAllocator, BufferSource};
7use crate::{Device, DeviceHolder};
8use anyhow::{Error, ensure};
9use async_trait::async_trait;
10use block_protocol::{ReadOptions, WriteFlags, WriteOptions};
11use fuchsia_sync::Mutex;
12use std::ops::Range;
13use std::sync::atomic::{AtomicBool, Ordering};
14
15pub enum Op {
16 Read,
17 Write,
18 Flush,
19}
20
21pub trait Observer: Send + Sync {
22 fn barrier(&self) {}
23}
24
25#[derive(Debug, Default, Clone)]
26struct Inner {
27 data: Vec<u8>,
28 blocks_written_since_last_barrier: Vec<usize>,
29}
30
31pub struct FakeDevice {
33 allocator: BufferAllocator,
34 inner: Mutex<Inner>,
35 closed: AtomicBool,
36 operation_closure: Box<dyn Fn(Op) -> Result<(), Error> + Send + Sync>,
37 read_only: AtomicBool,
38 poisoned: AtomicBool,
39 observer: Option<Box<dyn Observer>>,
40}
41
42const TRANSFER_HEAP_SIZE: usize = 64 * 1024 * 1024;
43
44impl FakeDevice {
45 pub fn new(block_count: u64, block_size: u32) -> Self {
46 let allocator =
47 BufferAllocator::new(block_size as usize, BufferSource::new(TRANSFER_HEAP_SIZE));
48 Self {
49 allocator,
50 inner: Mutex::new(Inner {
51 data: vec![0 as u8; block_count as usize * block_size as usize],
52 blocks_written_since_last_barrier: Vec::new(),
53 }),
54 closed: AtomicBool::new(false),
55 operation_closure: Box::new(|_: Op| Ok(())),
56 read_only: AtomicBool::new(false),
57 poisoned: AtomicBool::new(false),
58 observer: None,
59 }
60 }
61
62 pub fn set_observer(&mut self, observer: Box<dyn Observer>) {
63 self.observer = Some(observer);
64 }
65
66 pub fn set_op_callback(
69 &mut self,
70 cb: impl Fn(Op) -> Result<(), Error> + Send + Sync + 'static,
71 ) {
72 self.operation_closure = Box::new(cb);
73 }
74
75 pub fn from_image(
78 mut reader: impl std::io::Read,
79 block_size: u32,
80 ) -> Result<Self, std::io::Error> {
81 let mut data = Vec::new();
82 reader.read_to_end(&mut data)?;
83 Ok(Self::from_vec(data, block_size))
84 }
85
86 pub fn from_vec(data: Vec<u8>, block_size: u32) -> Self {
89 let allocator =
90 BufferAllocator::new(block_size as usize, BufferSource::new(TRANSFER_HEAP_SIZE));
91 Self {
92 allocator,
93 inner: Mutex::new(Inner { data, blocks_written_since_last_barrier: Vec::new() }),
94 closed: AtomicBool::new(false),
95 operation_closure: Box::new(|_| Ok(())),
96 read_only: AtomicBool::new(false),
97 poisoned: AtomicBool::new(false),
98 observer: None,
99 }
100 }
101}
102
103#[async_trait]
104impl Device for FakeDevice {
105 fn allocate_buffer(&self, size: usize) -> BufferFuture<'_> {
106 assert!(!self.closed.load(Ordering::Relaxed));
107 self.allocator.allocate_buffer(size)
108 }
109
110 fn block_size(&self) -> u32 {
111 self.allocator.block_size() as u32
112 }
113
114 fn block_count(&self) -> u64 {
115 self.inner.lock().data.len() as u64 / self.block_size() as u64
116 }
117
118 async fn read_with_opts(
119 &self,
120 offset: u64,
121 mut buffer: MutableBufferRef<'_>,
122 _read_opts: ReadOptions,
123 ) -> Result<(), Error> {
124 ensure!(!self.closed.load(Ordering::Relaxed));
125 (self.operation_closure)(Op::Read)?;
126 let offset = offset as usize;
127 assert_eq!(offset % self.allocator.block_size(), 0);
128 let inner = self.inner.lock();
129 let size = buffer.len();
130 ensure!(
131 offset + size <= inner.data.len(),
132 "offset: {} len: {} data.len: {}",
133 offset,
134 size,
135 inner.data.len()
136 );
137 buffer.copy_from_slice(&inner.data[offset..offset + size]);
138 Ok(())
139 }
140
141 async fn write_with_opts(
142 &self,
143 offset: u64,
144 buffer: BufferRef<'_>,
145 write_opts: WriteOptions,
146 ) -> Result<(), Error> {
147 ensure!(!self.closed.load(Ordering::Relaxed));
148 ensure!(!self.read_only.load(Ordering::Relaxed));
149 let mut inner = self.inner.lock();
150
151 if write_opts.flags.contains(WriteFlags::PRE_BARRIER) {
152 if let Some(observer) = &self.observer {
153 observer.barrier();
154 }
155 inner.blocks_written_since_last_barrier.clear();
156 }
157
158 (self.operation_closure)(Op::Write)?;
159 let offset = offset as usize;
160 assert_eq!(offset % self.allocator.block_size(), 0);
161
162 let size = buffer.len();
163 ensure!(
164 offset + size <= inner.data.len(),
165 "offset: {} len: {} data.len: {}",
166 offset,
167 size,
168 inner.data.len()
169 );
170 buffer.copy_to_slice(&mut inner.data[offset..offset + size]);
171 let first_block = offset / self.allocator.block_size();
172 for block in first_block..first_block + size / self.allocator.block_size() {
173 inner.blocks_written_since_last_barrier.push(block)
174 }
175 Ok(())
176 }
177
178 async fn trim(&self, range: Range<u64>) -> Result<(), Error> {
179 ensure!(!self.closed.load(Ordering::Relaxed));
180 ensure!(!self.read_only.load(Ordering::Relaxed));
181 assert_eq!(range.start % self.block_size() as u64, 0);
182 assert_eq!(range.end % self.block_size() as u64, 0);
183 let mut inner = self.inner.lock();
185 inner.data[range.start as usize..range.end as usize].fill(0xab);
186 Ok(())
187 }
188
189 async fn close(&self) -> Result<(), Error> {
190 self.closed.store(true, Ordering::Relaxed);
191 Ok(())
192 }
193
194 async fn flush(&self) -> Result<(), Error> {
195 self.inner.lock().blocks_written_since_last_barrier.clear();
196 (self.operation_closure)(Op::Flush)
197 }
198
199 fn reopen(&self, read_only: bool) {
200 self.closed.store(false, Ordering::Relaxed);
201 self.read_only.store(read_only, Ordering::Relaxed);
202 }
203
204 fn is_read_only(&self) -> bool {
205 self.read_only.load(Ordering::Relaxed)
206 }
207
208 fn supports_trim(&self) -> bool {
209 true
210 }
211
212 fn snapshot(&self) -> Result<DeviceHolder, Error> {
213 let allocator =
214 BufferAllocator::new(self.block_size() as usize, BufferSource::new(TRANSFER_HEAP_SIZE));
215 Ok(DeviceHolder::new(Self {
216 allocator,
217 inner: Mutex::new(self.inner.lock().clone()),
218 closed: AtomicBool::new(false),
219 operation_closure: Box::new(|_: Op| Ok(())),
220 read_only: AtomicBool::new(false),
221 poisoned: AtomicBool::new(false),
222 observer: None,
223 }))
224 }
225
226 fn discard_random_since_last_flush(&self) -> Result<(), Error> {
227 let bs = self.allocator.block_size();
228 use rand::RngExt as _;
229 let mut rng = rand::rng();
230 let mut guard = self.inner.lock();
231 let Inner { data, blocks_written_since_last_barrier, .. } = &mut *guard;
232 log::info!("Discarding from {blocks_written_since_last_barrier:?}");
233 let mut discarded = Vec::new();
234 for block in blocks_written_since_last_barrier.drain(..) {
235 if rng.random() {
236 data[block * bs..(block + 1) * bs].fill(0xaf);
237 discarded.push(block);
238 }
239 }
240 log::info!("Discarded {discarded:?}");
241 Ok(())
242 }
243
244 fn poison(&self) -> Result<(), Error> {
247 self.poisoned.store(true, Ordering::Relaxed);
248 Ok(())
249 }
250}
251
252impl Drop for FakeDevice {
253 fn drop(&mut self) {
254 if self.poisoned.load(Ordering::Relaxed) {
255 panic!("This device was poisoned to crash whomever is holding a reference here.");
256 }
257 }
258}
259
260#[cfg(test)]
261mod tests {
262 use super::FakeDevice;
263 use crate::Device;
264 use block_protocol::{WriteFlags, WriteOptions};
265
266 const TEST_DEVICE_BLOCK_SIZE: usize = 512;
267
268 #[fuchsia::test(threads = 10)]
269 async fn test_discard_random_with_barriers() {
270 let device = FakeDevice::new(8192, TEST_DEVICE_BLOCK_SIZE as u32);
271 for _ in 0..1000 {
273 let mut data = vec![0; 7 * TEST_DEVICE_BLOCK_SIZE];
274 rand::fill(&mut data[..]);
275 let indices = [1, 2, 3, 4, 3, 5, 6];
277 for i in 0..indices.len() {
278 let mut buffer = device.allocate_buffer(TEST_DEVICE_BLOCK_SIZE).await;
279 if i == 2 || i == 5 {
280 buffer.copy_from_slice(
281 &data[indices[i] * TEST_DEVICE_BLOCK_SIZE
282 ..indices[i] * TEST_DEVICE_BLOCK_SIZE + TEST_DEVICE_BLOCK_SIZE],
283 );
284 device
285 .write_with_opts(
286 i as u64 * TEST_DEVICE_BLOCK_SIZE as u64,
287 buffer.as_ref(),
288 WriteOptions { flags: WriteFlags::PRE_BARRIER, ..Default::default() },
289 )
290 .await
291 .expect("Failed to write to FakeDevice");
292 } else {
293 buffer.copy_from_slice(
294 &data[indices[i] * TEST_DEVICE_BLOCK_SIZE
295 ..indices[i] * TEST_DEVICE_BLOCK_SIZE + TEST_DEVICE_BLOCK_SIZE],
296 );
297 device
298 .write_with_opts(
299 i as u64 * TEST_DEVICE_BLOCK_SIZE as u64,
300 buffer.as_ref(),
301 WriteOptions::default(),
302 )
303 .await
304 .expect("Failed to write to FakeDevice");
305 }
306 }
307 device.discard_random_since_last_flush().expect("failed to randomly discard writes");
308 let mut discard = false;
309 let mut discard_2 = false;
310 for i in 0..7 {
311 let mut read_buffer = device.allocate_buffer(TEST_DEVICE_BLOCK_SIZE).await;
312 device
313 .read(i as u64 * TEST_DEVICE_BLOCK_SIZE as u64, read_buffer.as_mut())
314 .await
315 .expect("failed to read from FakeDevice");
316 let mut read_data = vec![0u8; TEST_DEVICE_BLOCK_SIZE];
317 read_buffer.copy_to_slice(&mut read_data);
318 let expected_data = &data[indices[i] * TEST_DEVICE_BLOCK_SIZE
319 ..indices[i] * TEST_DEVICE_BLOCK_SIZE + TEST_DEVICE_BLOCK_SIZE];
320 if i < 2 {
321 if expected_data != &read_data[..] {
322 discard = true;
323 }
324 } else if i < 5 {
325 if discard == true {
326 assert_ne!(expected_data, &read_data[..]);
327 discard_2 = true;
328 } else if expected_data != &read_data[..] {
329 discard_2 = true;
330 }
331 } else if discard_2 == true {
332 assert_ne!(expected_data, &read_data[..]);
333 }
334 }
335 }
336 }
337}