MCPcopy Create free account
hub / github.com/apache/arrow / ReadMessageAsync

Function ReadMessageAsync

cpp/src/arrow/ipc/message.cc:458–510  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

456}
457
458Future<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
512Status AlignStream(io::InputStream* stream, int32_t alignment) {
513 ARROW_ASSIGN_OR_RAISE(int64_t position, stream->Tell());

Callers 1

Calls 9

SliceBufferFunction · 0.85
IOErrorFunction · 0.85
ThenMethod · 0.80
InvalidFunction · 0.50
next_required_sizeMethod · 0.45
ReadAsyncMethod · 0.45
sizeMethod · 0.45
ConsumeMethod · 0.45
stateMethod · 0.45

Tested by

no test coverage detected