| 136 | } |
| 137 | |
| 138 | async fn read_header<R>( |
| 139 | reader: &mut R, |
| 140 | file_size: u64, |
| 141 | header_size_hint: Option<u64>, |
| 142 | ) -> Result<(Header, u64), AvroError> |
| 143 | where |
| 144 | R: AsyncFileReader, |
| 145 | { |
| 146 | let mut decoder = HeaderDecoder::default(); |
| 147 | let mut position = 0; |
| 148 | loop { |
| 149 | let range_to_fetch = position |
| 150 | ..(position + header_size_hint.unwrap_or(DEFAULT_HEADER_SIZE_HINT)).min(file_size); |
| 151 | |
| 152 | // Maybe EOF after the header, no actual data |
| 153 | if range_to_fetch.is_empty() { |
| 154 | break; |
| 155 | } |
| 156 | |
| 157 | let current_data = reader |
| 158 | .get_bytes(range_to_fetch.clone()) |
| 159 | .await |
| 160 | .map_err(|err| { |
| 161 | AvroError::General(format!( |
| 162 | "Error fetching Avro header from file reader: {err}" |
| 163 | )) |
| 164 | })?; |
| 165 | if current_data.is_empty() { |
| 166 | return Err(AvroError::EOF( |
| 167 | "Unexpected EOF while fetching header data".into(), |
| 168 | )); |
| 169 | } |
| 170 | |
| 171 | let read = current_data.len(); |
| 172 | let decoded = decoder.decode(¤t_data)?; |
| 173 | if decoded != read { |
| 174 | position += decoded as u64; |
| 175 | break; |
| 176 | } |
| 177 | position += read as u64; |
| 178 | } |
| 179 | |
| 180 | decoder |
| 181 | .flush() |
| 182 | .map(|header| (header, position)) |
| 183 | .ok_or_else(|| AvroError::EOF("Unexpected EOF while reading Avro header".into())) |
| 184 | } |
| 185 | |
| 186 | impl<R: AsyncFileReader> ReaderBuilder<R> { |
| 187 | /// Build the asynchronous Avro reader with the provided parameters. |