Merge remote-tracking branch 'nccl/master' into develop

This commit is contained in:
BertanDogancay
2024-10-02 09:29:22 -05:00
54 changed files with 2158 additions and 964 deletions
+168 -114
View File
@@ -202,6 +202,7 @@ ncclResult_t ncclGetUniqueId_impl(ncclUniqueId* out) {
void NCCL_NO_OPTIMIZE commPoison(ncclComm_t comm) {
// Important that this does not trash intraComm0.
comm->rank = comm->cudaDev = comm->busId = comm->nRanks = -1;
comm->startMagic = comm->endMagic = 0;
}
RCCL_PARAM(KernelCollTraceEnable, "KERNEL_COLL_TRACE_ENABLE", 0);
@@ -491,7 +492,6 @@ static ncclResult_t dmaBufSupported(struct ncclComm* comm) {
ncclResult_t ncclCommEnsureReady(ncclComm_t comm) {
/* comm must be ready, or error will be reported */
ncclResult_t ret = ncclSuccess;
if (__atomic_load_n(comm->abortFlag, __ATOMIC_RELAXED)) {
ncclGroupJobAbort(comm->groupJob);
} else {
@@ -572,6 +572,7 @@ static ncclResult_t commAlloc(struct ncclComm* comm, struct ncclComm* parent, in
ncclMemoryPoolConstruct(&comm->memPool_ncclProxyOp);
ncclMemoryPoolConstruct(&comm->memPool_ncclPointerList);
ncclMemoryPoolConstruct(&comm->memPool_ncclNvlsHandleList);
ncclMemoryPoolConstruct(&comm->memPool_ncclCollnetHandleList);
comm->groupNext = reinterpret_cast<struct ncclComm*>(0x1);
comm->preconnectNext = reinterpret_cast<struct ncclComm*>(0x1);
@@ -848,9 +849,8 @@ static ncclResult_t computeBuffSizes(struct ncclComm* comm) {
comm->buffSizes[p] = envs[p] != -2 ? envs[p] : defaults[p];
}
// MNNVL support
if (!comm->MNNVL && comm->nNodes > 1) comm->p2pChunkSize = ncclParamP2pNetChunkSize();
else if (comm->MNNVL || ncclTopoPathAllNVLink(comm->topo)) comm->p2pChunkSize = ncclParamP2pNvlChunkSize();
if (comm->nNodes > 1) comm->p2pChunkSize = ncclParamP2pNetChunkSize();
else if (ncclTopoPathAllNVLink(comm->topo)) comm->p2pChunkSize = ncclParamP2pNvlChunkSize();
else comm->p2pChunkSize = ncclParamP2pPciChunkSize();
// Make sure P2P chunksize is not larger than coll chunksize.
@@ -872,16 +872,38 @@ NCCL_PARAM(CollNetNodeThreshold, "COLLNET_NODE_THRESHOLD", 2);
NCCL_PARAM(NvbPreconnect, "NVB_PRECONNECT", 0);
NCCL_PARAM(AllocP2pNetLLBuffers, "ALLOC_P2P_NET_LL_BUFFERS", 0);
static ncclResult_t collNetInitRailRankMap(ncclComm_t comm) {
int rank = comm->rank;
uint64_t nonHeadMask = (1ull << comm->localRanks) - 1;
comm->collNetDenseToUserRank = ncclMemoryStackAlloc<int>(&comm->memPermanent, comm->nRanks);
comm->collNetUserToDenseRank = ncclMemoryStackAlloc<int>(&comm->memPermanent, comm->nRanks);
// initialize collNetUserToDenseRank[rank]
comm->collNetUserToDenseRank[rank] = -1;
for (int h = 0; h < comm->collNetHeadsNum; h++) {
nonHeadMask ^= 1ull << comm->rankToLocalRank[comm->collNetHeads[h]];
if (comm->collNetHeads[h] == rank) { comm->collNetUserToDenseRank[rank] = h; break; }
}
if (comm->collNetUserToDenseRank[rank] == -1) {
comm->collNetUserToDenseRank[rank] = __builtin_popcountll(nonHeadMask & ((1ull << comm->localRank) - 1));
}
comm->collNetUserToDenseRank[rank] += comm->node * comm->localRanks;
NCCLCHECK(bootstrapAllGather(comm->bootstrap, comm->collNetUserToDenseRank, sizeof(int)));
for (int r = 0; r < comm->nRanks; r++) {
comm->collNetDenseToUserRank[comm->collNetUserToDenseRank[r]] = r;
}
return ncclSuccess;
}
static ncclResult_t collNetTrySetup(ncclComm_t comm, ncclComm_t parent, struct ncclTopoGraph* collNetGraph) {
ncclResult_t ret = ncclSuccess;
int* heads = NULL;
int rank = comm->rank;
int collNetSetupFail = 0;
int highestTypes[NCCL_MAX_LOCAL_RANKS] = { TRANSPORT_P2P };
// Find all head ranks
int nHeads = collNetGraph->nChannels;
int nHeadsUnique = 0;
int headsUnique[NCCL_MAX_LOCAL_RANKS];
int* headsUnique = NULL;
int highestTransportType0, highestTransportType1;
char line[1024];
bool share;
@@ -892,27 +914,26 @@ static ncclResult_t collNetTrySetup(ncclComm_t comm, ncclComm_t parent, struct n
};
struct collnetShareInfo* infos = NULL;
NCCLCHECKGOTO(ncclCalloc(&heads, nHeads), ret, fail);
NCCLCHECKGOTO(ncclCalloc(&headsUnique, collNetGraph->nChannels), ret, fail);
{ uint64_t mask = 0;
// Head GPU index is always 0
for (int c = 0; c < nHeads; c++) {
heads[c] = collNetGraph->intra[c * comm->localRanks + 0];
assert(comm->rankToNode[heads[c]] == comm->node);
for (int c = 0; c < collNetGraph->nChannels; c++) {
int head = collNetGraph->intra[c * comm->localRanks + 0];
assert(comm->rankToNode[head] == comm->node);
uint64_t mask0 = mask;
mask |= 1ull<<comm->rankToLocalRank[heads[c]];
if (mask != mask0) headsUnique[nHeadsUnique++] = heads[c];
mask |= 1ull<<comm->rankToLocalRank[head];
if (mask != mask0) headsUnique[nHeadsUnique++] = head;
}
}
comm->collNetHeads = heads;
comm->collNetHeadsNum = nHeads;
comm->collNetHeadsUniqueNum = nHeadsUnique;
comm->collNetHeads = headsUnique;
comm->collNetHeadsNum = nHeadsUnique;
if (parent && parent->collNetSupport && parent->config.splitShare && parent->nNodes == comm->nNodes) {
NCCLCHECKGOTO(ncclCalloc(&infos, comm->nRanks), ret, fail);
/* check whether child can share collnet resources of parent. Since parent builds each collnet communicator
* based on heads with the same head position in each node, as long as the collnet heads of child comm
* can match parent's heads, we can let child communicator share parent's collnet resources. */
for (int h = 0; h < nHeads; ++h) {
for (int h = 0; h < nHeadsUnique; ++h) {
int prev = INT_MIN;
struct collnetShareInfo* myinfo;
@@ -920,7 +941,7 @@ static ncclResult_t collNetTrySetup(ncclComm_t comm, ncclComm_t parent, struct n
myinfo = infos + comm->rank;
memset(myinfo, 0, sizeof(struct collnetShareInfo));
/* find the child head position in parent collnet heads. */
if (heads[h] == comm->rank) {
if (headsUnique[h] == comm->rank) {
myinfo->headPosition = -1;
myinfo->isMaster = 1;
for (int th = 0; th < parent->collNetHeadsNum; ++th)
@@ -946,10 +967,11 @@ static ncclResult_t collNetTrySetup(ncclComm_t comm, ncclComm_t parent, struct n
if (share) {
if (myinfo->isMaster) {
comm->collNetSharedRes = parent->collNetSharedRes;
comm->collNetChannels = std::min(comm->nChannels, parent->collNetSharedRes->nChannels);
for (int c = 0; c < comm->collNetChannels; ++c)
for (int c = 0; c < comm->nChannels; ++c)
NCCLCHECKGOTO(initCollnetChannel(comm, c, parent, true), ret, fail);
}
NCCLCHECKGOTO(collNetInitRailRankMap(comm), ret, fail);
} else {
/* TODO: CX-6 and CX-7 both do not support multiple sharp resources per process, if child comm cannot
* share the sharp resource from parent, we cannot use sharp in this case. This restriction might be
@@ -965,35 +987,19 @@ static ncclResult_t collNetTrySetup(ncclComm_t comm, ncclComm_t parent, struct n
} else {
/* this allocated buffer will be freed on proxy side */
NCCLCHECK(ncclCalloc(&comm->collNetSharedRes, 1));
comm->collNetChannels = comm->collNetSharedRes->nChannels = comm->nChannels;
comm->collNetSharedRes->nChannels = comm->nChannels;
comm->collNetSharedRes->buffSize = comm->buffSizes[NCCL_PROTO_SIMPLE];
comm->collNetDenseToUserRank = ncclMemoryStackAlloc<int>(&comm->memPermanent, comm->nRanks);
comm->collNetUserToDenseRank = ncclMemoryStackAlloc<int>(&comm->memPermanent, comm->nRanks);
{ // initialize collNetUserToDenseRank[rank]
uint64_t nonHeadMask = (1ull<<comm->localRanks)-1;
comm->collNetUserToDenseRank[rank] = -1;
for (int h=0; h < nHeadsUnique; h++) {
nonHeadMask ^= 1ull<<comm->rankToLocalRank[headsUnique[h]];
if (headsUnique[h] == rank) { comm->collNetUserToDenseRank[rank] = h; break; }
}
if (comm->collNetUserToDenseRank[rank] == -1) {
comm->collNetUserToDenseRank[rank] = __builtin_popcountll(nonHeadMask & ((1ull<<comm->localRank)-1));
}
comm->collNetUserToDenseRank[rank] += comm->node*comm->localRanks;
}
NCCLCHECK(bootstrapAllGather(comm->bootstrap, comm->collNetUserToDenseRank, sizeof(int)));
for (int r=0; r < comm->nRanks; r++) {
comm->collNetDenseToUserRank[comm->collNetUserToDenseRank[r]] = r;
}
NCCLCHECKGOTO(collNetInitRailRankMap(comm), ret, fail);
for (int c = 0; c < comm->collNetChannels; c++) {
for (int c = 0; c < comm->nChannels; c++) {
struct ncclChannel* channel = comm->channels + c;
NCCLCHECKGOTO(initCollnetChannel(comm, c, parent, false), ret, fail);
for (int h = 0; h < nHeads; h++) {
const int head = heads[h];
collNetSetupFail |= ncclTransportCollNetSetup(comm, collNetGraph, channel, head, head, h, collNetRecv);
if (!collNetSetupFail) collNetSetupFail |= ncclTransportCollNetSetup(comm, collNetGraph, channel, head, head, h, collNetSend);
for (int h = 0; h < nHeadsUnique; h++) {
const int head = headsUnique[h];
ncclConnect connect;
collNetSetupFail |= ncclTransportCollNetSetup(comm, collNetGraph, channel, head, head, h, collNetRecv, &connect);
if (!collNetSetupFail) collNetSetupFail |= ncclTransportCollNetSetup(comm, collNetGraph, channel, head, head, h, collNetSend, &connect);
}
// Verify CollNet setup across ranks after trying the first channel
if (c == 0) {
@@ -1015,7 +1021,7 @@ static ncclResult_t collNetTrySetup(ncclComm_t comm, ncclComm_t parent, struct n
bool isHead = false;
matrix = nullptr;
NCCLCHECKGOTO(ncclCalloc(&matrix, comm->nRanks), ret, matrix_end);
for (int h = 0; h < nHeads; h++) isHead |= (heads[h] == comm->rank);
for (int h = 0; h < nHeadsUnique; h++) isHead |= (headsUnique[h] == comm->rank);
if (isHead) {
for (int ty=0; ty < ncclNumTypes; ty++) {
for (int i=0; i < 4; i++) {
@@ -1105,7 +1111,72 @@ fail:
}
// MNNVL: Flag to indicate whether to enable Multi-Node NVLink
NCCL_PARAM(MNNVL, "MNNVL", -2);
NCCL_PARAM(MNNVLEnable, "MNNVL_ENABLE", 2);
#if CUDART_VERSION >= 11030
#include <cuda.h>
#include "cudawrap.h"
// Determine if MNNVL support is available
static int checkMNNVL(struct ncclComm* comm) {
ncclResult_t ret = ncclSuccess;
// MNNVL requires cuMem to be enabled
if (!ncclCuMemEnable()) return 0;
// MNNVL also requires FABRIC handle support
int cudaDev;
int flag = 0;
CUdevice currentDev;
CUDACHECK(cudaGetDevice(&cudaDev));
CUCHECK(cuDeviceGet(&currentDev, cudaDev));
// Ignore error if CU_DEVICE_ATTRIBUTE_HANDLE_TYPE_FABRIC_SUPPORTED is not supported
(void) CUPFN(cuDeviceGetAttribute(&flag, CU_DEVICE_ATTRIBUTE_HANDLE_TYPE_FABRIC_SUPPORTED, currentDev));;
if (!flag) return 0;
// Check that all ranks have initialized the fabric fully
for (int i = 0; i < comm->nRanks; i++) {
if (comm->peerInfo[i].fabricInfo.state != NVML_GPU_FABRIC_STATE_COMPLETED) return 0;
}
// Determine our MNNVL domain/clique
NCCLCHECKGOTO(ncclCalloc(&comm->clique.ranks, comm->nRanks), ret, fail);
comm->clique.id = comm->peerInfo[comm->rank].fabricInfo.cliqueId;
for (int i = 0; i < comm->nRanks; i++) {
nvmlGpuFabricInfoV_t *fabricInfo1 = &comm->peerInfo[comm->rank].fabricInfo;
nvmlGpuFabricInfoV_t *fabricInfo2 = &comm->peerInfo[i].fabricInfo;
// Check if the cluster UUID and cliqueId match
// A zero UUID means we don't have MNNVL fabric info - disable MNNVL
if ((((long *)&fabricInfo2->clusterUuid)[0]|((long *)fabricInfo2->clusterUuid)[1]) == 0) goto fail;
if ((memcmp(fabricInfo1->clusterUuid, fabricInfo2->clusterUuid, NVML_GPU_FABRIC_UUID_LEN) == 0) &&
(fabricInfo1->cliqueId == fabricInfo2->cliqueId)) {
if (i == comm->rank) {
comm->cliqueRank = comm->clique.size;
}
comm->clique.ranks[comm->clique.size++] = i;
}
}
// Determine whether to enable MNNVL or not
comm->MNNVL = ncclParamMNNVLEnable() == 2 ? comm->clique.size > 1 : ncclParamMNNVLEnable();
INFO(NCCL_INIT, "MNNVL %d cliqueId %x cliqueSize %d cliqueRank %d ", comm->MNNVL, comm->clique.id, comm->clique.size, comm->cliqueRank);
if (comm->MNNVL) {
// Force the CUMEM handle type to be FABRIC for MNNVL
ncclCuMemHandleType = CU_MEM_HANDLE_TYPE_FABRIC;
}
return comm->MNNVL;
fail:
if (comm->clique.ranks) free(comm->clique.ranks);
return 0;
}
#else
static int checkMNNVL(struct ncclComm* comm) {
return 0;
}
#endif
static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* parent = NULL) {
// We use 2 AllGathers
@@ -1130,6 +1201,7 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
float bwInter;
int typeIntra;
int typeInter;
int crossNic;
};
struct allGatherInfo {
@@ -1172,61 +1244,19 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
// AllGather1 - end
#if CUDART_VERSION >= 11030
#include <cuda.h>
#include "cudawrap.h"
// MNNVL support
if (nNodes > 1) {
int cliqueSize = 0;
comm->MNNVL = 0;
// Determine the size of the MNNVL domain/clique
for (int i = 0; i < nranks; i++) {
nvmlGpuFabricInfoV_t *fabricInfo1 = &comm->peerInfo[rank].fabricInfo;
nvmlGpuFabricInfoV_t *fabricInfo2 = &comm->peerInfo[i].fabricInfo;
// Check that the Fabric state is fully initialized
if (fabricInfo2->state != NVML_GPU_FABRIC_STATE_COMPLETED) continue;
// Check that the cluster UUID and cliqueId match in each rank
// A zero UUID means we don't have MNNVL fabric info - disable MNNVL
if ((((long *)&fabricInfo2->clusterUuid)[0]|((long *)fabricInfo2->clusterUuid)[1]) == 0) continue;
if ((memcmp(fabricInfo1->clusterUuid, fabricInfo2->clusterUuid, NVML_GPU_FABRIC_UUID_LEN) == 0) &&
(fabricInfo1->cliqueId == fabricInfo2->cliqueId)) {
cliqueSize++;
}
}
// Determine whether this is a MNNVL system
comm->MNNVL = ncclParamMNNVL() < 0 ? cliqueSize == comm->nRanks : ncclParamMNNVL();
// MNNVL requires cuMem to be enabled
if (!ncclCuMemEnable()) comm->MNNVL = 0;
if (comm->MNNVL) {
// MNNVL also requires FABRIC handle support
int cudaDev;
int flag = 0;
CUdevice currentDev;
CUDACHECK(cudaGetDevice(&cudaDev));
CUCHECK(cuDeviceGet(&currentDev, cudaDev));
// Ignore error if CU_DEVICE_ATTRIBUTE_HANDLE_TYPE_FABRIC_SUPPORTED is not supported
(void) CUPFN(cuDeviceGetAttribute(&flag, CU_DEVICE_ATTRIBUTE_HANDLE_TYPE_FABRIC_SUPPORTED, currentDev));;
if (!flag)
comm->MNNVL = 0;
else
// Force the handle type to be FABRIC for MNNVL
ncclCuMemHandleType = CU_MEM_HANDLE_TYPE_FABRIC;
}
if (ncclParamMNNVL() == 1 && !comm->MNNVL) {
WARN("MNNVL is not supported on this system");
ret = ncclSystemError;
goto fail;
}
if (nNodes > 1 && !checkMNNVL(comm) && ncclParamMNNVLEnable() == 1) {
// Return an error if the user specifically requested MNNVL support
WARN("MNNVL is not supported on this system");
ret = ncclSystemError;
goto fail;
}
#endif
do {
// Compute intra-process ranks
int intraProcRank0 = -1, intraProcRank = -1, intraProcRanks = 0;
for (int i = 0; i < nranks; i++) comm->minCompCap = std::min(comm->minCompCap, comm->peerInfo[rank].cudaCompCap);
for (int i = 0; i < nranks; i++) comm->maxCompCap = std::max(comm->maxCompCap, comm->peerInfo[rank].cudaCompCap);
for (int i = 0; i < nranks; i++) comm->minCompCap = std::min(comm->minCompCap, comm->peerInfo[i].cudaCompCap);
for (int i = 0; i < nranks; i++) comm->maxCompCap = std::max(comm->maxCompCap, comm->peerInfo[i].cudaCompCap);
comm->nvlsRegSupport = 1;
for (int i = 0; i < nranks; i++) {
@@ -1252,6 +1282,10 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
}
}
}
// Buffer Registration is not supported with MNNVL
if (comm->MNNVL) comm->nvlsRegSupport = 0;
TRACE(NCCL_INIT,"pidHash[%d] %lx intraProcRank %d intraProcRanks %d intraProcRank0 %d",
rank, comm->peerInfo[rank].pidHash, intraProcRank, intraProcRanks, intraProcRank0);
if (intraProcRank == -1 || intraProcRank0 == -1 || comm->peerInfo[intraProcRank0].comm == NULL) {
@@ -1465,6 +1499,7 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
allGather3Data[rank].graphInfo[a].bwInter = graphs[a]->bwInter;
allGather3Data[rank].graphInfo[a].typeIntra = graphs[a]->typeIntra;
allGather3Data[rank].graphInfo[a].typeInter = graphs[a]->typeInter;
allGather3Data[rank].graphInfo[a].crossNic = graphs[a]->crossNic;
}
comm->nChannels = std::min(treeGraph.nChannels, ringGraph.nChannels);
@@ -1543,10 +1578,11 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
graphs[a]->bwInter = std::min(allGather3Data[i].graphInfo[a].bwInter, graphs[a]->bwInter);
graphs[a]->typeIntra = std::max(allGather3Data[i].graphInfo[a].typeIntra, graphs[a]->typeIntra);
graphs[a]->typeInter = std::max(allGather3Data[i].graphInfo[a].typeInter, graphs[a]->typeInter);
graphs[a]->crossNic = std::max(allGather3Data[i].graphInfo[a].crossNic, graphs[a]->crossNic);
}
if (graphs[NCCL_ALGO_COLLNET_CHAIN]->nChannels == 0) comm->collNetSupport = 0;
if (graphs[NCCL_ALGO_NVLS]->nChannels == 0) comm->nvlsSupport = 0;
}
if (graphs[NCCL_ALGO_COLLNET_CHAIN]->nChannels == 0) comm->collNetSupport = 0;
if (graphs[NCCL_ALGO_NVLS]->nChannels == 0) comm->nvlsSupport = comm->nvlsChannels = 0;
comm->nChannels = treeGraph.nChannels = ringGraph.nChannels =
(comm->topo->nodes[GPU].count != comm->topo->nRanks && comm->topo->nodes[NET].count)
@@ -1564,18 +1600,23 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
INFO(NCCL_INIT, "Communicator has %d nodes which is less than CollNet node threshold %d, disabling CollNet", comm->nNodes, collNetNodeThreshold);
comm->collNetSupport = 0;
}
comm->collNetRegSupport = true;
for (int n=0; n<comm->nNodes; n++) {
if (comm->nodeRanks[n].localRanks > NCCL_MAX_DIRECT_ARITY+1) {
WARN("CollNet currently only supports up to %d GPUs per node, disabling CollNet", NCCL_MAX_DIRECT_ARITY+1);
comm->collNetSupport = 0;
break;
}
if (comm->nodeRanks[n].localRanks > 1) {
// As long as there is more than 1 rank on any node, we need to disable collnet reg
comm->collNetRegSupport = false;
}
}
}
NCCLCHECKGOTO(ncclCalloc(&rings, nranks*MAXCHANNELS), ret, fail);
NCCLCHECKGOTO(ncclTopoPostset(comm, nodesFirstRank, nodesTreePatterns, allTopoRanks, rings, graphs, nc), ret, fail);
NCCLCHECKGOTO(ncclTopoPostset(comm, nodesFirstRank, nodesTreePatterns, allTopoRanks, rings, graphs, parent, nc), ret, fail);
if (comm->topo->treeDefined) NCCLCHECK(ncclTreeBasePostset(comm, &treeGraph));
// AllGather3 - end
@@ -1679,7 +1720,7 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
// Compute time models for algorithm and protocol combinations
NCCLCHECKGOTO(ncclTopoTuneModel(comm, comm->minCompCap, comm->maxCompCap, graphs), ret, fail);
INFO(NCCL_INIT, "%d coll channels, %d collnet channels, %d nvls channels, %d p2p channels, %d p2p channels per peer", comm->nChannels, comm->collNetChannels, comm->nvlsChannels, comm->p2pnChannels, comm->p2pnChannelsPerPeer);
INFO(NCCL_INIT, "%d coll channels, %d collnet channels, %d nvls channels, %d p2p channels, %d p2p channels per peer", comm->nChannels, comm->nChannels, comm->nvlsChannels, comm->p2pnChannels, comm->p2pnChannelsPerPeer);
do { // Setup p2p structures in comm->tasks
struct ncclTasks* tasks = &comm->tasks;
@@ -1792,7 +1833,7 @@ static ncclResult_t initTransportsRank(struct ncclComm* comm, struct ncclComm* p
}
/* Local intra-node barrier */
NCCLCHECKGOTO(bootstrapBarrier(comm->bootstrap, comm->localRankToRank, comm->localRank, comm->localRanks, comm->localRankToRank[0]), ret, fail);
NCCLCHECKGOTO(bootstrapIntraNodeBarrier(comm->bootstrap, comm->localRankToRank, comm->localRank, comm->localRanks, comm->localRankToRank[0]), ret, fail);
// We should have allocated all buffers, collective fifos, ... we can
// restore the affinity.
@@ -1949,7 +1990,13 @@ static ncclResult_t ncclCommInitRankFunc(struct ncclAsyncJob* job_) {
comm->cudaArch = cudaArch;
comm->commHash = getHash(job->commId.internal, NCCL_UNIQUE_ID_BYTES);
INFO(NCCL_INIT,"comm %p rank %d nranks %d cudaDev %d busId %lx commId 0x%llx - Init START", comm, comm->rank, comm->nRanks, comm->cudaDev, comm->busId, (unsigned long long)hashUniqueId(job->commId));
if (job->parent) {
INFO(NCCL_INIT,"ncclCommSplit comm %p rank %d nranks %d cudaDev %d busId %lx parent %p color %d key %d commId 0x%llx - Init START",
comm, comm->rank, comm->nRanks, comm->cudaDev, comm->busId, job->parent, job->color, job->key, (unsigned long long)hashUniqueId(job->commId));
} else {
INFO(NCCL_INIT,"ncclCommInitRank comm %p rank %d nranks %d cudaDev %d busId %lx commId 0x%llx - Init START",
comm, comm->rank, comm->nRanks, comm->cudaDev, comm->busId, (unsigned long long)hashUniqueId(job->commId));
}
NCCLCHECKGOTO(initTransportsRank(comm, job->parent), res, fail);
@@ -2006,9 +2053,9 @@ static ncclResult_t ncclCommInitRankFunc(struct ncclAsyncJob* job_) {
#endif
}
NCCLCHECKGOTO(ncclLoadTunerPlugin(&comm->tuner), res, fail);
NCCLCHECKGOTO(ncclTunerPluginLoad(&comm->tuner), res, fail);
if (comm->tuner) {
NCCLCHECK(comm->tuner->init(comm->nRanks, comm->nNodes, ncclDebugLog));
NCCLCHECK(comm->tuner->init(comm->nRanks, comm->nNodes, ncclDebugLog, &comm->tunerContext));
}
// update communicator state
@@ -2025,8 +2072,13 @@ static ncclResult_t ncclCommInitRankFunc(struct ncclAsyncJob* job_) {
comm, comm->nRanks, (unsigned long long)hashUniqueId(job->commId), comm->rank, comm->cudaDev);
}
INFO(NCCL_INIT,"comm %p rank %d nranks %d cudaDev %d busId %lx commId 0x%llx localSize %zi used %ld bytes on core %d - Init COMPLETE", comm, comm->rank, comm->nRanks, comm->cudaDev, comm->busId, (unsigned long long)hashUniqueId(job->commId), maxLocalSizeBytes, allocTracker[comm->cudaDev].totalAllocSize, sched_getcpu());
if (job->parent) {
INFO(NCCL_INIT,"ncclCommSplit comm %p rank %d nranks %d cudaDev %d busId %lx parent %p color %d key %d commId 0x%llx localSize %zi used %ld bytes on core %d - Init COMPLETE",
comm, comm->rank, comm->nRanks, comm->cudaDev, comm->busId, job->parent, job->color, job->key, (unsigned long long)hashUniqueId(job->commId), maxLocalSizeBytes, allocTracker[comm->cudaDev].totalAllocSize, sched_getcpu());
} else {
INFO(NCCL_INIT,"ncclCommInitRank comm %p rank %d nranks %d cudaDev %d busId %lx commId 0x%llx localSize %zi used %ld bytes on core %d - Init COMPLETE",
comm, comm->rank, comm->nRanks, comm->cudaDev, comm->busId, (unsigned long long)hashUniqueId(job->commId), maxLocalSizeBytes, allocTracker[comm->cudaDev].totalAllocSize, sched_getcpu());
}
exit:
if (job->newcomm) {
/* assign it to user pointer. */
@@ -2236,6 +2288,7 @@ static ncclResult_t ncclCommInitRankDev(ncclComm_t* newcomm, int nranks, ncclUni
}
NCCLCHECKGOTO(ncclCalloc(&comm, 1), res, fail);
comm->startMagic = comm->endMagic = NCCL_MAGIC; // Used to detect comm corruption.
NCCLCHECKGOTO(ncclCudaHostCalloc((uint32_t**)&comm->abortFlag, 1), res, fail);
NCCLCHECKGOTO(ncclCalloc((uint32_t**)&comm->abortFlagRefCount, 1), res, fail);
*comm->abortFlagRefCount = 1;
@@ -2434,8 +2487,8 @@ static ncclResult_t commCleanup(ncclComm_t comm) {
}
if (comm->tuner != NULL) {
NCCLCHECK(comm->tuner->destroy());
NCCLCHECK(ncclCloseTunerPlugin(&comm->tuner));
NCCLCHECK(comm->tuner->destroy(comm->tunerContext));
NCCLCHECK(ncclTunerPluginUnload(&comm->tuner));
}
NCCLCHECK(commFree(comm));
@@ -2688,7 +2741,7 @@ ncclResult_t ncclCommSplit_impl(ncclComm_t comm, int color, int key, ncclComm_t
ncclResult_t res = ncclSuccess;
NCCLCHECK(ncclGroupStartInternal());
NCCLCHECKGOTO(PtrCheck(comm, "CommSplit", "comm"), res, fail);
NCCLCHECKGOTO(CommCheck(comm, "CommSplit", "comm"), res, fail);
NCCLCHECKGOTO(PtrCheck(newcomm, "CommSplit", "newcomm"), res, fail);
NCCLCHECKGOTO(ncclCommEnsureReady(comm), res, fail);
@@ -2698,6 +2751,7 @@ ncclResult_t ncclCommSplit_impl(ncclComm_t comm, int color, int key, ncclComm_t
INFO(NCCL_INIT, "Rank %d has color with NCCL_SPLIT_NOCOLOR, not creating a new communicator", comm->rank);
} else {
NCCLCHECKGOTO(ncclCalloc(&childComm, 1), res, fail);
childComm->startMagic = childComm->endMagic = NCCL_MAGIC;
if (comm->config.splitShare) {
childComm->abortFlag = comm->abortFlag;
childComm->abortFlagRefCount = comm->abortFlagRefCount;
@@ -2770,7 +2824,7 @@ const char* ncclGetLastError_impl(ncclComm_t comm) {
NCCL_API(ncclResult_t, ncclCommGetAsyncError, ncclComm_t comm, ncclResult_t *asyncError);
ncclResult_t ncclCommGetAsyncError_impl(ncclComm_t comm, ncclResult_t *asyncError) {
NCCLCHECK(PtrCheck(comm, "ncclGetAsyncError", "comm"));
NCCLCHECK(CommCheck(comm, "ncclGetAsyncError", "comm"));
NCCLCHECK(PtrCheck(asyncError, "ncclGetAsyncError", "asyncError"));
*asyncError = __atomic_load_n(&comm->asyncResult, __ATOMIC_ACQUIRE);
@@ -2782,7 +2836,7 @@ NCCL_API(ncclResult_t, ncclCommCount, const ncclComm_t comm, int* count);
ncclResult_t ncclCommCount_impl(const ncclComm_t comm, int* count) {
NVTX3_FUNC_RANGE_IN(nccl_domain);
NCCLCHECK(PtrCheck(comm, "CommCount", "comm"));
NCCLCHECK(CommCheck(comm, "CommCount", "comm"));
NCCLCHECK(PtrCheck(count, "CommCount", "count"));
/* init thread must be joined before we access the attributes of comm. */
@@ -2796,7 +2850,7 @@ NCCL_API(ncclResult_t, ncclCommCuDevice, const ncclComm_t comm, int* devid);
ncclResult_t ncclCommCuDevice_impl(const ncclComm_t comm, int* devid) {
NVTX3_FUNC_RANGE_IN(nccl_domain);
NCCLCHECK(PtrCheck(comm, "CommCuDevice", "comm"));
NCCLCHECK(CommCheck(comm, "CommCuDevice", "comm"));
NCCLCHECK(PtrCheck(devid, "CommCuDevice", "devid"));
NCCLCHECK(ncclCommEnsureReady(comm));
@@ -2809,7 +2863,7 @@ NCCL_API(ncclResult_t, ncclCommUserRank, const ncclComm_t comm, int* rank);
ncclResult_t ncclCommUserRank_impl(const ncclComm_t comm, int* rank) {
NVTX3_FUNC_RANGE_IN(nccl_domain);
NCCLCHECK(PtrCheck(comm, "CommUserRank", "comm"));
NCCLCHECK(CommCheck(comm, "CommUserRank", "comm"));
NCCLCHECK(PtrCheck(rank, "CommUserRank", "rank"));
NCCLCHECK(ncclCommEnsureReady(comm));
@@ -2848,7 +2902,7 @@ ncclResult_t ncclMemAlloc_impl(void **ptr, size_t size) {
if (mcSupport) {
memprop.type = CU_MEM_ALLOCATION_TYPE_PINNED;
memprop.location.type = CU_MEM_LOCATION_TYPE_DEVICE;
memprop.requestedHandleTypes = NVLS_CU_MEM_HANDLE_TYPE;
memprop.requestedHandleTypes = ncclCuMemHandleType;
memprop.location.id = currentDev;
// Query device to see if RDMA support is available
CUCHECK(cuDeviceGetAttribute(&flag, CU_DEVICE_ATTRIBUTE_GPU_DIRECT_RDMA_SUPPORTED, currentDev));
@@ -2860,7 +2914,7 @@ ncclResult_t ncclMemAlloc_impl(void **ptr, size_t size) {
mcprop.size = size;
/* device cnt is a dummy value right now, it might affect mc granularity in the future. */
mcprop.numDevices = dcnt;
mcprop.handleTypes = NVLS_CU_MEM_HANDLE_TYPE;
mcprop.handleTypes = ncclCuMemHandleType;
mcprop.flags = 0;
CUCHECK(cuMulticastGetGranularity(&mcGran, &mcprop, CU_MULTICAST_GRANULARITY_RECOMMENDED));