Try to read the next [`RecordBatch`] from the provided [`Buffer`] [`Buffer::advance`] will be called on `buffer` for any consumed bytes. The push-based interface facilitates integration with sources that yield arbitrarily delimited bytes ranges, such as a chunked byte stream received from object storage ``` # use arrow_array::RecordBatch; # use arrow_buffer::Buffer; # use arrow_ipc::reader::Str
(&mut self, buffer: &mut Buffer)
| 157 | /// } |
| 158 | /// ``` |
| 159 | pub fn decode(&mut self, buffer: &mut Buffer) -> Result<Option<RecordBatch>, ArrowError> { |
| 160 | while !buffer.is_empty() { |
| 161 | match &mut self.state { |
| 162 | DecoderState::Header { |
| 163 | buf, |
| 164 | read, |
| 165 | continuation, |
| 166 | } => { |
| 167 | let offset_buf = &mut buf[*read as usize..]; |
| 168 | let to_read = buffer.len().min(offset_buf.len()); |
| 169 | offset_buf[..to_read].copy_from_slice(&buffer[..to_read]); |
| 170 | *read += to_read as u8; |
| 171 | buffer.advance(to_read); |
| 172 | if *read == 4 { |
| 173 | if !*continuation && buf == &CONTINUATION_MARKER { |
| 174 | *continuation = true; |
| 175 | *read = 0; |
| 176 | continue; |
| 177 | } |
| 178 | let size = u32::from_le_bytes(*buf); |
| 179 | |
| 180 | if size == 0 { |
| 181 | self.state = DecoderState::Finished; |
| 182 | continue; |
| 183 | } |
| 184 | self.state = DecoderState::Message { size }; |
| 185 | } |
| 186 | } |
| 187 | DecoderState::Message { size } => { |
| 188 | let len = *size as usize; |
| 189 | if self.buf.is_empty() && buffer.len() > len { |
| 190 | let message = MessageBuffer::try_new(buffer.slice_with_length(0, len))?; |
| 191 | self.state = DecoderState::Body { message }; |
| 192 | buffer.advance(len); |
| 193 | continue; |
| 194 | } |
| 195 | |
| 196 | let to_read = buffer.len().min(len - self.buf.len()); |
| 197 | self.buf.extend_from_slice(&buffer[..to_read]); |
| 198 | buffer.advance(to_read); |
| 199 | if self.buf.len() == len { |
| 200 | let message = MessageBuffer::try_new(std::mem::take(&mut self.buf).into())?; |
| 201 | self.state = DecoderState::Body { message }; |
| 202 | } |
| 203 | } |
| 204 | DecoderState::Body { message } => { |
| 205 | let message = message.as_ref(); |
| 206 | let body_length = message.bodyLength() as usize; |
| 207 | |
| 208 | let body = if self.buf.is_empty() && buffer.len() >= body_length { |
| 209 | let body = buffer.slice_with_length(0, body_length); |
| 210 | buffer.advance(body_length); |
| 211 | body |
| 212 | } else { |
| 213 | let to_read = buffer.len().min(body_length - self.buf.len()); |
| 214 | self.buf.extend_from_slice(&buffer[..to_read]); |
| 215 | buffer.advance(to_read); |
| 216 |