(&mut self, cx: &mut Context<'_>)
| 252 | |
| 253 | #[forbid(clippy::question_mark_used)] |
| 254 | fn read_next(&mut self, cx: &mut Context<'_>) -> Poll<Option<Result<RecordBatch, AvroError>>> { |
| 255 | loop { |
| 256 | match mem::replace(&mut self.reader_state, ReaderState::InvalidState) { |
| 257 | ReaderState::Idle { reader } => { |
| 258 | let range = self.range.clone(); |
| 259 | if range.start >= range.end { |
| 260 | return self.finish_with_error(AvroError::InvalidArgument(format!( |
| 261 | "Invalid range specified for Avro file: start {} >= end {}, file_size: {}", |
| 262 | range.start, range.end, self.file_size |
| 263 | ))); |
| 264 | } |
| 265 | |
| 266 | let future = Self::fetch_bytes(reader, range).boxed(); |
| 267 | self.reader_state = ReaderState::FetchingData { |
| 268 | future, |
| 269 | next_behaviour: FetchNextBehaviour::ReadSyncMarker, |
| 270 | }; |
| 271 | } |
| 272 | ReaderState::FetchingData { |
| 273 | mut future, |
| 274 | next_behaviour, |
| 275 | } => { |
| 276 | let (reader, data_chunk) = match future.poll_unpin(cx) { |
| 277 | Poll::Ready(Ok(data)) => data, |
| 278 | Poll::Ready(Err(e)) => return self.finish_with_error(e), |
| 279 | Poll::Pending => { |
| 280 | self.reader_state = ReaderState::FetchingData { |
| 281 | future, |
| 282 | next_behaviour, |
| 283 | }; |
| 284 | return Poll::Pending; |
| 285 | } |
| 286 | }; |
| 287 | |
| 288 | match next_behaviour { |
| 289 | FetchNextBehaviour::ReadSyncMarker => { |
| 290 | let sync_marker_pos = data_chunk |
| 291 | .windows(16) |
| 292 | .position(|slice| slice == self.sync_marker); |
| 293 | let block_start = match sync_marker_pos { |
| 294 | Some(pos) => pos + 16, // Move past the sync marker |
| 295 | None => { |
| 296 | // Sync marker not found, valid if we arbitrarily split the file at its end. |
| 297 | self.reader_state = ReaderState::Finished; |
| 298 | return Poll::Ready(None); |
| 299 | } |
| 300 | }; |
| 301 | |
| 302 | self.reader_state = ReaderState::DecodingBlock { |
| 303 | reader, |
| 304 | data: data_chunk.slice(block_start..), |
| 305 | }; |
| 306 | } |
| 307 | FetchNextBehaviour::DecodeVLQHeader => { |
| 308 | let mut data = data_chunk; |
| 309 | |
| 310 | // Feed bytes one at a time until we reach Data state (VLQ header complete) |
| 311 | while !matches!(self.block_decoder.state(), BlockDecoderState::Data) { |
no test coverage detected