| 145 | } |
| 146 | |
| 147 | void TRemoteQueryProcessor::RunMaster(const TVector<TNetworkAddress>& baseSearcherAddrs, unsigned short masterListenPort) { |
| 148 | CHROMIUM_TRACE_FUNCTION(); |
| 149 | Y_ASSERT(Requester.Get() == nullptr); |
| 150 | BaseSearcherAddrs = baseSearcherAddrs; |
| 151 | LastCounts.resize(BaseSearcherAddrs.ysize(), TAtomicWrap(0)); |
| 152 | |
| 153 | SetRequester(CreateRequester( |
| 154 | masterListenPort, |
| 155 | [this](const TGUID& canceledReq) { QueryCancelCallback(canceledReq); }, |
| 156 | [this](TAutoPtr<TNetworkRequest>& nlReq) { IncomingQueryCallback(nlReq); }, |
| 157 | [this](TAutoPtr<TNetworkResponse> response) { ReplyCallback(response); })); |
| 158 | |
| 159 | MasterAddress = TNetworkAddress(HostName(), Requester->GetListenPort()); |
| 160 | DEBUG_LOG << "Listening on port: " << Requester->GetListenPort() << Endl; |
| 161 | // init base searchers |
| 162 | DEBUG_LOG << "Init base searchers" << Endl; |
| 163 | TIntrusivePtr<TMetaRequester> mr(new TMetaRequester(this)); |
| 164 | int searcherCount = GetSlaveCount(); |
| 165 | for (int i = 0; i < searcherCount; ++i) { |
| 166 | THostInitData initData; |
| 167 | initData.CompId = i; |
| 168 | initData.MasterAddress = MasterAddress; |
| 169 | initData.BaseSearcherAddrs = BaseSearcherAddrs; |
| 170 | TVector<char> cmdData; |
| 171 | SerializeToMem(&cmdData, initData); |
| 172 | mr->AddQuery(i, "init", &cmdData); |
| 173 | } |
| 174 | mr->GetResults(nullptr); |
| 175 | |
| 176 | // run_ping |
| 177 | TArray2D<TVector<float>> delayMatrixData; |
| 178 | |
| 179 | // try to reuse previous ping stats |
| 180 | bool needRunPing = true; |
| 181 | if (NFs::Exists(DELAY_MATRIX_NAME)) { |
| 182 | TDelayData data; |
| 183 | SerializeFromFile(DELAY_MATRIX_NAME, data); |
| 184 | if (data.BaseSearcherAddrs == BaseSearcherAddrs) { |
| 185 | DEBUG_LOG << "Reusing ping times from " << DELAY_MATRIX_NAME << Endl; |
| 186 | needRunPing = false; |
| 187 | delayMatrixData.Swap(data.DelayMatrixData); |
| 188 | } |
| 189 | } |
| 190 | |
| 191 | if (needRunPing) { |
| 192 | DEBUG_LOG << "Run ping times collection" << Endl; |
| 193 | int pingCount = 0, iterCount = 0; |
| 194 | while (pingCount < searcherCount * searcherCount * 10 && iterCount < 10) { |
| 195 | for (int i = 0; i < searcherCount; ++i) { |
| 196 | mr->AddQuery(i, "run_ping", nullptr); |
| 197 | } |
| 198 | TVector<TVector<char>> res; |
| 199 | mr->GetResults(&res); |
| 200 | |
| 201 | delayMatrixData.SetSizes(searcherCount, searcherCount); |
| 202 | Y_ASSERT(res.ysize() == searcherCount); |
| 203 | for (int srcCompId = 0; srcCompId < searcherCount; ++srcCompId) { |
| 204 | TAllPingResults stats; |
no test coverage detected