| 406 | SRQ->PostReceive(FakeRecvMem->GetMemRegion(), 0, FakeRecvMem->GetData(), FakeRecvMem->GetSize()); |
| 407 | } |
| 408 | void Sync() { |
| 409 | Y_ASSERT(WasFlushed && "Have to call Flush() before data fill & Sync()"); |
| 410 | char* myData = (char*)MemBlock[CurrentBuffer]->GetData(); |
| 411 | |
| 412 | ui64 recvMask = FutureRecvMask; |
| 413 | FutureRecvMask = 0; |
| 414 | int recvDebt = 0; |
| 415 | for (int z = 0; z < Iterations.ysize(); ++z) { |
| 416 | const TIteration& iter = Iterations[z]; |
| 417 | for (int k = 0; k < iter.OutList.ysize(); ++k) { |
| 418 | const TSend& ss = iter.OutList[k]; |
| 419 | const TBlockInfo& remoteBlk = ss.RemoteBlocks[CurrentBuffer]; |
| 420 | ss.QP->PostRDMAWriteImm(remoteBlk.Addr + ss.DstOffset, remoteBlk.Key, ss.ImmData, |
| 421 | MemBlock[CurrentBuffer]->GetMemRegion(), 0, myData + ss.SrcOffset, ss.Length); |
| 422 | ++ActiveRDMACount; |
| 423 | //printf("-> %d, imm %d (%" PRId64 " bytes)\n", ss.DstRank, ss.ImmData, ss.Length); |
| 424 | //printf("send %d\n", ss.SrcOffset); |
| 425 | } |
| 426 | ibv_wc wc; |
| 427 | while ((recvMask & iter.RecvMask) != iter.RecvMask) { |
| 428 | int rv = CQ->Poll(&wc, 1); |
| 429 | if (rv > 0) { |
| 430 | Y_ABORT_UNLESS(wc.status == IBV_WC_SUCCESS, "AllGather::Sync fail, status %d", (int)wc.status); |
| 431 | if (wc.opcode == IBV_WC_RECV_RDMA_WITH_IMM) { |
| 432 | //printf("Got %d\n", wc.imm_data); |
| 433 | ++recvDebt; |
| 434 | ui64 newBit = ui64(1) << wc.imm_data; |
| 435 | if (recvMask & newBit) { |
| 436 | Y_ABORT_UNLESS((FutureRecvMask & newBit) == 0, "data from 2 Sync() ahead is impossible"); |
| 437 | FutureRecvMask |= newBit; |
| 438 | } else { |
| 439 | recvMask |= newBit; |
| 440 | } |
| 441 | } else if (wc.opcode == IBV_WC_RDMA_WRITE) { |
| 442 | --ActiveRDMACount; |
| 443 | } else { |
| 444 | Y_ASSERT(0); |
| 445 | } |
| 446 | } else { |
| 447 | if (recvDebt > 0) { |
| 448 | PostRecv(); |
| 449 | --recvDebt; |
| 450 | } |
| 451 | } |
| 452 | } |
| 453 | for (int k = 0; k < iter.ReduceList.ysize(); ++k) { |
| 454 | const TReduce& rr = iter.ReduceList[k]; |
| 455 | ReduceOp->Reduce(myData + rr.DstOffset, myData + rr.SrcOffset, DataSize); |
| 456 | //printf("Merge %d -> %d (%d bytes)\n", rr.SrcOffset, rr.DstOffset, DataSize); |
| 457 | } |
| 458 | //printf("Iteration %d done\n", z); |
| 459 | } |
| 460 | while (recvDebt > 0) { |
| 461 | PostRecv(); |
| 462 | --recvDebt; |
| 463 | } |
| 464 | CurrentOffset = ReadyOffset; |
| 465 | WasFlushed = false; |
nothing calls this directly
no test coverage detected