| 731 | } |
| 732 | |
| 733 | void SchedulerWrapper::InitializeScheduler() { |
| 734 | DCHECK(scheduler_ == nullptr); |
| 735 | DCHECK_GT(plan_.cluster().NumHosts(), 0) << "Cannot initialize scheduler with 0 " |
| 736 | << "hosts."; |
| 737 | const Host& scheduler_host = plan_.cluster().hosts()[0]; |
| 738 | string scheduler_backend_id = scheduler_host.ip; |
| 739 | cluster_membership_mgr_.reset( |
| 740 | new ClusterMembershipMgr(scheduler_backend_id, nullptr, &metrics_)); |
| 741 | cluster_membership_mgr_->SetLocalBeDescFn( |
| 742 | [scheduler_host]() { return BuildBackendDescriptor(scheduler_host); }); |
| 743 | Status status = cluster_membership_mgr_->Init(); |
| 744 | DCHECK(status.ok()) << "Cluster membership manager init failed in test"; |
| 745 | scheduler_.reset(new Scheduler(&metrics_, nullptr)); |
| 746 | // Initialize the cluster membership manager |
| 747 | SendFullMembershipMap(); |
| 748 | } |
| 749 | |
| 750 | void SchedulerWrapper::AddHostToTopicDelta(const Host& host, TTopicDelta* delta) const { |
| 751 | DCHECK_GT(host.be_port, 0) << "Host cannot be added to scheduler without a running " |
nothing calls this directly
no test coverage detected