[CCF Archive] Store object type eviction policy submission #3

Closed
kancel wants to merge 382 commits from kancel:ccf-archive-pr2746 into main
2 changed files with 364 additions and 75 deletions
Showing only changes of commit 208d105e6c - Show all commits

View File

@ -33,6 +33,17 @@
#include "transfer_metadata.h"
#include "transport/transport.h"
static bool checkCudaErrorReturn(cudaError_t result, const char *message) {
if (result != cudaSuccess) {
LOG(ERROR) << message << " (Error code: " << result << " - "
<< cudaGetErrorString(result) << ")" << std::endl;
return false;
}
return true;
}
namespace mooncake {
namespace {
struct CudaStreamNVLinkRAII {
cudaStream_t stream_;
@ -67,24 +78,17 @@ struct CudaEventNVLinkRAII {
static thread_local CudaEventNVLinkRAII tl_nvlink_sync_event;
} // anonymous namespace
static bool checkCudaErrorReturn(cudaError_t result, const char *message) {
if (result != cudaSuccess) {
LOG(ERROR) << message << " (Error code: " << result << " - "
<< cudaGetErrorString(result) << ")" << std::endl;
return false;
}
return true;
}
namespace mooncake {
typedef Transport::Slice Slice;
using Slice = Transport::Slice;
/// Submit batched memcpy operations using cudaMemcpyBatchAsync when available
/// (CUDA 12.8+), falling back to per-slice cudaMemcpyAsync otherwise.
/// Uses cudaMemcpySrcAccessOrderStream attribute to respect stream access
/// ordering semantics. Individual slice errors are tracked so that slices
/// whose memcpy failed are marked as FAILED while successful ones are POSTED.
/// Uses cudaMemcpySrcAccessOrderStream to ensure source data visibility for
/// P2P copies. This attribute is REQUIRED — without it, the GPU does not
/// insert the necessary memory barriers for P2P access, causing segfaults.
/// The caller must also establish GPU-level stream synchronization
/// (via cudaEventRecord + cudaStreamWaitEvent) before calling this function.
/// Individual slice errors are tracked so that slices whose memcpy failed
/// are marked as FAILED while successfully submitted ones are POSTED.
static void submitBatchMemcpy(const std::vector<Slice *> &slices,
const std::vector<void *> &srcs,
const std::vector<void *> &dsts,
@ -113,15 +117,18 @@ static void submitBatchMemcpy(const std::vector<Slice *> &slices,
(void)logged_once;
#if CUDART_VERSION >= 12080
// Use srcAccessOrderStream for P2P copies. The GPU-level dependency
// established by cudaEventRecord + cudaStreamWaitEvent in the caller
// ensures the source data is visible through nvlink_stream's access
// order, satisfying srcAccessOrderStream's requirement.
// srcAccessOrderStream is REQUIRED for P2P copies — without it, the GPU
// does not insert necessary memory barriers for cross-device access,
// resulting in segmentation faults. The caller also establishes a
// GPU-level dependency via cudaEventRecord(cudaStreamPerThread) +
// cudaStreamWaitEvent(nvlink_stream) to ensure source data is coherent
// before the memcpy starts; these two mechanisms are complementary.
cudaMemcpyAttributes attr{};
attr.srcAccessOrder = cudaMemcpySrcAccessOrderStream;
size_t attrs_idx = 0;
// cudaMemcpyBatchAsync in CUDA 12.8 takes non-const size_t* for sizes
std::vector<size_t> mutable_sizes(sizes);
size_t fail_idx = count;
#endif
#if CUDART_VERSION >= 13000
@ -129,26 +136,56 @@ static void submitBatchMemcpy(const std::vector<Slice *> &slices,
const_cast<const void **>(srcs.data()),
mutable_sizes.data(), static_cast<size_t>(count),
&attr, &attrs_idx, 1, stream);
#elif CUDART_VERSION >= 12080
{
size_t fail_idx = count;
err = cudaMemcpyBatchAsync(
const_cast<void **>(dsts.data()), const_cast<void **>(srcs.data()),
mutable_sizes.data(), static_cast<size_t>(count), &attr, &attrs_idx,
1, &fail_idx, stream);
if (err != cudaSuccess) {
if (fail_idx < count) {
LOG(ERROR) << "IntraNodeNvlinkTransport: cudaMemcpyBatchAsync "
<< "failed at index " << fail_idx
<< " (src=" << srcs[fail_idx]
<< ", dst=" << dsts[fail_idx]
<< ", size=" << sizes[fail_idx]
<< "): " << cudaGetErrorString(err);
} else {
LOG(ERROR) << "IntraNodeNvlinkTransport: cudaMemcpyBatchAsync "
<< "failed: " << cudaGetErrorString(err);
if (err != cudaSuccess) {
LOG(ERROR) << "IntraNodeNvlinkTransport: cudaMemcpyBatchAsync "
<< "failed: " << cudaGetErrorString(err);
// CUDA >= 13.0 does not return fail_idx; conservatively mark all
// as FAILED since we cannot determine which copies succeeded.
for (size_t i = 0; i < count; ++i) {
if (slices[i]->status == Slice::PENDING) {
slices[i]->markFailed();
}
}
} else {
for (size_t i = 0; i < count; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
}
#elif CUDART_VERSION >= 12080
err = cudaMemcpyBatchAsync(const_cast<void **>(dsts.data()),
const_cast<void **>(srcs.data()),
mutable_sizes.data(), static_cast<size_t>(count),
&attr, &attrs_idx, 1, &fail_idx, stream);
if (err != cudaSuccess) {
if (fail_idx < count) {
LOG(ERROR) << "IntraNodeNvlinkTransport: cudaMemcpyBatchAsync "
<< "failed at index " << fail_idx
<< " (src=" << srcs[fail_idx]
<< ", dst=" << dsts[fail_idx]
<< ", size=" << sizes[fail_idx]
<< "): " << cudaGetErrorString(err);
} else {
LOG(ERROR) << "IntraNodeNvlinkTransport: cudaMemcpyBatchAsync "
<< "failed: " << cudaGetErrorString(err);
}
// Copies [0, fail_idx) were submitted successfully → POSTED.
// Copy [fail_idx] failed → FAILED.
// Copies (fail_idx, count) were never submitted → FAILED.
for (size_t i = 0; i < fail_idx; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
for (size_t i = fail_idx; i < count; ++i) {
if (slices[i]->status == Slice::PENDING) {
slices[i]->markFailed();
}
}
} else {
for (size_t i = 0; i < count; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
}
#else
// Fallback for CUDA < 12.8: submit each memcpy individually
@ -167,20 +204,6 @@ static void submitBatchMemcpy(const std::vector<Slice *> &slices,
}
return; // Slice states already set above
#endif
// For cudaMemcpyBatchAsync paths, update slice states based on result
if (err != cudaSuccess) {
for (size_t i = 0; i < count; ++i) {
if (slices[i]->status == Slice::PENDING) {
slices[i]->markFailed();
}
}
} else {
for (size_t i = 0; i < count; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
}
}
static int getNumDevices() {
@ -404,11 +427,21 @@ Status IntraNodeNvlinkTransport::getTransferStatus(BatchID batch_id,
std::to_string(batch_id));
}
auto &task = batch_desc.task_list[task_id];
// Poll POSTED slices for async completion via cudaStreamQuery
// Poll POSTED slices for async completion via cudaStreamQuery.
// Cache the query result per stream to avoid redundant driver calls,
// since multiple slices typically share the same CUDA stream.
std::unordered_map<cudaStream_t, cudaError_t> stream_status_cache;
for (auto *slice : task.slice_list) {
if (slice && slice->status == Slice::POSTED) {
cudaStream_t stream = (cudaStream_t)slice->local.cuda_stream;
cudaError_t cuda_err = cudaStreamQuery(stream);
auto it = stream_status_cache.find(stream);
cudaError_t cuda_err;
if (it == stream_status_cache.end()) {
cuda_err = cudaStreamQuery(stream);
stream_status_cache[stream] = cuda_err;
} else {
cuda_err = it->second;
}
if (cuda_err == cudaSuccess) {
slice->markSuccess();
} else if (cuda_err != cudaErrorNotReady) {

View File

@ -24,6 +24,7 @@
#include <cstdint>
#include <iomanip>
#include <memory>
#include <vector>
#include "common.h"
#include "common/serialization.h"
@ -42,6 +43,168 @@ static bool checkCudaErrorReturn(cudaError_t result, const char *message) {
}
namespace mooncake {
namespace {
struct CudaStreamNVLinkRAII {
cudaStream_t stream_;
CudaStreamNVLinkRAII() : stream_(nullptr) {
auto err = cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking);
if (err != cudaSuccess) {
LOG(FATAL) << "Failed to create NVLink CUDA stream: " << err
<< " - " << cudaGetErrorString(err);
}
}
~CudaStreamNVLinkRAII() { cudaStreamDestroy(stream_); }
};
static thread_local CudaStreamNVLinkRAII tl_nvlink_stream;
/// Thread-local CUDA event for GPU-level stream synchronization.
/// Used to establish a GPU-visible dependency between cudaStreamPerThread
/// and nvlink_stream, which is required by cudaMemcpyBatchAsync's
/// srcAccessOrderStream attribute for cross-stream P2P copies.
struct CudaEventNVLinkRAII {
cudaEvent_t event_;
CudaEventNVLinkRAII() {
auto err = cudaEventCreateWithFlags(&event_, cudaEventDisableTiming);
if (err != cudaSuccess) {
LOG(FATAL) << "Failed to create NVLink CUDA sync event: " << err
<< " - " << cudaGetErrorString(err);
}
}
~CudaEventNVLinkRAII() { cudaEventDestroy(event_); }
};
static thread_local CudaEventNVLinkRAII tl_nvlink_sync_event;
} // anonymous namespace
using Slice = Transport::Slice;
/// Submit batched memcpy operations using cudaMemcpyBatchAsync when available
/// (CUDA 12.8+), falling back to per-slice cudaMemcpyAsync otherwise.
/// Uses cudaMemcpySrcAccessOrderStream to ensure source data visibility for
/// P2P copies. This attribute is REQUIRED — without it, the GPU does not
/// insert the necessary memory barriers for P2P access, causing segfaults.
/// The caller must also establish GPU-level stream synchronization
/// (via cudaEventRecord + cudaStreamWaitEvent) before calling this function.
/// Individual slice errors are tracked so that slices whose memcpy failed
/// are marked as FAILED while successfully submitted ones are POSTED.
static void submitBatchMemcpy(const std::vector<Slice *> &slices,
const std::vector<void *> &srcs,
const std::vector<void *> &dsts,
const std::vector<size_t> &sizes,
cudaStream_t stream) {
if (slices.empty()) return;
const size_t count = slices.size();
cudaError_t err = cudaSuccess;
// Log the active memcpy path once per process lifetime
static const bool logged_once = [] {
#if CUDART_VERSION >= 13000
LOG(INFO) << "NvlinkTransport: using cudaMemcpyBatchAsync "
<< "(CUDA >= 13.0 path)";
#elif CUDART_VERSION >= 12080
LOG(INFO) << "NvlinkTransport: using cudaMemcpyBatchAsync "
<< "(CUDA >= 12.8 path)";
#else
LOG(INFO) << "NvlinkTransport: using per-slice cudaMemcpyAsync "
<< "(CUDA < 12.8 fallback path)";
#endif
return true;
}();
(void)logged_once;
#if CUDART_VERSION >= 12080
// srcAccessOrderStream is REQUIRED for P2P copies — without it, the GPU
// does not insert necessary memory barriers for cross-device access,
// resulting in segmentation faults. The caller also establishes a
// GPU-level dependency via cudaEventRecord(cudaStreamPerThread) +
// cudaStreamWaitEvent(nvlink_stream) to ensure source data is coherent
// before the memcpy starts; these two mechanisms are complementary.
cudaMemcpyAttributes attr{};
attr.srcAccessOrder = cudaMemcpySrcAccessOrderStream;
size_t attrs_idx = 0;
// cudaMemcpyBatchAsync in CUDA 12.8 takes non-const size_t* for sizes
std::vector<size_t> mutable_sizes(sizes);
size_t fail_idx = count;
#endif
#if CUDART_VERSION >= 13000
err = cudaMemcpyBatchAsync(const_cast<const void **>(dsts.data()),
const_cast<const void **>(srcs.data()),
mutable_sizes.data(), static_cast<size_t>(count),
&attr, &attrs_idx, 1, stream);
if (err != cudaSuccess) {
LOG(ERROR) << "NvlinkTransport: cudaMemcpyBatchAsync "
<< "failed: " << cudaGetErrorString(err);
// CUDA >= 13.0 does not return fail_idx; conservatively mark all
// as FAILED since we cannot determine which copies succeeded.
for (size_t i = 0; i < count; ++i) {
if (slices[i]->status == Slice::PENDING) {
slices[i]->markFailed();
}
}
} else {
for (size_t i = 0; i < count; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
}
#elif CUDART_VERSION >= 12080
err = cudaMemcpyBatchAsync(const_cast<void **>(dsts.data()),
const_cast<void **>(srcs.data()),
mutable_sizes.data(), static_cast<size_t>(count),
&attr, &attrs_idx, 1, &fail_idx, stream);
if (err != cudaSuccess) {
if (fail_idx < count) {
LOG(ERROR) << "NvlinkTransport: cudaMemcpyBatchAsync "
<< "failed at index " << fail_idx
<< " (src=" << srcs[fail_idx]
<< ", dst=" << dsts[fail_idx]
<< ", size=" << sizes[fail_idx]
<< "): " << cudaGetErrorString(err);
} else {
LOG(ERROR) << "NvlinkTransport: cudaMemcpyBatchAsync "
<< "failed: " << cudaGetErrorString(err);
}
// Copies [0, fail_idx) were submitted successfully → POSTED.
// Copy [fail_idx] failed → FAILED.
// Copies (fail_idx, count) were never submitted → FAILED.
for (size_t i = 0; i < fail_idx; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
for (size_t i = fail_idx; i < count; ++i) {
if (slices[i]->status == Slice::PENDING) {
slices[i]->markFailed();
}
}
} else {
for (size_t i = 0; i < count; ++i) {
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
}
#else
// Fallback for CUDA < 12.8: submit each memcpy individually
for (size_t i = 0; i < count; ++i) {
auto single_err = cudaMemcpyAsync(dsts[i], srcs[i], sizes[i],
cudaMemcpyDefault, stream);
if (single_err != cudaSuccess) {
LOG(ERROR) << "NvlinkTransport: cudaMemcpyAsync failed at "
<< "index " << i << ": "
<< cudaGetErrorString(single_err);
slices[i]->markFailed();
continue;
}
slices[i]->status = Slice::POSTED;
slices[i]->local.cuda_stream = (void *)stream;
}
return; // Slice states already set above
#endif
}
static int getNumDevices() {
static int cached_num_devices = -1;
if (cached_num_devices == -1) {
@ -203,6 +366,41 @@ Status NvlinkTransport::submitTransfer(
size_t task_id = batch_desc.task_list.size();
batch_desc.task_list.resize(task_id + entries.size());
// Synchronize with the caller's CUDA stream before issuing any memcpy.
// PyTorch uses cudaStreamPerThread (per-thread default stream), NOT the
// legacy default stream (nullptr). Recording an event on
// cudaStreamPerThread and making nvlink_stream wait for it ensures that all
// PyTorch GPU operations on source/dest buffers complete before the NVLink
// memcpy starts. This GPU-level dependency is also required by
// cudaMemcpyBatchAsync's srcAccessOrderStream attribute, which needs
// the source data to be visible through the nvlink_stream's access order.
//
// Do NOT use the legacy default stream (nullptr) for cudaEventRecord,
// as it would trigger implicit synchronization with blocking streams and
// could cause deadlocks.
cudaStream_t stream = tl_nvlink_stream.stream_;
cudaError_t sync_err =
cudaEventRecord(tl_nvlink_sync_event.event_, cudaStreamPerThread);
if (sync_err != cudaSuccess) {
LOG(ERROR) << "NvlinkTransport: cudaEventRecord on "
"cudaStreamPerThread failed: "
<< cudaGetErrorString(sync_err);
return Status::Context("cudaEventRecord failed: " +
std::string(cudaGetErrorString(sync_err)));
}
sync_err = cudaStreamWaitEvent(stream, tl_nvlink_sync_event.event_, 0);
if (sync_err != cudaSuccess) {
LOG(ERROR) << "NvlinkTransport: cudaStreamWaitEvent failed: "
<< cudaGetErrorString(sync_err);
return Status::Context("cudaStreamWaitEvent failed: " +
std::string(cudaGetErrorString(sync_err)));
}
// Phase 1: Prepare slices and collect memcpy parameters
std::vector<void *> dsts, srcs;
std::vector<size_t> sizes;
std::vector<Slice *> slices;
for (auto &request : entries) {
TransferTask &task = batch_desc.task_list[task_id];
++task_id;
@ -221,20 +419,25 @@ Status NvlinkTransport::submitTransfer(
slice->task = &task;
slice->target_id = request.target_id;
slice->status = Slice::PENDING;
slice->ts = getCurrentTimeInNano();
task.slice_list.push_back(slice);
__sync_fetch_and_add(&task.slice_count, 1);
cudaError_t err;
if (slice->opcode == TransferRequest::READ)
err = cudaMemcpy(slice->source_addr, (void *)slice->local.dest_addr,
slice->length, cudaMemcpyDefault);
else
err = cudaMemcpy((void *)slice->local.dest_addr, slice->source_addr,
slice->length, cudaMemcpyDefault);
if (err != cudaSuccess)
slice->markFailed();
else
slice->markSuccess();
void *src = (request.opcode == TransferRequest::READ)
? (void *)slice->local.dest_addr
: (void *)slice->source_addr;
void *dst = (request.opcode == TransferRequest::READ)
? slice->source_addr
: (void *)slice->local.dest_addr;
srcs.push_back(src);
dsts.push_back(dst);
sizes.push_back(slice->length);
slices.push_back(slice);
}
// Phase 2: Submit all memcpy operations
submitBatchMemcpy(slices, srcs, dsts, sizes, stream);
return Status::OK();
}
@ -248,6 +451,29 @@ Status NvlinkTransport::getTransferStatus(BatchID batch_id, size_t task_id,
std::to_string(batch_id));
}
auto &task = batch_desc.task_list[task_id];
// Poll POSTED slices for async completion via cudaStreamQuery.
// Cache the query result per stream to avoid redundant driver calls,
// since multiple slices typically share the same CUDA stream.
std::unordered_map<cudaStream_t, cudaError_t> stream_status_cache;
for (auto *slice : task.slice_list) {
if (slice && slice->status == Slice::POSTED) {
cudaStream_t stream = (cudaStream_t)slice->local.cuda_stream;
auto it = stream_status_cache.find(stream);
cudaError_t cuda_err;
if (it == stream_status_cache.end()) {
cuda_err = cudaStreamQuery(stream);
stream_status_cache[stream] = cuda_err;
} else {
cuda_err = it->second;
}
if (cuda_err == cudaSuccess) {
slice->markSuccess();
} else if (cuda_err != cudaErrorNotReady) {
slice->markFailed();
}
// cudaErrorNotReady means still in progress, keep POSTED
}
}
status.transferred_bytes = task.transferred_bytes;
uint64_t success_slice_count = task.success_slice_count;
uint64_t failed_slice_count = task.failed_slice_count;
@ -266,6 +492,31 @@ Status NvlinkTransport::getTransferStatus(BatchID batch_id, size_t task_id,
Status NvlinkTransport::submitTransferTask(
const std::vector<TransferTask *> &task_list) {
// Synchronize with the caller's CUDA stream before issuing any memcpy.
// See submitTransfer() for detailed rationale.
cudaStream_t stream = tl_nvlink_stream.stream_;
cudaError_t sync_err =
cudaEventRecord(tl_nvlink_sync_event.event_, cudaStreamPerThread);
if (sync_err != cudaSuccess) {
LOG(ERROR) << "NvlinkTransport: cudaEventRecord on "
"cudaStreamPerThread failed: "
<< cudaGetErrorString(sync_err);
return Status::Context("cudaEventRecord failed: " +
std::string(cudaGetErrorString(sync_err)));
}
sync_err = cudaStreamWaitEvent(stream, tl_nvlink_sync_event.event_, 0);
if (sync_err != cudaSuccess) {
LOG(ERROR) << "NvlinkTransport: cudaStreamWaitEvent failed: "
<< cudaGetErrorString(sync_err);
return Status::Context("cudaStreamWaitEvent failed: " +
std::string(cudaGetErrorString(sync_err)));
}
// Phase 1: Prepare slices and collect memcpy parameters
std::vector<void *> dsts, srcs;
std::vector<size_t> sizes;
std::vector<Slice *> slices;
for (size_t index = 0; index < task_list.size(); ++index) {
assert(task_list[index]);
auto &task = *task_list[index];
@ -286,20 +537,25 @@ Status NvlinkTransport::submitTransferTask(
slice->task = &task;
slice->target_id = request.target_id;
slice->status = Slice::PENDING;
slice->ts = getCurrentTimeInNano();
task.slice_list.push_back(slice);
__sync_fetch_and_add(&task.slice_count, 1);
cudaError_t err;
if (slice->opcode == TransferRequest::READ)
err = cudaMemcpy(slice->source_addr, (void *)slice->local.dest_addr,
slice->length, cudaMemcpyDefault);
else
err = cudaMemcpy((void *)slice->local.dest_addr, slice->source_addr,
slice->length, cudaMemcpyDefault);
if (err != cudaSuccess)
slice->markFailed();
else
slice->markSuccess();
void *src = (request.opcode == TransferRequest::READ)
? (void *)slice->local.dest_addr
: (void *)slice->source_addr;
void *dst = (request.opcode == TransferRequest::READ)
? slice->source_addr
: (void *)slice->local.dest_addr;
srcs.push_back(src);
dsts.push_back(dst);
sizes.push_back(slice->length);
slices.push_back(slice);
}
// Phase 2: Submit all memcpy operations
submitBatchMemcpy(slices, srcs, dsts, sizes, stream);
return Status::OK();
}