| 635 | } |
| 636 | |
| 637 | RemoteQueryExecutor::ReadResult RemoteQueryExecutor::processPacket(Packet packet) |
| 638 | { |
| 639 | switch (packet.type) |
| 640 | { |
| 641 | case Protocol::Server::MergeTreeReadTaskRequest: |
| 642 | chassert(packet.request.has_value()); |
| 643 | processMergeTreeReadTaskRequest(packet.request.value()); |
| 644 | return ReadResult(ReadResult::Type::ParallelReplicasToken); |
| 645 | |
| 646 | case Protocol::Server::MergeTreeAllRangesAnnouncement: |
| 647 | chassert(packet.announcement.has_value()); |
| 648 | processMergeTreeInitialReadAnnouncement(packet.announcement.value()); |
| 649 | return ReadResult(ReadResult::Type::ParallelReplicasToken); |
| 650 | |
| 651 | case Protocol::Server::ReadTaskRequest: |
| 652 | processReadTaskRequest(); |
| 653 | break; |
| 654 | case Protocol::Server::PartUUIDs: |
| 655 | LOG_WARNING( |
| 656 | log, |
| 657 | "The remote server has sent no longer supported packet (Server::PartUUIDs). allow_experimental_query_deduplication feature " |
| 658 | "has been deprecated. Consider upgrading the remote server ({})", |
| 659 | connections->dumpAddresses()); |
| 660 | break; |
| 661 | case Protocol::Server::Data: |
| 662 | /// Note: `packet.block.rows() > 0` means it's a header block. |
| 663 | /// We can actually return it, and the first call to RemoteQueryExecutor::read |
| 664 | /// will return earlier. We should consider doing it. |
| 665 | if (!packet.block.empty() && (packet.block.rows() > 0)) |
| 666 | return ReadResult(adaptBlockStructure(packet.block, *header)); |
| 667 | break; /// If the block is empty - we will receive other packets before EndOfStream. |
| 668 | |
| 669 | case Protocol::Server::Exception: |
| 670 | got_exception_from_replica = true; |
| 671 | packet.exception->rethrow(); |
| 672 | break; |
| 673 | |
| 674 | case Protocol::Server::EndOfStream: |
| 675 | if (!connections->hasActiveConnections()) |
| 676 | { |
| 677 | finished = true; |
| 678 | /// TODO: Replace with Type::Finished |
| 679 | return ReadResult(Block{}); |
| 680 | } |
| 681 | break; |
| 682 | |
| 683 | case Protocol::Server::Progress: |
| 684 | /** We use the progress from a remote server. |
| 685 | * We also include in ProcessList, |
| 686 | * and we use it to check |
| 687 | * constraints (for example, the minimum speed of query execution) |
| 688 | * and quotas (for example, the number of lines to read). |
| 689 | */ |
| 690 | if (progress_callback) |
| 691 | progress_callback(packet.progress); |
| 692 | break; |
| 693 | |
| 694 | case Protocol::Server::ProfileInfo: |
nothing calls this directly
no test coverage detected