Skip to main content

storage_device/
fake_device.rs

1// Copyright 2021 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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
31/// A Device backed by a memory buffer.
32pub 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    /// Sets a callback that will run at the beginning of read, write, and flush which will forward
67    /// any errors, and proceed on Ok().
68    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    /// Creates a fake block device from an image (which can be anything that implements
76    /// std::io::Read).  The size of the device is determined by how much data is read.
77    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    /// Creates a fake block device from a `Vec`. The size of the device is determined by the size
87    /// of the `Vec`.
88    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        // Blast over the range to simulate it being used for something else.
184        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    /// Sets the poisoned state for the device. A poisoned device will panic the thread that
245    /// performs Drop on it.
246    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        // Loop 100 times to catch errors.
272        for _ in 0..1000 {
273            let mut data = vec![0; 7 * TEST_DEVICE_BLOCK_SIZE];
274            rand::fill(&mut data[..]);
275            // Ensure that barriers work with overwrites.
276            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}