Merge pull request #1559 from BertanDogancay/2.23

[SYNC] 2.23.4-1
This commit is contained in:
Bertan Dogancay
2025-03-28 17:06:56 -04:00
committed by GitHub
92 changed files with 7322 additions and 2168 deletions
+380 -121
View File
@@ -18,6 +18,7 @@
#include "channel.h"
#include "rocmwrap.h"
#include "rccl_vars.h"
#include "profiler.h"
#include "transport.h"
#include "common.h"
#include "api_trace.h"
@@ -160,6 +161,10 @@ static void addWorkBatchToPlan(
if (newBatch || extendBatch) {
if (!newBatch) batch->nextExtends = extendBatch; // Extending the previous batch.
struct ncclWorkBatchList* batchNode = ncclMemoryStackAlloc<ncclWorkBatchList>(&comm->memScoped);
// Coverity thinks that ncclIntruQueueEnqueue will access chan->workBatchQueue->tail, which might
// be NULL. But that code is guarded by chan->workBatchQueue->head not being NULL, in which
// case tail won't be NULL either.
// coverity[var_deref_model:FALSE]
ncclIntruQueueEnqueue(&chan->workBatchQueue, batchNode);
batch = &batchNode->batch;
batch->nextExtends = 0;
@@ -280,7 +285,29 @@ static ncclResult_t cleanupIpc(struct ncclComm* comm, struct ncclCommCallback* c
return ncclSuccess;
}
static ncclResult_t registerIntraNodeBuffers(
static ncclResult_t registerCheckP2PConnection(struct ncclComm* comm, struct ncclConnector* conn, struct ncclTopoGraph* graph, int peer, bool* needReg) {
if (conn->connected) {
if (conn->conn.flags & (NCCL_IPC_READ | NCCL_IPC_WRITE | NCCL_DIRECT_READ | NCCL_DIRECT_WRITE)) {
*needReg = true;
} else {
// network connection
*needReg = false;
}
} else {
struct ncclPeerInfo* peerInfo = &comm->peerInfo[peer];
struct ncclPeerInfo* myInfo = &comm->peerInfo[comm->rank];
int canConnect = 0;
NCCLCHECK(ncclTransports[0]->canConnect(&canConnect, comm, graph, myInfo, peerInfo));
if (canConnect) {
*needReg = true;
} else {
*needReg = false;
}
}
return ncclSuccess;
}
static ncclResult_t registerCollBuffers(
struct ncclComm* comm, struct ncclTaskColl* info,
void* outRegBufSend[NCCL_MAX_LOCAL_RANKS],
void* outRegBufRecv[NCCL_MAX_LOCAL_RANKS],
@@ -291,8 +318,10 @@ static ncclResult_t registerIntraNodeBuffers(
info->regBufType = NCCL_REGULAR_BUFFER;
*regNeedConnect = true;
if (!(ncclParamLocalRegister() || (comm->planner.persistent && ncclParamGraphRegister()))) goto exit;
#if CUDART_VERSION >= 11030
if ((info->algorithm == NCCL_ALGO_NVLS || info->algorithm == NCCL_ALGO_NVLS_TREE) && comm->nvlsRegSupport) {
if (info->algorithm == NCCL_ALGO_NVLS || info->algorithm == NCCL_ALGO_NVLS_TREE) {
if (!comm->nvlsRegSupport || info->opDev.op == ncclDevPreMulSum) goto exit;
bool regBufUsed = false;
const void *sendbuff = info->sendbuff;
void *recvbuff = info->recvbuff;
@@ -325,60 +354,6 @@ static ncclResult_t registerIntraNodeBuffers(
}
info->regBufType = NCCL_NVLS_REG_BUFFER;
}
} else if (info->algorithm == NCCL_ALGO_COLLNET_DIRECT && // limited to CollNetDirect for now
comm->intraHighestTransportType == TRANSPORT_P2P && // only when all ranks can p2p each other
comm->intraRanks < comm->localRanks && // only with inter-process & intra-node peers
comm->planner.persistent && 0) {
/* Disable CollnetDirect registration since it does not support cuMem* allocated memory. */
int localRank = comm->localRank;
cudaPointerAttributes sattr, rattr;
CUDACHECK(cudaPointerGetAttributes(&sattr, info->sendbuff));
CUDACHECK(cudaPointerGetAttributes(&rattr, info->recvbuff));
if (sattr.type != cudaMemoryTypeDevice || rattr.type != cudaMemoryTypeDevice) return ncclSuccess;
if (CUPFN(cuMemGetAddressRange) == nullptr) return ncclSuccess;
struct HandlePair {
cudaIpcMemHandle_t ipc[2]; // {send, recv}
size_t offset[2]; // {send, recv}
};
struct HandlePair handles[NCCL_MAX_LOCAL_RANKS];
CUDACHECKGOTO(cudaIpcGetMemHandle(&handles[localRank].ipc[0], (void*)info->sendbuff), result, fallback);
CUDACHECKGOTO(cudaIpcGetMemHandle(&handles[localRank].ipc[1], (void*)info->recvbuff), result, fallback);
void *baseSend, *baseRecv;
size_t size;
CUCHECK(cuMemGetAddressRange((CUdeviceptr *)&baseSend, &size, (CUdeviceptr)info->sendbuff));
handles[localRank].offset[0] = (char*)info->sendbuff - (char*)baseSend;
CUCHECK(cuMemGetAddressRange((CUdeviceptr *)&baseRecv, &size, (CUdeviceptr)info->recvbuff));
handles[localRank].offset[1] = (char*)info->recvbuff - (char*)baseRecv;
NCCLCHECK(bootstrapIntraNodeAllGather(comm->bootstrap, comm->localRankToRank, comm->localRank, comm->localRanks, handles, sizeof(struct HandlePair)));
// Open handles locally
for (int i=0; i < comm->localRanks; i++) {
if (i == localRank) { // Skip self
outRegBufSend[i] = nullptr;
outRegBufRecv[i] = nullptr;
} else {
for (int sr=0; sr < 2; sr++) {
// Get base address of mapping
void* base;
CUDACHECK(cudaIpcOpenMemHandle(&base, handles[i].ipc[sr], cudaIpcMemLazyEnablePeerAccess));
// Get real buffer address by adding offset in the mapping
(sr == 0 ? outRegBufSend : outRegBufRecv)[i] = (char*)base + handles[i].offset[sr];
// Enqueue reminder to close memory handle
struct ncclIpcCleanupCallback* cb = (struct ncclIpcCleanupCallback*)malloc(sizeof(struct ncclIpcCleanupCallback));
cb->base.fn = cleanupIpc;
cb->ptr = base;
ncclIntruQueueEnqueue(cleanupQueue, &cb->base);
info->nCleanupQueueElts += 1;
}
}
}
info->regBufType = NCCL_IPC_REG_BUFFER;
} else if ((info->algorithm == NCCL_ALGO_COLLNET_DIRECT || info->algorithm == NCCL_ALGO_COLLNET_CHAIN) && comm->collNetRegSupport && info->opDev.op != ncclDevPreMulSum && info->opDev.op != ncclDevSumPostDiv) {
size_t elementSize = ncclTypeSize(info->datatype);
size_t sendbuffSize = elementSize*ncclFuncSendCount(info->func, comm->nRanks, info->count);
@@ -397,27 +372,200 @@ static ncclResult_t registerIntraNodeBuffers(
}
if ((sendRegBufFlag == 0 || recvRegBufFlag == 0) && comm->planner.persistent && ncclParamGraphRegister()) {
ncclCollnetGraphRegisterBuffer(comm, info->sendbuff, sendbuffSize, collNetSend, &sendRegBufFlag, &sendHandle, cleanupQueue, &info->nCleanupQueueElts);
info->sendMhandle = sendHandle;
if (sendRegBufFlag) {
if (!sendRegBufFlag) {
ncclCollnetGraphRegisterBuffer(comm, info->sendbuff, sendbuffSize, collNetSend, &sendRegBufFlag, &sendHandle, cleanupQueue, &info->nCleanupQueueElts);
info->sendMhandle = sendHandle;
}
if (sendRegBufFlag && !recvRegBufFlag) {
ncclCollnetGraphRegisterBuffer(comm, info->recvbuff, recvbuffSize, collNetRecv, &recvRegBufFlag, &recvHandle, cleanupQueue, &info->nCleanupQueueElts);
info->recvMhandle = recvHandle;
}
}
if (sendRegBufFlag && recvRegBufFlag) {
info->nMaxChannels = std::max(comm->config.minCTAs, std::min(comm->config.maxCTAs, 1));
info->nMaxChannels = 1;
info->regBufType = NCCL_COLLNET_REG_BUFFER;
if (sendRegBufFlag == 1 && recvRegBufFlag == 1) {
INFO(NCCL_REG, "rank %d successfully registered collNet sendbuff %p (handle %p), sendbuff size %ld, recvbuff %p (handle %p), recvbuff size %ld", comm->rank, info->sendbuff, sendHandle, sendbuffSize, info->recvbuff, recvHandle, recvbuffSize);
}
}
} else if (comm->intraNodeP2pSupport && info->protocol == NCCL_PROTO_SIMPLE) {
// IPC buffer registration
if (info->func == ncclFuncReduceScatter) goto exit;
if (info->algorithm == NCCL_ALGO_RING && ((info->func == ncclFuncAllReduce && info->sendbuff == info->recvbuff) || info->func == ncclFuncReduce)) goto exit;
if ((info->algorithm == NCCL_ALGO_TREE || info->algorithm == NCCL_ALGO_COLLNET_CHAIN) && info->sendbuff == info->recvbuff) goto exit;
if (info->func == ncclFuncAllGather && info->algorithm == NCCL_ALGO_PAT) goto exit;
int peerRanks[NCCL_MAX_LOCAL_RANKS];
int nPeers = 0;
size_t elementSize = ncclTypeSize(info->datatype);
size_t sendbuffSize = elementSize*ncclFuncSendCount(info->func, comm->nRanks, info->count);
size_t recvbuffSize = elementSize*ncclFuncRecvCount(info->func, comm->nRanks, info->count);
int regBufFlag = 0;
memset(peerRanks, 0xff, sizeof(int) * NCCL_MAX_LOCAL_RANKS);
if (info->algorithm == NCCL_ALGO_COLLNET_DIRECT) {
struct ncclChannel* channel = comm->channels;
for (int r = 0; r < NCCL_MAX_DIRECT_ARITY; ++r) {
for (int updown = 0; updown < 2; ++updown) {
int peer;
if (updown == 0)
peer = channel->collnetDirect.up[r];
else
peer = channel->collnetDirect.down[r];
if (peer != -1) {
struct ncclConnector* peerConn = &channel->peers[peer]->recv[0];
bool needReg = false;
NCCLCHECK(registerCheckP2PConnection(comm, peerConn, &comm->graphs[NCCL_ALGO_COLLNET_DIRECT], peer, &needReg));
if (needReg) {
bool found = false;
for (int p = 0; p < nPeers; ++p) {
if (peerRanks[p] == peer) {
found = true;
break;
}
}
if (!found) peerRanks[nPeers++] = peer;
}
}
}
}
if (nPeers > 0) {
if (ncclParamLocalRegister())
ncclIpcLocalRegisterBuffer(comm, info->sendbuff, sendbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->sendbuffOffset, &info->sendbuffRmtAddrs);
if (!regBufFlag && comm->planner.persistent && ncclParamGraphRegister()) {
ncclIpcGraphRegisterBuffer(comm, info->sendbuff, sendbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->sendbuffOffset, &info->sendbuffRmtAddrs, cleanupQueue, &info->nCleanupQueueElts);
}
if (regBufFlag) {
if (ncclParamLocalRegister())
ncclIpcLocalRegisterBuffer(comm, info->recvbuff, recvbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->recvbuffOffset, &info->recvbuffRmtAddrs);
if (!regBufFlag && comm->planner.persistent && ncclParamGraphRegister()) {
ncclIpcGraphRegisterBuffer(comm, info->recvbuff, recvbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->recvbuffOffset, &info->recvbuffRmtAddrs, cleanupQueue, &info->nCleanupQueueElts);
}
}
}
if (regBufFlag) {
info->regBufType = NCCL_IPC_REG_BUFFER;
}
} else if (info->algorithm == NCCL_ALGO_RING) {
struct ncclReg* recvRegRecord;
NCCLCHECK(ncclRegFind(comm, info->recvbuff, recvbuffSize, &recvRegRecord));
if (recvRegRecord == NULL) goto exit;
for (int c = 0; c < comm->nChannels; ++c) {
struct ncclChannel* channel = comm->channels + c;
for (int r = 0; r < 2; ++r) {
bool needReg = false;
int peer;
struct ncclConnector* peerConn;
// P2P transport
if (r == 0)
peer = channel->ring.prev;
else
peer = channel->ring.next;
peerConn = &channel->peers[peer]->recv[0];
NCCLCHECK(registerCheckP2PConnection(comm, peerConn, &comm->graphs[NCCL_ALGO_RING], peer, &needReg));
if (needReg) {
bool found = false;
for (int p = 0; p < nPeers; ++p) {
if (peerRanks[p] == peer) {
found = true;
break;
}
}
if (!found) peerRanks[nPeers++] = peer;
}
}
}
if (nPeers > 0) {
if (ncclParamLocalRegister()) {
ncclIpcLocalRegisterBuffer(comm, info->recvbuff, recvbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->recvbuffOffset, &info->recvbuffRmtAddrs);
}
if (!regBufFlag && comm->planner.persistent && ncclParamGraphRegister()) {
ncclIpcGraphRegisterBuffer(comm, info->recvbuff, recvbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->recvbuffOffset, &info->recvbuffRmtAddrs, cleanupQueue, &info->nCleanupQueueElts);
}
}
if (regBufFlag) {
info->regBufType = NCCL_IPC_REG_BUFFER;
}
} else if (info->algorithm == NCCL_ALGO_TREE || info->algorithm == NCCL_ALGO_COLLNET_CHAIN) {
struct ncclReg* recvRegRecord;
NCCLCHECK(ncclRegFind(comm, info->recvbuff, recvbuffSize, &recvRegRecord));
if (recvRegRecord == NULL) goto exit;
for (int c = 0; c < comm->nChannels; ++c) {
struct ncclChannel* channel = comm->channels + c;
struct ncclTree* tree = NULL;
int peers[NCCL_MAX_TREE_ARITY + 1];
if (info->algorithm == NCCL_ALGO_TREE)
tree = &channel->tree;
else
tree = &channel->collnetChain;
for (int p = 0; p < NCCL_MAX_TREE_ARITY; ++p) peers[p] = tree->down[p];
peers[NCCL_MAX_TREE_ARITY] = tree->up;
for (int p = 0; p < NCCL_MAX_TREE_ARITY + 1; ++p) {
int peer = peers[p];
bool peerNeedReg = false;
struct ncclConnector* recvConn = NULL;
// P2P transport
if (peer == -1 || peer == comm->nRanks) continue;
recvConn = &channel->peers[peer]->recv[0];
NCCLCHECK(registerCheckP2PConnection(comm, recvConn, &comm->graphs[info->algorithm], peer, &peerNeedReg));
if (peerNeedReg) {
bool found = false;
for (int pindex = 0; pindex < nPeers; ++pindex) {
if (peerRanks[pindex] == peer) {
found = true;
break;
}
}
if (!found) peerRanks[nPeers++] = peer;
}
}
}
if (nPeers > 0) {
if (ncclParamLocalRegister()) {
ncclIpcLocalRegisterBuffer(comm, info->recvbuff, recvbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->recvbuffOffset, &info->recvbuffRmtAddrs);
}
if (!regBufFlag && comm->planner.persistent && ncclParamGraphRegister()) {
ncclIpcGraphRegisterBuffer(comm, info->recvbuff, recvbuffSize, peerRanks, nPeers, NCCL_IPC_COLLECTIVE, &regBufFlag, &info->recvbuffOffset, &info->recvbuffRmtAddrs, cleanupQueue, &info->nCleanupQueueElts);
}
}
if (regBufFlag) {
info->regBufType = NCCL_IPC_REG_BUFFER;
}
}
if (info->regBufType == NCCL_IPC_REG_BUFFER && comm->nNodes == 1 && 16 < info->nMaxChannels && info->nMaxChannels <= 24) {
info->nMaxChannels = 16;
}
}
fallback:
#endif
exit:
return result;
}
static ncclResult_t registerP2pBuffer(struct ncclComm* comm, void* userbuff, int peerRank, size_t size, int* regFlag, void** regAddr, struct ncclIntruQueue<struct ncclCommCallback, &ncclCommCallback::next>* cleanupQueue) {
ncclResult_t ret = ncclSuccess;
uintptr_t offset = 0;
uintptr_t* peerRmtAddrs = NULL;
*regFlag = 0;
if (ncclParamLocalRegister()) {
ncclIpcLocalRegisterBuffer(comm, userbuff, size, &peerRank, 1, NCCL_IPC_SENDRECV, regFlag, &offset, &peerRmtAddrs);
}
if (*regFlag == 0 && comm->planner.persistent && ncclParamGraphRegister()) {
ncclIpcGraphRegisterBuffer(comm, userbuff, size, &peerRank, 1, NCCL_IPC_SENDRECV, regFlag, &offset, &peerRmtAddrs, reinterpret_cast<void*>(cleanupQueue), NULL);
}
if (*regFlag)
*regAddr = (void*)((uintptr_t)peerRmtAddrs + offset);
return ret;
}
static ncclResult_t getCollNetSupport(struct ncclComm* comm, struct ncclTaskColl* task, int* collNetSupport);
static ncclResult_t getAlgoInfo(
struct ncclComm* comm, struct ncclTaskColl* task,
@@ -545,7 +693,7 @@ ncclResult_t ncclPrepareTasks(struct ncclComm* comm, bool* algoNeedConnect, bool
void* regBufSend[NCCL_MAX_LOCAL_RANKS];
void* regBufRecv[NCCL_MAX_LOCAL_RANKS];
bool regNeedConnect = true;
registerIntraNodeBuffers(comm, task, regBufSend, regBufRecv, &planner->collCleanupQueue, &regNeedConnect);
registerCollBuffers(comm, task, regBufSend, regBufRecv, &planner->collCleanupQueue, &regNeedConnect);
if (comm->runtimeConn && comm->initAlgoChannels[task->algorithm] == false) {
if (task->algorithm == NCCL_ALGO_NVLS_TREE && comm->initAlgoChannels[NCCL_ALGO_NVLS] == false && regNeedConnect == true) {
@@ -562,6 +710,10 @@ ncclResult_t ncclPrepareTasks(struct ncclComm* comm, bool* algoNeedConnect, bool
struct ncclDevWorkColl devWork = {};
devWork.sendbuff = (void*)task->sendbuff;
devWork.recvbuff = (void*)task->recvbuff;
devWork.sendbuffOffset = task->sendbuffOffset;
devWork.recvbuffOffset = task->recvbuffOffset;
devWork.sendbuffRmtAddrs = task->sendbuffRmtAddrs;
devWork.recvbuffRmtAddrs = task->recvbuffRmtAddrs;
devWork.root = task->root;
devWork.nWarps = task->nWarps;
devWork.redOpArg = task->opDev.scalarArg;
@@ -574,35 +726,13 @@ ncclResult_t ncclPrepareTasks(struct ncclComm* comm, bool* algoNeedConnect, bool
struct ncclWorkList* workNode;
switch (task->regBufType) {
case NCCL_REGULAR_BUFFER:
case NCCL_IPC_REG_BUFFER:
case NCCL_COLLNET_REG_BUFFER:
{ workNode = ncclMemoryStackAllocInlineArray<ncclWorkList, ncclDevWorkColl>(&comm->memScoped, 1);
workNode->workType = ncclDevWorkTypeColl;
workNode->size = sizeof(struct ncclDevWorkColl);
memcpy((void*)(workNode+1), (void*)&devWork, workNode->size);
} break;
case NCCL_IPC_REG_BUFFER:
{ struct ncclDevWorkCollReg workReg = {};
workReg.coll = devWork;
struct ncclChannel *channel0 = &comm->channels[0];
for (int i=0; i < NCCL_MAX_DIRECT_ARITY; i++) {
int peer = channel0->collnetDirect.down[i];
if (peer == -1) break;
int j = comm->rankToLocalRank[peer]; // Get intra-node slot
workReg.dnInputs[i] = regBufSend[j]; // Input buffer of leaf peer
workReg.dnOutputs[i] = regBufRecv[j]; // Output buffer of leaf peer
}
for (int i=0; i < NCCL_MAX_DIRECT_ARITY; i++) {
int peer = channel0->collnetDirect.up[i];
if (peer == -1) break;
int j = comm->rankToLocalRank[peer];
// Output buffer of root peer
workReg.upOutputs[i] = regBufRecv[j];
}
workNode = ncclMemoryStackAllocInlineArray<ncclWorkList, ncclDevWorkCollReg>(&comm->memScoped, 1);
workNode->workType = ncclDevWorkTypeCollReg;
workNode->size = sizeof(struct ncclDevWorkCollReg);
memcpy((void*)(workNode+1), (void*)&workReg, workNode->size);
} break;
case NCCL_NVLS_REG_BUFFER:
{ struct ncclDevWorkCollReg workReg = {};
workReg.coll = devWork; // C++ struct assignment
@@ -639,6 +769,7 @@ static ncclResult_t scheduleCollTasksToPlan(
int nChannels[2*2] = {0, 0, 0, 0}; // [collnet][nvls]
int const nMaxChannels[2*2] = {comm->nChannels, comm->nvlsChannels, // [collnet][nvls]
comm->nChannels, comm->nvlsChannels};
constexpr size_t MinTrafficPerChannel = 512; // Traffic as minimal
do {
size_t workBytes = 0;
struct ncclTaskColl* task = ncclIntruQueueHead(&planner->collTaskQueue);
@@ -650,7 +781,7 @@ static ncclResult_t scheduleCollTasksToPlan(
nPlanColls += 1;
workBytes += workNode->size;
int kind = 2*task->isCollnet + task->isNvls;
trafficBytes[kind] += task->trafficBytes;
trafficBytes[kind] += std::max(MinTrafficPerChannel, task->trafficBytes);
nChannels[kind] += task->nMaxChannels;
nChannels[kind] = std::min(nChannels[kind], nMaxChannels[kind]);
task = task->next;
@@ -660,7 +791,6 @@ static ncclResult_t scheduleCollTasksToPlan(
} while (0);
int kindPrev = -1;
constexpr size_t MinTrafficPerChannel = 512;
size_t trafficPerChannel = 0;
int channelId = 0;
size_t currentTraffic = 0;
@@ -700,14 +830,16 @@ static ncclResult_t scheduleCollTasksToPlan(
for (int c=devWork->channelLo; c <= (int)devWork->channelHi; c++) {
proxyOp.channelId = c;
proxyOp.opCount = proxyOpId;
proxyOp.task.coll = task;
proxyOp.rank = comm->rank;
addWorkBatchToPlan(comm, plan, c, workNode->workType, task->devFuncId, plan->workBytes);
NCCLCHECK(addProxyOpIfNeeded(comm, plan, &proxyOp));
}
} else { // not task->isCollnet
constexpr size_t cellSize = 16;
int trafficPerByte = ncclFuncTrafficPerByte(task->func, comm->nRanks);
size_t cellSize = divUp(divUp(MinTrafficPerChannel, (size_t)trafficPerByte), 16) * 16;
int elementsPerCell = cellSize/elementSize;
size_t cells = divUp(task->count*elementSize, cellSize);
int trafficPerByte = ncclFuncTrafficPerByte(task->func, comm->nRanks);
size_t trafficPerElement = elementSize*trafficPerByte;
size_t trafficPerCell = cellSize*trafficPerByte;
size_t cellsPerChannel = std::min(cells, divUp(trafficPerChannel, trafficPerCell));
@@ -715,7 +847,7 @@ static ncclResult_t scheduleCollTasksToPlan(
if (channelId+1 == nMaxChannels[kind]) { // On last channel everything goes to "lo"
cellsLo = cells;
} else {
cellsLo = std::min(cells, (trafficPerChannel-currentTraffic)/trafficPerCell);
cellsLo = std::min(cells, divUp((trafficPerChannel-currentTraffic),trafficPerCell));
}
int nMidChannels = (cells-cellsLo)/cellsPerChannel;
size_t cellsHi = (cells-cellsLo)%cellsPerChannel;
@@ -783,12 +915,12 @@ static ncclResult_t scheduleCollTasksToPlan(
// Update the current channel and vacant traffic budget.
if (countHi != 0) {
channelId += nChannels-1;
currentTraffic = countHi*trafficPerElement;
currentTraffic = cellsHi*elementsPerCell*trafficPerElement;
} else if (nMidChannels != 0) {
channelId += nChannels;
currentTraffic = 0;
} else {
currentTraffic += countLo*trafficPerElement;
currentTraffic += cellsLo*elementsPerCell*trafficPerElement;
}
if (currentTraffic >= trafficPerChannel && channelId+1 != nMaxChannels[kind]) {
@@ -808,6 +940,8 @@ static ncclResult_t scheduleCollTasksToPlan(
}
proxyOp->channelId = c;
proxyOp->opCount = proxyOpId;
proxyOp->task.coll = task;
proxyOp->rank = comm->rank;
proxyOp->connIndex = 0;
if (task->protocol == NCCL_PROTO_SIMPLE && task->algorithm == NCCL_ALGO_RING) {
if (comm->useIntraNet && nBytes > rcclParamIntraNetThreshold()) {
@@ -815,6 +949,9 @@ static ncclResult_t scheduleCollTasksToPlan(
}
}
addWorkBatchToPlan(comm, plan, c, workNode->workType, task->devFuncId, plan->workBytes);
// Coverity reports "proxyOp->connection" as being possibly uninitialized. It's hard to
// determine if that's actually true but it's also not clear if that would be an issue.
// coverity[uninit_use_in_call:FALSE]
NCCLCHECK(addProxyOpIfNeeded(comm, plan, proxyOp));
}
}
@@ -859,6 +996,7 @@ static ncclResult_t scheduleCollTasksToPlan(
ncclIntruQueueDequeue(&planner->collWorkQueue);
nPlanColls -= 1;
planner->nTasksColl -= 1;
ncclIntruQueueEnqueue(&plan->collTaskQueue, task);
ncclIntruQueueEnqueue(&plan->workQueue, workNode);
plan->workBytes += workNode->size;
}
@@ -878,7 +1016,8 @@ static ncclResult_t addP2pToPlan(
int nChannelsMin, int nChannelsMax, int p2pRound,
int sendRank, void* sendAddr, ssize_t sendBytes,
int recvRank, void* recvAddr, ssize_t recvBytes,
uint64_t sendOpCount, uint64_t recvOpCount
uint64_t sendOpCount, uint64_t recvOpCount,
struct ncclTaskP2p** p2pTasks
) {
int connIndex[2] = {1, 1};
bool selfSend = (sendRank == comm->rank);
@@ -921,7 +1060,8 @@ static ncclResult_t addP2pToPlan(
int chunkSize[2];
int chunkDataSize[2];
int chunkDataSize_u32fp8[2];
bool registered[2];
bool registered[2] = {false, false};
bool ipcRegistered[2] = {false, false};
for (int dir=0; dir < 2; dir++) { // 0=recv, 1=send
if (bytes[dir] != -1) protoLL[dir] &= bytes[dir] <= thresholdLL;
@@ -945,11 +1085,29 @@ static ncclResult_t addP2pToPlan(
chunkSize[dir] = chunkDataSize[dir];
if (protocol[dir] == NCCL_PROTO_LL) chunkSize[dir] *= 2;
registered[dir] = false;
if (bytes[dir] > 0 && network[dir] && proxySameProcess[dir] && protocol[dir] == NCCL_PROTO_SIMPLE) {
struct ncclReg* regRecord;
NCCLCHECK(ncclRegFind(comm, addrs[dir], bytes[dir], &regRecord));
registered[dir] = (regRecord && regRecord->nDevs);
if (network[dir]) {
if (bytes[dir] > 0 && proxySameProcess[dir] && protocol[dir] == NCCL_PROTO_SIMPLE) {
struct ncclReg* regRecord;
NCCLCHECK(ncclRegFind(comm, addrs[dir], bytes[dir], &regRecord));
registered[dir] = regRecord && regRecord->nDevs;
}
} else if (bytes[dir] > 0 && addrs[dir] && protocol[dir] == NCCL_PROTO_SIMPLE && !selfSend) {
int peerRank = dir ? sendRank : recvRank;
int regFlag = 0;
int channelId = ncclP2pChannelForPart(comm->p2pnChannels, base, 0, nChannelsMax, comm->nNodes);
struct ncclChannelPeer** channelPeers = comm->channels[channelId].peers;
struct ncclConnector* conn = dir ? &channelPeers[peerRank]->send[connIndex[dir]]
: &channelPeers[peerRank]->recv[connIndex[dir]];
void* regAddr = NULL;
if (conn->conn.flags & (NCCL_IPC_WRITE | NCCL_IPC_READ | NCCL_DIRECT_WRITE | NCCL_DIRECT_READ)) {
// We require users registering buffers on both sides
NCCLCHECK(registerP2pBuffer(comm, addrs[dir], peerRank, bytes[dir], &regFlag, &regAddr, &plan->cleanupQueue));
if (regFlag) {
if (dir == 0 && conn->conn.flags & (NCCL_IPC_WRITE | NCCL_DIRECT_WRITE)) recvAddr = regAddr;
else if (dir == 1 && conn->conn.flags & (NCCL_IPC_READ | NCCL_DIRECT_READ)) sendAddr = regAddr;
}
}
ipcRegistered[dir] = regFlag ? true : false;
}
if (bytes[dir] == -1) nChannels[dir] = 0;
@@ -979,6 +1137,7 @@ static ncclResult_t addP2pToPlan(
work->nSendChannels = nChannels[1];
work->sendProtoLL = protoLL[1];
work->sendRegistered = registered[1];
work->sendIpcReg = ipcRegistered[1];
work->sendChunkSize_u32fp8 = chunkDataSize_u32fp8[1];
work->sendRank = sendRank;
work->sendAddr = sendAddr;
@@ -988,6 +1147,7 @@ static ncclResult_t addP2pToPlan(
work->nRecvChannels = nChannels[0];
work->recvProtoLL = protoLL[0];
work->recvRegistered = registered[0];
work->recvIpcReg = ipcRegistered[0];
work->recvChunkSize_u32fp8 = chunkDataSize_u32fp8[0];
work->recvRank = recvRank;
work->recvAddr = recvAddr;
@@ -1008,6 +1168,9 @@ static ncclResult_t addP2pToPlan(
op->pattern = dir ? ncclPatternSend : ncclPatternRecv;
op->chunkSize = chunkSize[dir];
op->reg = registered[dir];
op->coll = p2pTasks[dir] ? p2pTasks[dir]->func : 0;
op->task.p2p = p2pTasks[dir];
op->rank = comm->rank;
op->connIndex = connIndex[dir];
// The following are modified per channel part in addWorkToChannels():
// op->buffer, op->nbytes, op->nsteps = ...;
@@ -1130,14 +1293,16 @@ static ncclResult_t scheduleP2pTasksToPlan(
if (!testBudget(budget, plan->nWorkBatches+nChannelsMax, plan->workBytes + sizeof(struct ncclDevWorkP2p))) {
return ncclSuccess;
}
NCCLCHECK(addP2pToPlan(comm, plan, nChannelsMin, nChannelsMax, round, sendRank, sendBuff, sendBytes, recvRank, recvBuff, recvBytes,
send ? send->opCount : 0, recv ? recv->opCount : 0));
struct ncclTaskP2p* p2pTasks[2] = { recv, send };
NCCLCHECK(addP2pToPlan(comm, plan, nChannelsMin, nChannelsMax, round, sendRank, sendBuff, sendBytes, recvRank, recvBuff, recvBytes, send ? send->opCount : 0, recv ? recv->opCount : 0, p2pTasks));
if (send != nullptr) {
ncclIntruQueueDequeue(&peers[sendRank].sendQueue);
ncclIntruQueueEnqueue(&plan->p2pTaskQueue, send);
comm->planner.nTasksP2p -= 1;
}
if (recv != nullptr) {
ncclIntruQueueDequeue(&peers[recvRank].recvQueue);
ncclIntruQueueEnqueue(&plan->p2pTaskQueue, recv);
comm->planner.nTasksP2p -= 1;
}
}
@@ -1190,29 +1355,43 @@ static void waitWorkFifoAvailable(struct ncclComm* comm, uint32_t desiredProduce
}
}
namespace {
struct uploadWork_cleanup_t {
struct ncclCommEventCallback base;
void *hostBuf;
};
ncclResult_t uploadWork_cleanup_fn(
struct ncclComm* comm, struct ncclCommEventCallback* cb
) {
struct uploadWork_cleanup_t* me = (struct uploadWork_cleanup_t*)cb;
free(me->hostBuf);
CUDACHECK(cudaEventDestroy(me->base.event));
return ncclSuccess;
}
}
static ncclResult_t uploadWork(struct ncclComm* comm, struct ncclKernelPlan* plan) {
size_t workBytes = plan->workBytes;
size_t batchBytes = plan->nWorkBatches*sizeof(struct ncclDevWorkBatch);
void* fifoBuf;
void* fifoBufHost;
uint32_t fifoCursor, fifoMask;
switch (plan->workStorageType) {
case ncclDevWorkStorageTypeArgs:
plan->kernelArgs->workBuf = nullptr;
fifoBuf = (void*)plan->kernelArgs;
fifoBufHost = (void*)plan->kernelArgs;
fifoCursor = sizeof(ncclDevKernelArgs) + batchBytes;
fifoMask = ~0u;
break;
case ncclDevWorkStorageTypeFifo:
fifoBuf = comm->workFifoBuf;
fifoBufHost = comm->workFifoBuf;
fifoCursor = comm->workFifoProduced;
fifoMask = comm->workFifoBytes-1;
waitWorkFifoAvailable(comm, fifoCursor + workBytes);
plan->kernelArgs->workBuf = comm->workFifoBufDev;
break;
case ncclDevWorkStorageTypePersistent:
ncclMemoryStackPush(&comm->memScoped);
fifoBuf = ncclMemoryStackAlloc(&comm->memScoped, workBytes, /*align=*/16);
fifoBufHost = aligned_alloc(16, workBytes); // We rely on 16-byte alignment
fifoCursor = 0;
fifoMask = ~0u;
break;
@@ -1234,7 +1413,7 @@ static ncclResult_t uploadWork(struct ncclComm* comm, struct ncclKernelPlan* pla
// Write the channel-shared work structs.
struct ncclWorkList* workNode = ncclIntruQueueHead(&plan->workQueue);
while (workNode != nullptr) {
char* dst = (char*)fifoBuf;
char* dst = (char*)fifoBufHost;
char* src = (char*)(workNode+1);
for (int n = workNode->size; n != 0; n -= 16) {
memcpy(
@@ -1254,11 +1433,39 @@ static ncclResult_t uploadWork(struct ncclComm* comm, struct ncclKernelPlan* pla
if (comm->workFifoBufGdrHandle != nullptr) wc_store_fence();
break;
case ncclDevWorkStorageTypePersistent:
NCCLCHECK(ncclCudaMalloc(&plan->workBufPersistent, workBytes));
plan->kernelArgs->workBuf = plan->workBufPersistent;
NCCLCHECK(ncclCudaMemcpy(plan->workBufPersistent, fifoBuf, workBytes));
ncclMemoryStackPop(&comm->memScoped);
break;
{ ncclResult_t result = ncclSuccess;
cudaStreamCaptureMode mode = cudaStreamCaptureModeRelaxed;
void* fifoBufDev = nullptr;
CUDACHECK(cudaThreadExchangeStreamCaptureMode(&mode));
// Acquire deviceStream to gain access to deviceStream.cudaStream. Since the
// user's graph will be launched later, and it also acquires the deviceStream,
// it will observe this upload.
NCCLCHECKGOTO(ncclStrongStreamAcquireUncaptured(&comm->sharedRes->deviceStream), result, finish_scope);
CUDACHECKGOTO(cudaMallocAsync(&fifoBufDev, workBytes, comm->memPool, comm->sharedRes->deviceStream.cudaStream), result, finish_scope);
plan->workBufPersistent = fifoBufDev;
plan->kernelArgs->workBuf = fifoBufDev;
CUDACHECKGOTO(cudaMemcpyAsync(fifoBufDev, fifoBufHost, workBytes, cudaMemcpyDefault, comm->sharedRes->deviceStream.cudaStream), result, finish_scope);
cudaEvent_t memcpyDone;
CUDACHECKGOTO(cudaEventCreateWithFlags(&memcpyDone, cudaEventDisableTiming), result, finish_scope);
CUDACHECKGOTO(cudaEventRecord(memcpyDone, comm->sharedRes->deviceStream.cudaStream), result, finish_scope);
struct uploadWork_cleanup_t* cleanup;
NCCLCHECK(ncclCalloc(&cleanup, 1));
cleanup->base.fn = uploadWork_cleanup_fn;
cleanup->base.event = memcpyDone;
cleanup->hostBuf = fifoBufHost;
ncclIntruQueueEnqueue(&comm->eventCallbackQueue, &cleanup->base);
NCCLCHECKGOTO(ncclStrongStreamRelease(ncclCudaGraphNone(), &comm->sharedRes->deviceStream), result, finish_scope);
NCCLCHECKGOTO(ncclCommPollEventCallbacks(comm), result, finish_scope);
finish_scope:
CUDACHECK(cudaThreadExchangeStreamCaptureMode(&mode));
if (result != ncclSuccess) return result;
} break;
default: break;
}
return ncclSuccess;
@@ -1272,6 +1479,11 @@ static ncclResult_t uploadProxyOps(struct ncclComm* comm, struct ncclKernelPlan*
struct ncclProxyOp* op = ncclIntruQueueHead(&plan->proxyOpQueue);
while (op != nullptr) {
op->profilerContext = comm->profilerContext;
op->eActivationMask = op->coll <= ncclFuncAllReduce ? op->task.coll->eActivationMask : op->task.p2p->eActivationMask;
op->taskEventHandle = op->coll <= ncclFuncAllReduce ? op->task.coll->eventHandle : op->task.p2p->eventHandle;
ncclProfilerAddPidToProxyOp(op);
uint64_t oldId = op->opCount;
// Ignoring the bottom tag bit, opCount's are zero-based within plan so
// translate them to the tip of the comm's history.
@@ -1306,8 +1518,12 @@ static ncclResult_t uploadProxyOps(struct ncclComm* comm, struct ncclKernelPlan*
}
static ncclResult_t hostStreamPlanTask(struct ncclComm* comm, struct ncclKernelPlan* plan) {
NCCLCHECK(ncclProfilerStartGroupEvent(plan));
NCCLCHECK(ncclProfilerStartTaskEvents(plan));
NCCLCHECK(uploadProxyOps(comm, plan));
NCCLCHECK(ncclProxyStart(comm));
NCCLCHECK(ncclProfilerStopTaskEvents(plan));
NCCLCHECK(ncclProfilerStopGroupEvent(plan));
if (!plan->persistent) {
// Notify main thread of our reclaiming. This will reclaim plan concurrently.
ncclIntruQueueMpscEnqueue(&comm->callbackQueue, &plan->reclaimer);
@@ -1376,7 +1592,7 @@ ncclResult_t ncclLaunchPrepare(struct ncclComm* comm) {
plan->comm = comm;
plan->reclaimer.fn = reclaimPlan;
plan->persistent = persistent;
// uploadWork() promotes ncclDevWorkStorageType[Fifo|Buf]->Args if the work can fit.
// finishPlan() promotes ncclDevWorkStorageType[Fifo|Persistent]->Args if the work can fit.
plan->workStorageType = persistent ? ncclDevWorkStorageTypePersistent
: ncclDevWorkStorageTypeFifo;
@@ -1655,10 +1871,15 @@ static ncclResult_t updateCollCostTable(
for (int a=0; a<NCCL_NUM_ALGORITHMS; a++) {
if ((a == NCCL_ALGO_COLLNET_DIRECT || a == NCCL_ALGO_COLLNET_CHAIN) && collNetSupport != 1) continue;
// CollNetDirect is only supported for up to 8 local GPUs
if (a == NCCL_ALGO_COLLNET_DIRECT && comm->maxLocalRanks > NCCL_MAX_DIRECT_ARITY+1) continue;
if ((a == NCCL_ALGO_NVLS || a == NCCL_ALGO_NVLS_TREE) && nvlsSupport != 1 && info->func != ncclFuncAllGather) continue;
if (a == NCCL_ALGO_NVLS && collNetSupport != 1 && comm->nNodes > 1) continue;
/* now we only support single-node NVLS allgather and reducescatter */
if (a == NCCL_ALGO_NVLS && (info->func == ncclFuncAllGather || info->func == ncclFuncReduceScatter) && comm->nNodes > 1) continue;
/* Tree reduceScatter doesn't support scaling yet */
if (a == NCCL_ALGO_PAT && info->func == ncclFuncReduceScatter
&& (info->opDev.op == ncclDevPreMulSum || info->opDev.op == ncclDevSumPostDiv)) continue;
for (int p=0; p<NCCL_NUM_PROTOCOLS; p++) {
if (p == NCCL_PROTO_LL128 && !(comm->topo->type & RCCL_TOPO_XGMI_ALL)) {
table[a][p] = NCCL_ALGO_PROTO_IGNORE;
@@ -1714,6 +1935,8 @@ static ncclResult_t topoGetAlgoInfo(
info->protocol = protocol;
float time = minTime;
// Yes, we are first assigning and then testing if protocol is sane, but that's OK in this case.
// coverity[check_after_sink]
if (info->algorithm == NCCL_ALGO_UNDEF || info->protocol == NCCL_PROTO_UNDEF) {
if (backupAlgo == NCCL_ALGO_UNDEF || backupProto == NCCL_PROTO_UNDEF) {
WARN("Error : no algorithm/protocol available");
@@ -1749,7 +1972,7 @@ static ncclResult_t topoGetAlgoInfo(
#endif
}
#endif
if (comm->rank == 0) INFO(NCCL_TUNING, "%ld Bytes -> Algo %d proto %d time %f", nBytes, info->algorithm, info->protocol, time);
if (comm->rank == 0) INFO(NCCL_TUNING, "%s: %ld Bytes -> Algo %d proto %d time %f", ncclFuncToString(info->func), nBytes, info->algorithm, info->protocol, time);
if (simInfo) simInfo->estimatedTime = time;
TRACE(NCCL_COLL, "%ld Bytes -> Algo %d proto %d time %f", nBytes, info->algorithm, info->protocol, time);
@@ -1822,6 +2045,7 @@ static ncclResult_t topoGetAlgoInfo(
info->nMaxChannels = nc;
}
if (info->algorithm == NCCL_ALGO_TREE) nt = NCCL_MAX_NTHREADS; // Tree now uses all threads always.
if (info->algorithm == NCCL_ALGO_PAT) nt = NCCL_MAX_NTHREADS;
info->nWarps = nt/WARP_SIZE;
return ncclSuccess;
}
@@ -1872,8 +2096,15 @@ static ncclResult_t calcCollChunking(
pattern = info->algorithm == NCCL_ALGO_TREE ? ncclPatternTreeUp : ncclPatternPipelineTo;
break;
case ncclFuncReduceScatter:
pattern =
info->algorithm == NCCL_ALGO_PAT ? ncclPatternPatUp :
info->algorithm == NCCL_ALGO_NVLS ? ncclPatternNvls :
info->algorithm == NCCL_ALGO_COLLNET_DIRECT ? ncclPatternCollnetDirect :
ncclPatternRing;
break;
case ncclFuncAllGather:
pattern =
info->algorithm == NCCL_ALGO_PAT ? ncclPatternPatDown :
info->algorithm == NCCL_ALGO_NVLS ? ncclPatternNvls :
info->algorithm == NCCL_ALGO_COLLNET_DIRECT ? ncclPatternCollnetDirect :
ncclPatternRing;
@@ -1900,6 +2131,8 @@ static ncclResult_t calcCollChunking(
case ncclPatternTreeUp:
case ncclPatternTreeDown:
case ncclPatternTreeUpDown:
case ncclPatternPatUp:
case ncclPatternPatDown:
case ncclPatternPipelineFrom:
case ncclPatternPipelineTo:
case ncclPatternCollnetChain:
@@ -1962,13 +2195,17 @@ static ncclResult_t calcCollChunking(
int maxChunkSize = comm->nvlsChunkSize;
if (comm->nNodes > 1 && comm->bandwidths[ncclFuncAllReduce][NCCL_ALGO_NVLS][NCCL_PROTO_SIMPLE] < 150) maxChunkSize = 32768;
if (chunkSize > maxChunkSize) chunkSize = maxChunkSize;
// Use uint64_t so that concurrentOps*chunkSize*X does not overflow
// Use uint64_t so that concurrentOps*chunkSize*X does not overflow.
// However, nChannels * comm->channels[0].nvls.nHeads should easily fit in 32 bits.
// coverity[overflow_before_widen]
uint64_t concurrentOps = nChannels * comm->channels[0].nvls.nHeads;
if ((nBytes < (64 * (concurrentOps * chunkSize))) && (chunkSize > 65536)) chunkSize = 65536;
if ((nBytes < (8 * (concurrentOps * chunkSize))) && (chunkSize > 32768)) chunkSize = 32768;
if ((nBytes < (2 * (concurrentOps * chunkSize))) && (chunkSize > 16384)) chunkSize = 16384;
} else if (info->algorithm == NCCL_ALGO_NVLS_TREE) {
// Use uint64_t so that concurrentOps*chunkSize*X does not overflow
// Use uint64_t so that concurrentOps*chunkSize*X does not overflow.
// However, nChannels * comm->channels[0].nvls.nHeads should easily fit in 32 bits.
// coverity[overflow_before_widen]
uint64_t concurrentOps = nChannels * comm->channels[0].nvls.nHeads;
chunkSize = comm->nvlsChunkSize;
int maxChunkSize = (int)ncclParamNvlsTreeMaxChunkSize();
@@ -1982,14 +2219,21 @@ static ncclResult_t calcCollChunking(
int nNodes = comm->nNodes;
float ppn = comm->nRanks / (float)nNodes;
float nstepsLL128 = 1+log2i(nNodes) + 0.1*ppn;
// Yes, we are OK with the division on the left side of the < operand being integer.
// coverity[integer_division]
while (nBytes / (nChannels*chunkSize) < nstepsLL128*64/ppn && chunkSize > 131072) chunkSize /= 2;
// coverity[integer_division]
while (nBytes / (nChannels*chunkSize) < nstepsLL128*16/ppn && chunkSize > 32768) chunkSize /= 2;
} else if (info->func == ncclFuncAllGather && info->algorithm == NCCL_ALGO_PAT) {
while (chunkSize*nChannels*32 > nBytes && chunkSize > 65536) chunkSize /= 2;
} else if (info->func == ncclFuncReduceScatter && info->algorithm == NCCL_ALGO_PAT) {
while (chunkSize*nChannels*16 > nBytes && chunkSize > 65536) chunkSize /= 2;
}
// Compute directFlags of work struct.
if (info->algorithm == NCCL_ALGO_COLLNET_DIRECT) {
// Set direct direction for broadcast-gather (read or write)
*outDirectFlags = (nBytes/nChannels <= 1024*1024) ? NCCL_DIRECT_WRITE : NCCL_DIRECT_READ;
*outDirectFlags = (nBytes/nChannels <= 1024 * 4) ? NCCL_DIRECT_READ : NCCL_DIRECT_WRITE;
} else {
*outDirectFlags = 0;
}
@@ -2038,6 +2282,10 @@ static ncclResult_t calcCollChunking(
}
}
if (pattern == ncclPatternPatUp || pattern == ncclPatternPatDown) {
proxyOp->nbytes = DIVUP(nBytes, nChannels);
}
*outChunkSize = chunkSize;
return ncclSuccess;
}
@@ -2069,6 +2317,7 @@ static ncclResult_t hostToDevRedOp(
opFull->proxyOp = op;
int nbits = 8*ncclTypeSize(datatype);
if (nbits <= 0) return ncclInvalidArgument;
uint64_t allBits = uint64_t(-1)>>(64-nbits);
uint64_t signBit = allBits^(allBits>>1);
@@ -2154,6 +2403,9 @@ static ncclResult_t taskAppend(struct ncclComm* comm, struct ncclInfo* info) {
ncclGroupCommJoin(info->comm);
struct ncclTaskP2p* p2p = ncclMemoryStackAlloc<struct ncclTaskP2p>(&comm->memScoped);
p2p->buff = (void*)info->recvbuff;
p2p->count = info->count;
p2p->datatype = info->datatype;
p2p->root = info->root;
p2p->bytes = nBytes;
p2p->opCount = comm->opCount;
ncclIntruQueueEnqueue(
@@ -2245,7 +2497,7 @@ static ncclResult_t taskAppend(struct ncclComm* comm, struct ncclInfo* info) {
while (true) {
if (l == nullptr) { // Got to the end, this must be a new stream.
struct ncclCudaGraph graph;
NCCLCHECK(ncclCudaGetCapturingGraph(&graph, info->stream))
NCCLCHECK(ncclCudaGetCapturingGraph(&graph, info->stream));
if (planner->streams != nullptr && !ncclCudaGraphSame(planner->capturingGraph, graph)) {
WARN("Streams given to a communicator within a NCCL group must either be all uncaptured or all captured by the same graph.");
return ncclInvalidUsage;
@@ -2297,7 +2549,7 @@ exit:
NCCLCHECK(ncclGroupEndInternal());
/* if depth is 1, ncclGroupEndInternal() will trigger group ops. The state can change
* so we have to check state here. */
if (info->comm && !info->comm->config.blocking) { NCCLCHECK(ncclCommGetAsyncError(info->comm, &ret)) };
if (info->comm && !info->comm->config.blocking) { NCCLCHECK(ncclCommGetAsyncError(info->comm, &ret)); }
return ret;
fail:
if (info->comm && !info->comm->config.blocking) (void) ncclCommSetAsyncError(info->comm, ret);
@@ -2315,7 +2567,8 @@ ncclResult_t ncclRedOpCreatePreMulSum_impl(ncclRedOp_t *op, void *scalar, ncclDa
int cap = 2*comm->userRedOpCapacity;
if (cap < 4) cap = 4;
ncclUserRedOp *ops = new ncclUserRedOp[cap];
std::memcpy(ops, comm->userRedOps, comm->userRedOpCapacity*sizeof(ncclUserRedOp));
if (comm->userRedOpCapacity > 0)
std::memcpy(ops, comm->userRedOps, comm->userRedOpCapacity*sizeof(ncclUserRedOp));
for(int ix=comm->userRedOpCapacity; ix < cap; ix++)
ops[ix].freeNext = ix + 1;
delete[] comm->userRedOps;
@@ -2331,8 +2584,10 @@ ncclResult_t ncclRedOpCreatePreMulSum_impl(ncclRedOp_t *op, void *scalar, ncclDa
user->datatype = datatype;
user->opFull.op = ncclDevPreMulSum;
if (residence == ncclScalarHostImmediate) {
int size = ncclTypeSize(datatype);
if (size < 1) return ncclInternalError;
user->opFull.scalarArgIsPtr = false;
std::memcpy(&user->opFull.scalarArg, scalar, ncclTypeSize(datatype));
std::memcpy(&user->opFull.scalarArg, scalar, size);
} else {
user->opFull.scalarArgIsPtr = true;
user->opFull.scalarArg = reinterpret_cast<uint64_t>(scalar);
@@ -2349,6 +2604,10 @@ ncclResult_t ncclRedOpDestroy_impl(ncclRedOp_t op, ncclComm_t comm) {
WARN("ncclRedOpDestroy : operator is a NCCL builtin.");
return ncclInvalidArgument;
}
// int(ncclMaxRedOp) < int(op) will always be false due to the sizes of
// the datatypes involved, and that's by design. We keep the check though
// just as a reminder.
// coverity[result_independent_of_operands]
if (int(op) < 0 || int(ncclMaxRedOp) < int(op)) {
WARN("ncclRedOpDestroy : operator is garbage.");
return ncclInvalidArgument;