| 2622 | } |
| 2623 | |
| 2624 | void AdmissionController::AddPoolAndPerHostStatsUpdates( |
| 2625 | vector<TTopicDelta>* topic_updates) { |
| 2626 | // local_stats_ are updated eagerly except for backend_mem_reserved (which isn't used |
| 2627 | // for local admission control decisions). Update that now before sending local_stats_. |
| 2628 | for (auto& entry : pool_stats_) { |
| 2629 | entry.second.UpdateMemTrackerStats(); |
| 2630 | } |
| 2631 | if (pools_for_updates_.empty()) { |
| 2632 | // No pool updates means no changes to host stats as well, so just return. |
| 2633 | return; |
| 2634 | } |
| 2635 | topic_updates->push_back(TTopicDelta()); |
| 2636 | TTopicDelta& topic_delta = topic_updates->back(); |
| 2637 | topic_delta.topic_name = request_queue_topic_name_; |
| 2638 | for (const string& pool_name: pools_for_updates_) { |
| 2639 | DCHECK(pool_stats_.find(pool_name) != pool_stats_.end()); |
| 2640 | PoolStats* stats = GetPoolStats(pool_name); |
| 2641 | VLOG_ROW << "Sending topic update " << stats->DebugString(); |
| 2642 | topic_delta.topic_entries.push_back(TTopicItem()); |
| 2643 | TTopicItem& topic_item = topic_delta.topic_entries.back(); |
| 2644 | topic_item.key = MakePoolTopicKey(pool_name, host_id_); |
| 2645 | Status status = |
| 2646 | thrift_serializer_.SerializeToString(&stats->local_stats(), &topic_item.value); |
| 2647 | if (!status.ok()) { |
| 2648 | LOG(WARNING) << "Failed to serialize query pool stats: " << status.GetDetail(); |
| 2649 | topic_delta.topic_entries.pop_back(); |
| 2650 | } |
| 2651 | } |
| 2652 | pools_for_updates_.clear(); |
| 2653 | |
| 2654 | // Now add the host stats |
| 2655 | topic_delta.topic_entries.push_back(TTopicItem()); |
| 2656 | TTopicItem& topic_item = topic_delta.topic_entries.back(); |
| 2657 | topic_item.key = Substitute("$0$1", TOPIC_KEY_STAT_PREFIX, host_id_); |
| 2658 | TPerHostStatsUpdate update; |
| 2659 | for (const auto& elem : host_stats_) { |
| 2660 | update.per_host_stats.emplace_back(); |
| 2661 | TPerHostStatsUpdateElement& inserted_elem = update.per_host_stats.back(); |
| 2662 | inserted_elem.__set_host_addr(elem.first); |
| 2663 | inserted_elem.__set_stats(elem.second); |
| 2664 | } |
| 2665 | Status status = thrift_serializer_.SerializeToString(&update, &topic_item.value); |
| 2666 | if (!status.ok()) { |
| 2667 | LOG(WARNING) << "Failed to serialize host stats: " << status.GetDetail(); |
| 2668 | topic_delta.topic_entries.pop_back(); |
| 2669 | } |
| 2670 | } |
| 2671 | |
| 2672 | void AdmissionController::DequeueLoop() { |
| 2673 | unique_lock<mutex> lock(admission_ctrl_lock_); |
nothing calls this directly
no test coverage detected