| 689 | } |
| 690 | |
| 691 | std::unique_ptr<ral::cache::CacheData> CacheMachine::pullAnyCacheData(const std::vector<std::string> & messages) { |
| 692 | |
| 693 | if (messages.size() == 0){ |
| 694 | return nullptr; |
| 695 | } |
| 696 | |
| 697 | CodeTimer cacheEventTimer; |
| 698 | cacheEventTimer.start(); |
| 699 | |
| 700 | std::unique_ptr<message> message_data = waitingCache->get_or_wait_any(messages); |
| 701 | std::string message_id = message_data->get_message_id(); |
| 702 | |
| 703 | size_t num_rows = message_data->get_data().num_rows(); |
| 704 | size_t num_bytes = message_data->get_data().sizeInBytes(); |
| 705 | int dataType = static_cast<int>(message_data->get_data().get_type()); |
| 706 | std::unique_ptr<ral::cache::CacheData> output = message_data->release_data(); |
| 707 | |
| 708 | cacheEventTimer.stop(); |
| 709 | if(cache_events_logger) { |
| 710 | cache_events_logger->info("{ral_id}|{query_id}|{message_id}|{cache_id}|{num_rows}|{num_bytes}|{event_type}|{timestamp_begin}|{timestamp_end}|{description}", |
| 711 | "ral_id"_a=(ctx ? ctx->getNodeIndex(ral::communication::CommunicationData::getInstance().getSelfNode()) : -1), |
| 712 | "query_id"_a=(ctx ? ctx->getContextToken() : -1), |
| 713 | "message_id"_a=message_id, |
| 714 | "cache_id"_a=cache_id, |
| 715 | "num_rows"_a=num_rows, |
| 716 | "num_bytes"_a=num_bytes, |
| 717 | "event_type"_a="pullAnyCacheData", |
| 718 | "timestamp_begin"_a=cacheEventTimer.start_time(), |
| 719 | "timestamp_end"_a=cacheEventTimer.end_time(), |
| 720 | "description"_a="Pull from CacheMachine CacheData object type {}"_format(dataType)); |
| 721 | } |
| 722 | |
| 723 | return output; |
| 724 | } |
| 725 | |
| 726 | bool CacheMachine::has_data_in_index_now(size_t index){ |
| 727 | std::string message = this->cache_machine_name + "_" + std::to_string(index); |
no test coverage detected