Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions arrow-avro/src/errors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,13 @@ impl From<ArrowError> for AvroError {
}
}

#[cfg(feature = "object_store")]
impl From<object_store::Error> for AvroError {
fn from(e: object_store::Error) -> AvroError {
AvroError::External(Box::new(e))
}
}

impl From<AvroError> for io::Error {
fn from(e: AvroError) -> Self {
io::Error::other(e)
Expand Down
9 changes: 1 addition & 8 deletions arrow-avro/src/reader/async_reader/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -154,14 +154,7 @@ where
break;
}

let current_data = reader
.get_bytes(range_to_fetch.clone())
.await
.map_err(|err| {
AvroError::General(format!(
"Error fetching Avro header from file reader: {err}"
))
})?;
let current_data = reader.get_bytes(range_to_fetch.clone()).await?;
if current_data.is_empty() {
return Err(AvroError::EOF(
"Unexpected EOF while fetching header data".into(),
Expand Down
55 changes: 49 additions & 6 deletions arrow-avro/src/reader/async_reader/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ use futures::{FutureExt, TryFutureExt};
use object_store::ObjectStore;
use object_store::ObjectStoreExt;
use object_store::path::Path;
use std::error::Error;
use std::ops::Range;
use std::sync::Arc;
use tokio::runtime::Handle;
Expand Down Expand Up @@ -66,7 +65,7 @@ impl AvroObjectReader {
+ Send
+ 'static,
O: Send + 'static,
E: Error + Send + 'static,
E: Into<AvroError> + Send + 'static,
{
match &self.runtime {
Some(handle) => {
Expand All @@ -79,13 +78,11 @@ impl AvroObjectReader {
Err(e) => Err(AvroError::External(Box::new(e))),
Ok(p) => std::panic::resume_unwind(p),
},
|res| res.map_err(|e| AvroError::General(e.to_string())),
|res| res.map_err(Into::into),
)
.boxed()
}
None => f(&self.store, &self.path)
.map_err(|e| AvroError::General(e.to_string()))
.boxed(),
None => f(&self.store, &self.path).map_err(Into::into).boxed(),
}
}
}
Expand All @@ -105,3 +102,49 @@ impl AsyncFileReader for AvroObjectReader {
self.spawn(|store, path| async move { store.get_ranges(path, &ranges).await }.boxed())
}
}

#[cfg(test)]
mod tests {
use super::*;
use object_store::memory::InMemory;

fn find_object_store_error(err: &AvroError) -> Option<&object_store::Error> {
let mut source: Option<&(dyn std::error::Error + 'static)> = Some(err);
while let Some(e) = source {
if let Some(os) = e.downcast_ref::<object_store::Error>() {
return Some(os);
}
source = e.source();
}
None
}

#[tokio::test]
async fn test_get_bytes_preserves_object_store_error_source() {
let store = Arc::new(InMemory::new());
let mut reader = AvroObjectReader::new(store, Path::from("missing.avro"));
let err = reader.get_bytes(0..10).await.unwrap_err();
assert!(
matches!(
find_object_store_error(&err),
Some(object_store::Error::NotFound { .. })
),
"expected NotFound in source chain, got: {err}"
);
}

#[tokio::test]
async fn test_get_bytes_on_runtime_preserves_object_store_error_source() {
let store = Arc::new(InMemory::new());
let mut reader = AvroObjectReader::new(store, Path::from("missing.avro"))
.with_runtime(Handle::current());
let err = reader.get_bytes(0..10).await.unwrap_err();
assert!(
matches!(
find_object_store_error(&err),
Some(object_store::Error::NotFound { .. })
),
"expected NotFound in source chain, got: {err}"
);
}
}
Loading