2020pub mod local;
2121
2222use std:: collections:: HashMap ;
23+ use std:: error:: Error ;
2324use std:: fmt:: { self , Debug } ;
2425use std:: io:: Read ;
2526use std:: pin:: Pin ;
@@ -31,7 +32,7 @@ use futures::{AsyncRead, Stream, StreamExt};
3132
3233use local:: LocalFileSystem ;
3334
34- use crate :: error:: { DataFusionError , Result } ;
35+ use crate :: error:: { DataFusionError , Result as DataFusionResult } ;
3536
3637/// Object Reader for one file in an object store.
3738///
@@ -40,18 +41,23 @@ use crate::error::{DataFusionError, Result};
4041#[ async_trait]
4142pub trait ObjectReader : Send + Sync {
4243 /// Get reader for a part [start, start + length] in the file asynchronously
43- async fn chunk_reader ( & self , start : u64 , length : usize )
44- -> Result < Box < dyn AsyncRead > > ;
44+ async fn chunk_reader (
45+ & self ,
46+ start : u64 ,
47+ length : usize ,
48+ ) -> Result < Box < dyn AsyncRead > , Box < dyn Error + Send + Sync > > ;
4549
4650 /// Get reader for a part [start, start + length] in the file
4751 fn sync_chunk_reader (
4852 & self ,
4953 start : u64 ,
5054 length : usize ,
51- ) -> Result < Box < dyn Read + Send + Sync > > ;
55+ ) -> Result < Box < dyn Read + Send + Sync > , Box < dyn Error + Send + Sync > > ;
5256
5357 /// Get reader for the entire file
54- fn sync_reader ( & self ) -> Result < Box < dyn Read + Send + Sync > > {
58+ fn sync_reader (
59+ & self ,
60+ ) -> Result < Box < dyn Read + Send + Sync > , Box < dyn Error + Send + Sync > > {
5561 self . sync_chunk_reader ( 0 , self . length ( ) as usize )
5662 }
5763
@@ -114,29 +120,38 @@ impl std::fmt::Display for FileMeta {
114120
115121/// Stream of files listed from object store
116122pub type FileMetaStream =
117- Pin < Box < dyn Stream < Item = Result < FileMeta > > + Send + Sync + ' static > > ;
123+ Pin < Box < dyn Stream < Item = DataFusionResult < FileMeta > > + Send + Sync + ' static > > ;
118124
119125/// Stream of list entries obtained from object store
120126pub type ListEntryStream =
121- Pin < Box < dyn Stream < Item = Result < ListEntry > > + Send + Sync + ' static > > ;
127+ Pin < Box < dyn Stream < Item = DataFusionResult < ListEntry > > + Send + Sync + ' static > > ;
122128
123129/// Stream readers opened on a given object store
124- pub type ObjectReaderStream =
125- Pin < Box < dyn Stream < Item = Result < Arc < dyn ObjectReader > > > + Send + Sync + ' static > > ;
130+ pub type ObjectReaderStream = Pin <
131+ Box <
132+ dyn Stream < Item = Result < Arc < dyn ObjectReader > , Box < dyn Error + Send + Sync > > >
133+ + Send
134+ + Sync
135+ + ' static ,
136+ > ,
137+ > ;
126138
127139/// A ObjectStore abstracts access to an underlying file/object storage.
128140/// It maps strings (e.g. URLs, filesystem paths, etc) to sources of bytes
129141#[ async_trait]
130142pub trait ObjectStore : Sync + Send + Debug {
131143 /// Returns all the files in path `prefix`
132- async fn list_file ( & self , prefix : & str ) -> Result < FileMetaStream > ;
144+ async fn list_file (
145+ & self ,
146+ prefix : & str ,
147+ ) -> Result < FileMetaStream , Box < dyn Error + Send + Sync > > ;
133148
134149 /// Calls `list_file` with a suffix filter
135150 async fn list_file_with_suffix (
136151 & self ,
137152 prefix : & str ,
138153 suffix : & str ,
139- ) -> Result < FileMetaStream > {
154+ ) -> Result < FileMetaStream , Box < dyn Error + Send + Sync > > {
140155 let file_stream = self . list_file ( prefix) . await ?;
141156 let suffix = suffix. to_owned ( ) ;
142157 Ok ( Box :: pin ( file_stream. filter ( move |fr| {
@@ -154,10 +169,13 @@ pub trait ObjectStore: Sync + Send + Debug {
154169 & self ,
155170 prefix : & str ,
156171 delimiter : Option < String > ,
157- ) -> Result < ListEntryStream > ;
172+ ) -> Result < ListEntryStream , Box < dyn Error + Send + Sync > > ;
158173
159174 /// Get object reader for one file
160- fn file_reader ( & self , file : SizedFile ) -> Result < Arc < dyn ObjectReader > > ;
175+ fn file_reader (
176+ & self ,
177+ file : SizedFile ,
178+ ) -> Result < Arc < dyn ObjectReader > , Box < dyn Error + Send + Sync > > ;
161179}
162180
163181static LOCAL_SCHEME : & str = "file" ;
@@ -223,7 +241,7 @@ impl ObjectStoreRegistry {
223241 pub fn get_by_uri < ' a > (
224242 & self ,
225243 uri : & ' a str ,
226- ) -> Result < ( Arc < dyn ObjectStore > , & ' a str ) > {
244+ ) -> Result < ( Arc < dyn ObjectStore > , & ' a str ) , Box < dyn Error + Send + Sync > > {
227245 if let Some ( ( scheme, path) ) = uri. split_once ( "://" ) {
228246 let stores = self . object_stores . read ( ) . unwrap ( ) ;
229247 let store = stores
0 commit comments