Skip to main content

moirai_async/io/
positioned.rs

1//! Positioned asynchronous read primitives.
2//!
3//! These contracts complement [`super::AsyncRead`], whose stateful stream
4//! cursor is not sufficient for format readers that issue concurrent reads at
5//! explicit file offsets. They intentionally use `std::io::Result` so the
6//! runtime remains independent of any consumer's error hierarchy.
7
8use std::future::Future;
9use std::io;
10
11/// Read exactly the requested bytes at an absolute offset without changing a
12/// shared cursor.
13pub trait AsyncReadAt: Send + Sync {
14    /// Read exactly `buf.len()` bytes beginning at `offset`.
15    fn read_at(&self, offset: u64, buf: &mut [u8]) -> impl Future<Output = io::Result<()>> + Send;
16}
17
18/// Query the byte length of an asynchronous positioned-read source.
19pub trait AsyncLength: Send + Sync {
20    /// Return the total number of readable bytes.
21    fn len(&self) -> impl Future<Output = io::Result<u64>> + Send;
22
23    /// Whether the source has no readable bytes.
24    ///
25    /// Provided rather than required: a zero length is the only way to be
26    /// empty, so an implementor that answers `len` has already answered this.
27    /// It exists because a caller checking for an empty source should not have
28    /// to compare against zero itself, and because a `len` without an
29    /// `is_empty` is a trait-shape defect that `clippy::len_without_is_empty`
30    /// correctly rejects.
31    fn is_empty(&self) -> impl Future<Output = io::Result<bool>> + Send {
32        async move { Ok(self.len().await? == 0) }
33    }
34}
35
36/// A cloneable, read-only in-memory positioned source for tests and examples.
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct AsyncMemReader {
39    data: Vec<u8>,
40}
41
42impl AsyncMemReader {
43    /// Construct a reader from owned bytes.
44    #[must_use]
45    pub fn from_bytes(data: Vec<u8>) -> Self {
46        Self { data }
47    }
48
49    /// Construct an empty reader.
50    #[must_use]
51    pub fn new() -> Self {
52        Self::from_bytes(Vec::new())
53    }
54
55    /// Borrow the complete source contents.
56    #[must_use]
57    pub fn as_bytes(&self) -> &[u8] {
58        &self.data
59    }
60
61    /// Consume the reader and return its source contents.
62    #[must_use]
63    pub fn into_bytes(self) -> Vec<u8> {
64        self.data
65    }
66}
67
68impl Default for AsyncMemReader {
69    fn default() -> Self {
70        Self::new()
71    }
72}
73
74impl AsyncReadAt for AsyncMemReader {
75    async fn read_at(&self, offset: u64, buf: &mut [u8]) -> io::Result<()> {
76        if buf.is_empty() {
77            return Ok(());
78        }
79
80        let offset = usize::try_from(offset).map_err(|_| {
81            io::Error::new(
82                io::ErrorKind::InvalidInput,
83                "read offset does not fit usize",
84            )
85        })?;
86        let end = offset.checked_add(buf.len()).ok_or_else(|| {
87            io::Error::new(io::ErrorKind::InvalidInput, "read range overflows usize")
88        })?;
89        if end > self.data.len() {
90            return Err(io::Error::new(
91                io::ErrorKind::UnexpectedEof,
92                format!(
93                    "positioned read needs {end} bytes but source contains {}",
94                    self.data.len()
95                ),
96            ));
97        }
98
99        buf.copy_from_slice(&self.data[offset..end]);
100        Ok(())
101    }
102}
103
104impl AsyncLength for AsyncMemReader {
105    async fn len(&self) -> io::Result<u64> {
106        u64::try_from(self.data.len()).map_err(|_| {
107            io::Error::new(io::ErrorKind::InvalidData, "source length does not fit u64")
108        })
109    }
110}
111
112#[cfg(test)]
113mod tests {
114    use super::*;
115    use futures::executor::block_on;
116
117    #[test]
118    fn positioned_reader_returns_exact_bytes_and_length() {
119        let reader = AsyncMemReader::from_bytes(vec![10, 20, 30, 40]);
120        let mut output = [0; 2];
121
122        block_on(async {
123            reader
124                .read_at(1, &mut output)
125                .await
126                .expect("read must succeed");
127            assert_eq!(output, [20, 30]);
128            assert_eq!(reader.len().await.expect("length must succeed"), 4);
129        });
130    }
131
132    #[test]
133    fn positioned_reader_rejects_short_reads() {
134        let reader = AsyncMemReader::from_bytes(vec![1, 2]);
135        let mut output = [0; 2];
136
137        let error = block_on(reader.read_at(1, &mut output)).expect_err("read must fail");
138        assert_eq!(error.kind(), io::ErrorKind::UnexpectedEof);
139    }
140
141    #[test]
142    fn zero_length_reads_do_not_require_source_bytes() {
143        let reader = AsyncMemReader::new();
144        let mut output = [];
145
146        block_on(reader.read_at(u64::MAX, &mut output)).expect("empty read must succeed");
147    }
148
149    #[test]
150    fn empty_source_reports_is_empty() {
151        // The provided `is_empty` must agree with `len` at the boundary that
152        // matters: a source with no bytes, and one with exactly one.
153        block_on(async {
154            let empty = AsyncMemReader::new();
155            assert_eq!(empty.len().await.expect("len"), 0);
156            assert!(empty.is_empty().await.expect("is_empty"));
157
158            let one = AsyncMemReader::from_bytes(vec![0u8; 1]);
159            assert_eq!(one.len().await.expect("len"), 1);
160            assert!(!one.is_empty().await.expect("is_empty"));
161        });
162    }
163}