| 566 | } |
| 567 | |
| 568 | RemoteQueryExecutor::ReadResult RemoteQueryExecutor::readAsync() |
| 569 | { |
| 570 | #if defined(OS_LINUX) |
| 571 | if (!read_context) |
| 572 | { |
| 573 | LockAndBlocker lock(was_cancelled_mutex); |
| 574 | if (was_cancelled) |
| 575 | return ReadResult(Block()); |
| 576 | |
| 577 | read_context = std::make_unique<ReadContext>( |
| 578 | *this, |
| 579 | /*suspend_when_query_sent*/ false, |
| 580 | read_packet_type_separately); |
| 581 | } |
| 582 | |
| 583 | while (true) |
| 584 | { |
| 585 | LockAndBlocker lock(was_cancelled_mutex); |
| 586 | if (was_cancelled) |
| 587 | return ReadResult(Block()); |
| 588 | |
| 589 | if (packet_in_progress) |
| 590 | { |
| 591 | chassert(read_context->readPacketTypeSeparately()); |
| 592 | chassert(read_context->hasReadTillPacketType()); |
| 593 | |
| 594 | /// packet type is handled already, read and parse packet itself |
| 595 | if (!read_context->hasReadPacket() && !read_context->read()) |
| 596 | return ReadResult(read_context->getFileDescriptor()); |
| 597 | |
| 598 | packet_in_progress = false; |
| 599 | auto read_result = processPacket(read_context->getPacket()); |
| 600 | if (read_result.getType() == ReadResult::Type::Data || read_result.getType() == ReadResult::Type::ParallelReplicasToken) |
| 601 | return read_result; |
| 602 | } |
| 603 | |
| 604 | read_context->resume(); |
| 605 | |
| 606 | if (isReplicaUnavailable() || needToSkipUnavailableShard()) |
| 607 | { |
| 608 | /// We need to tell the coordinator not to wait for this replica. |
| 609 | /// But at this point it may lead to an incomplete result set, because |
| 610 | /// this replica committed to read some part of there data and then died. |
| 611 | if (extension && extension->parallel_reading_coordinator) |
| 612 | { |
| 613 | chassert(extension->parallel_reading_coordinator); |
| 614 | extension->parallel_reading_coordinator->markReplicaAsUnavailable(extension->replica_info->number_of_current_replica); |
| 615 | } |
| 616 | |
| 617 | return ReadResult(Block()); |
| 618 | } |
| 619 | |
| 620 | /// Check if packet is not ready yet. |
| 621 | if (read_context->isInProgress()) |
| 622 | return ReadResult(read_context->getFileDescriptor()); |
| 623 | |
| 624 | /// if reading separately packet header and body enabled, try to read packet itself this time |
| 625 | if (read_context->readPacketTypeSeparately() && !read_context->hasReadPacket() && !read_context->read()) |
no test coverage detected