| 456 | } |
| 457 | |
| 458 | Future<std::shared_ptr<Message>> ReadMessageAsync(int64_t offset, int32_t metadata_length, |
| 459 | int64_t body_length, |
| 460 | io::RandomAccessFile* file, |
| 461 | const io::IOContext& context) { |
| 462 | struct State { |
| 463 | std::unique_ptr<Message> result; |
| 464 | std::shared_ptr<MessageDecoderListener> listener; |
| 465 | std::shared_ptr<MessageDecoder> decoder; |
| 466 | }; |
| 467 | auto state = std::make_shared<State>(); |
| 468 | state->listener = std::make_shared<AssignMessageDecoderListener>(&state->result); |
| 469 | state->decoder = std::make_shared<MessageDecoder>(state->listener); |
| 470 | |
| 471 | if (metadata_length < state->decoder->next_required_size()) { |
| 472 | return Status::Invalid("metadata_length should be at least ", |
| 473 | state->decoder->next_required_size()); |
| 474 | } |
| 475 | return file->ReadAsync(context, offset, metadata_length + body_length) |
| 476 | .Then([=](std::shared_ptr<Buffer> metadata) -> Result<std::shared_ptr<Message>> { |
| 477 | if (metadata->size() < metadata_length) { |
| 478 | return Status::Invalid("Expected to read ", metadata_length, |
| 479 | " metadata bytes but got ", metadata->size()); |
| 480 | } |
| 481 | ARROW_RETURN_NOT_OK( |
| 482 | state->decoder->Consume(SliceBuffer(metadata, 0, metadata_length))); |
| 483 | switch (state->decoder->state()) { |
| 484 | case MessageDecoder::State::INITIAL: |
| 485 | return std::move(state->result); |
| 486 | case MessageDecoder::State::METADATA_LENGTH: |
| 487 | return Status::Invalid("metadata length is missing. File offset: ", offset, |
| 488 | ", metadata length: ", metadata_length); |
| 489 | case MessageDecoder::State::METADATA: |
| 490 | return Status::Invalid("flatbuffer size ", |
| 491 | state->decoder->next_required_size(), |
| 492 | " invalid. File offset: ", offset, |
| 493 | ", metadata length: ", metadata_length); |
| 494 | case MessageDecoder::State::BODY: { |
| 495 | auto body = SliceBuffer(metadata, metadata_length, body_length); |
| 496 | if (body->size() < state->decoder->next_required_size()) { |
| 497 | return Status::IOError("Expected to be able to read ", |
| 498 | state->decoder->next_required_size(), |
| 499 | " bytes for message body, got ", body->size()); |
| 500 | } |
| 501 | RETURN_NOT_OK(state->decoder->Consume(body)); |
| 502 | return std::move(state->result); |
| 503 | } |
| 504 | case MessageDecoder::State::EOS: |
| 505 | return Status::Invalid("Unexpected empty message in IPC file format"); |
| 506 | default: |
| 507 | return Status::Invalid("Unexpected state: ", state->decoder->state()); |
| 508 | } |
| 509 | }); |
| 510 | } |
| 511 | |
| 512 | Status AlignStream(io::InputStream* stream, int32_t alignment) { |
| 513 | ARROW_ASSIGN_OR_RAISE(int64_t position, stream->Tell()); |
no test coverage detected