Skip to main content

async_tiff/
reader.rs

1//! Abstractions for network reading.
2
3use std::fmt::Debug;
4use std::io::Read;
5use std::ops::Range;
6use std::sync::Arc;
7
8use async_trait::async_trait;
9use byteorder::{BigEndian, LittleEndian, ReadBytesExt};
10use bytes::buf::Reader;
11use bytes::{Buf, Bytes};
12use futures::TryFutureExt;
13
14use crate::error::AsyncTiffResult;
15
16/// The asynchronous interface used to read COG files
17///
18/// This was derived from the Parquet
19/// [`AsyncFileReader`](https://docs.rs/parquet/latest/parquet/arrow/async_reader/trait.AsyncFileReader.html)
20///
21/// Notes:
22///
23/// 1. [`ObjectReader`], available when the `object_store` crate feature
24///    is enabled, implements this interface for [`ObjectStore`].
25///
26/// 2. You can use [`TokioReader`] to implement [`AsyncFileReader`] for types that implement
27///    [`tokio::io::AsyncRead`] and [`tokio::io::AsyncSeek`], for example [`tokio::fs::File`].
28///
29/// [`ObjectStore`]: object_store::ObjectStore
30///
31/// [`tokio::fs::File`]: https://docs.rs/tokio/latest/tokio/fs/struct.File.html
32#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
33#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
34pub trait AsyncFileReader: Debug + Send + Sync + 'static {
35    /// Retrieve the bytes in `range` as part of a request for image data, not header metadata.
36    ///
37    /// This is also used as the default implementation of
38    /// [`MetadataFetch`][crate::metadata::MetadataFetch] if not overridden.
39    async fn get_bytes(&self, range: Range<u64>) -> AsyncTiffResult<Bytes>;
40
41    /// Retrieve multiple byte ranges as part of a request for image data, not header metadata. The
42    /// default implementation will call `get_bytes` sequentially
43    async fn get_byte_ranges(&self, ranges: Vec<Range<u64>>) -> AsyncTiffResult<Vec<Bytes>> {
44        let mut result = Vec::with_capacity(ranges.len());
45
46        for range in ranges.into_iter() {
47            let data = self.get_bytes(range).await?;
48            result.push(data);
49        }
50
51        Ok(result)
52    }
53}
54
55/// This allows Box<dyn AsyncFileReader + '_> to be used as an AsyncFileReader,
56#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
57#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
58impl AsyncFileReader for Box<dyn AsyncFileReader + '_> {
59    async fn get_bytes(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
60        self.as_ref().get_bytes(range).await
61    }
62
63    async fn get_byte_ranges(&self, ranges: Vec<Range<u64>>) -> AsyncTiffResult<Vec<Bytes>> {
64        self.as_ref().get_byte_ranges(ranges).await
65    }
66}
67
68/// This allows Arc<dyn AsyncFileReader + '_> to be used as an AsyncFileReader,
69#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
70#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
71impl AsyncFileReader for Arc<dyn AsyncFileReader + '_> {
72    async fn get_bytes(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
73        self.as_ref().get_bytes(range).await
74    }
75
76    async fn get_byte_ranges(&self, ranges: Vec<Range<u64>>) -> AsyncTiffResult<Vec<Bytes>> {
77        self.as_ref().get_byte_ranges(ranges).await
78    }
79}
80
81/// A wrapper for things that implement [AsyncRead] and [AsyncSeek] to also implement
82/// [AsyncFileReader].
83///
84/// This wrapper is needed because `AsyncRead` and `AsyncSeek` require mutable access to seek and
85/// read data, while the `AsyncFileReader` trait requires immutable access to read data.
86///
87/// This wrapper stores the inner reader in a `Mutex`.
88///
89/// [AsyncRead]: tokio::io::AsyncRead
90/// [AsyncSeek]: tokio::io::AsyncSeek
91#[cfg(feature = "tokio")]
92#[derive(Debug)]
93pub struct TokioReader<T: tokio::io::AsyncRead + tokio::io::AsyncSeek + Unpin + Send + Debug>(
94    tokio::sync::Mutex<T>,
95);
96
97#[cfg(feature = "tokio")]
98impl<T: tokio::io::AsyncRead + tokio::io::AsyncSeek + Unpin + Send + Debug> TokioReader<T> {
99    /// Create a new TokioReader from a reader.
100    pub fn new(inner: T) -> Self {
101        Self(tokio::sync::Mutex::new(inner))
102    }
103
104    async fn make_range_request(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
105        use std::io::SeekFrom;
106
107        use tokio::io::{AsyncReadExt, AsyncSeekExt};
108
109        use crate::error::AsyncTiffError;
110
111        let mut file = self.0.lock().await;
112
113        file.seek(SeekFrom::Start(range.start)).await?;
114
115        let to_read = range.end - range.start;
116        let mut buffer = Vec::with_capacity(to_read as usize);
117        let read = file.read(&mut buffer).await? as u64;
118        if read != to_read {
119            return Err(AsyncTiffError::EndOfFile(to_read, read));
120        }
121
122        Ok(buffer.into())
123    }
124}
125
126#[cfg(feature = "tokio")]
127#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
128#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
129impl<T: tokio::io::AsyncRead + tokio::io::AsyncSeek + Unpin + Send + Debug + 'static>
130    AsyncFileReader for TokioReader<T>
131{
132    async fn get_bytes(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
133        self.make_range_request(range).await
134    }
135}
136
137/// An AsyncFileReader that reads from an [`ObjectStore`][object_store::ObjectStore] instance.
138#[cfg(feature = "object_store")]
139#[derive(Clone, Debug)]
140pub struct ObjectReader {
141    store: Arc<dyn object_store::ObjectStore>,
142    path: object_store::path::Path,
143}
144
145#[cfg(feature = "object_store")]
146impl ObjectReader {
147    /// Creates a new [`ObjectReader`] for the provided [`ObjectStore`][object_store::ObjectStore]
148    /// and path.
149    pub fn new(store: Arc<dyn object_store::ObjectStore>, path: object_store::path::Path) -> Self {
150        Self { store, path }
151    }
152
153    async fn make_range_request(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
154        use object_store::ObjectStoreExt;
155
156        let range = range.start as _..range.end as _;
157        self.store
158            .get_range(&self.path, range)
159            .map_err(|e| e.into())
160            .await
161    }
162}
163
164#[cfg(feature = "object_store")]
165#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
166#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
167impl AsyncFileReader for ObjectReader {
168    async fn get_bytes(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
169        self.make_range_request(range).await
170    }
171
172    async fn get_byte_ranges(&self, ranges: Vec<Range<u64>>) -> AsyncTiffResult<Vec<Bytes>>
173    where
174        Self: Send,
175    {
176        let ranges = ranges
177            .into_iter()
178            .map(|r| r.start as _..r.end as _)
179            .collect::<Vec<_>>();
180        self.store
181            .get_ranges(&self.path, &ranges)
182            .await
183            .map_err(|e| e.into())
184    }
185}
186
187/// An AsyncFileReader that reads from a URL using reqwest.
188#[cfg(feature = "reqwest")]
189#[derive(Debug, Clone)]
190pub struct ReqwestReader {
191    client: reqwest::Client,
192    url: reqwest::Url,
193}
194
195#[cfg(feature = "reqwest")]
196impl ReqwestReader {
197    /// Construct a new ReqwestReader from a reqwest client and URL.
198    pub fn new(client: reqwest::Client, url: reqwest::Url) -> Self {
199        Self { client, url }
200    }
201
202    async fn make_range_request(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
203        let url = self.url.clone();
204        let client = self.client.clone();
205        // HTTP range is inclusive, so we need to subtract 1 from the end
206        let range = format!("bytes={}-{}", range.start, range.end - 1);
207        let response = client
208            .get(url)
209            .header("Range", range)
210            .send()
211            .await?
212            .error_for_status()?;
213        let bytes = response.bytes().await?;
214        Ok(bytes)
215    }
216}
217
218#[cfg(feature = "reqwest")]
219#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
220#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
221impl AsyncFileReader for ReqwestReader {
222    async fn get_bytes(&self, range: Range<u64>) -> AsyncTiffResult<Bytes> {
223        self.make_range_request(range).await
224    }
225}
226
227/// Endianness
228#[derive(Debug, Clone, Copy, PartialEq)]
229pub enum Endianness {
230    /// Little Endian
231    LittleEndian,
232    /// Big Endian
233    BigEndian,
234}
235
236impl Endianness {
237    /// Check if the endianness matches the native endianness of the host system.
238    ///
239    /// ```
240    /// use async_tiff::reader::Endianness;
241    ///
242    /// if cfg!(target_endian = "little") {
243    ///     assert!(Endianness::LittleEndian.is_native());
244    ///     assert!(!Endianness::BigEndian.is_native());
245    /// } else {
246    ///     assert!(Endianness::BigEndian.is_native());
247    ///     assert!(!Endianness::LittleEndian.is_native());
248    /// }
249    /// ```
250    pub fn is_native(&self) -> bool {
251        let native_endianness = if cfg!(target_endian = "little") {
252            Endianness::LittleEndian
253        } else {
254            Endianness::BigEndian
255        };
256
257        *self == native_endianness
258    }
259}
260
261pub(crate) struct EndianAwareReader {
262    reader: Reader<Bytes>,
263    endianness: Endianness,
264}
265
266impl EndianAwareReader {
267    pub(crate) fn new(bytes: Bytes, endianness: Endianness) -> Self {
268        Self {
269            reader: bytes.reader(),
270            endianness,
271        }
272    }
273
274    /// Read a u8 from the cursor, advancing the internal state by 1 byte.
275    pub(crate) fn read_u8(&mut self) -> AsyncTiffResult<u8> {
276        Ok(self.reader.read_u8()?)
277    }
278
279    /// Read a i8 from the cursor, advancing the internal state by 1 byte.
280    pub(crate) fn read_i8(&mut self) -> AsyncTiffResult<i8> {
281        Ok(self.reader.read_i8()?)
282    }
283
284    pub(crate) fn read_u16(&mut self) -> AsyncTiffResult<u16> {
285        match self.endianness {
286            Endianness::LittleEndian => Ok(self.reader.read_u16::<LittleEndian>()?),
287            Endianness::BigEndian => Ok(self.reader.read_u16::<BigEndian>()?),
288        }
289    }
290
291    pub(crate) fn read_i16(&mut self) -> AsyncTiffResult<i16> {
292        match self.endianness {
293            Endianness::LittleEndian => Ok(self.reader.read_i16::<LittleEndian>()?),
294            Endianness::BigEndian => Ok(self.reader.read_i16::<BigEndian>()?),
295        }
296    }
297
298    pub(crate) fn read_u32(&mut self) -> AsyncTiffResult<u32> {
299        match self.endianness {
300            Endianness::LittleEndian => Ok(self.reader.read_u32::<LittleEndian>()?),
301            Endianness::BigEndian => Ok(self.reader.read_u32::<BigEndian>()?),
302        }
303    }
304
305    pub(crate) fn read_i32(&mut self) -> AsyncTiffResult<i32> {
306        match self.endianness {
307            Endianness::LittleEndian => Ok(self.reader.read_i32::<LittleEndian>()?),
308            Endianness::BigEndian => Ok(self.reader.read_i32::<BigEndian>()?),
309        }
310    }
311
312    pub(crate) fn read_u64(&mut self) -> AsyncTiffResult<u64> {
313        match self.endianness {
314            Endianness::LittleEndian => Ok(self.reader.read_u64::<LittleEndian>()?),
315            Endianness::BigEndian => Ok(self.reader.read_u64::<BigEndian>()?),
316        }
317    }
318
319    pub(crate) fn read_i64(&mut self) -> AsyncTiffResult<i64> {
320        match self.endianness {
321            Endianness::LittleEndian => Ok(self.reader.read_i64::<LittleEndian>()?),
322            Endianness::BigEndian => Ok(self.reader.read_i64::<BigEndian>()?),
323        }
324    }
325
326    pub(crate) fn read_f32(&mut self) -> AsyncTiffResult<f32> {
327        match self.endianness {
328            Endianness::LittleEndian => Ok(self.reader.read_f32::<LittleEndian>()?),
329            Endianness::BigEndian => Ok(self.reader.read_f32::<BigEndian>()?),
330        }
331    }
332
333    pub(crate) fn read_f64(&mut self) -> AsyncTiffResult<f64> {
334        match self.endianness {
335            Endianness::LittleEndian => Ok(self.reader.read_f64::<LittleEndian>()?),
336            Endianness::BigEndian => Ok(self.reader.read_f64::<BigEndian>()?),
337        }
338    }
339
340    #[allow(dead_code)]
341    pub(crate) fn into_inner(self) -> (Reader<Bytes>, Endianness) {
342        (self.reader, self.endianness)
343    }
344}
345
346impl AsRef<[u8]> for EndianAwareReader {
347    fn as_ref(&self) -> &[u8] {
348        self.reader.get_ref().as_ref()
349    }
350}
351
352impl Read for EndianAwareReader {
353    fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
354        self.reader.read(buf)
355    }
356}